Describe the bug
A native plan with no JVM input, such as a native Parquet scan feeding a native sort, runs on a Tokio task that sends its batches to the Spark task thread over a channel (jni_api.rs:1152-1189). executePlan treats a closed channel as the end of the stream (None => Ok(-1), jni_api.rs:1214-1217), so a producer that was dropped looks the same as one that finished.
The runtime drops it when it shuts down. release_runtime (jni_api.rs:427-432) calls shutdown_timeout, which cancels every spawned task at its next yield. The producer's sender goes with it, blocking_recv on the task thread returns None, and the task finishes normally with whatever the plan had produced so far.
CometExecutorPlugin.shutdown calls NativeBase.releaseNative(). Executor.stop() is also the executor's JVM shutdown hook, and in 3.4.3, 3.5.8, 4.0.1 and 4.1.3 it calls threadPool.shutdown(), which neither waits for nor interrupts running tasks, then the plugins' shutdown(), and only then env.stop(). So when an executor gets a SIGTERM (YARN preemption, a Kubernetes eviction, a spot reclaim without decommissioning), a task on this path can report success with truncated output while the RpcEnv is still up.
What gets truncated depends on the consumer. A result task returns partial rows, a Spark writer over a Comet child commits a partial file, and a JVM shuffle writer writes a partial map output, which outlives the executor when an external shuffle service serves it. The native local shuffle writer fails instead, because it checks that the plan was drained before it publishes offsets (jni_api.rs:1413-1428). I haven't checked the Celeborn destination, which skips that check.
Steps to reproduce
A throwaway suite on main at 634e37d08 (Spark 4.1, local[4], 2g off-heap):
- Write 4.8M rows to 16 Parquet files.
- Run
spark.read.parquet(path).sortWithinPartitions("s").rdd.count() in a Future.
- Call
NativeBase.releaseNative() 1.5 s after four tasks have started.
The count comes back as 1,200,000 instead of 4,800,000, with no exception and no warning. The four tasks that were running log Finished task ... result sent to driver about 200 ms after the release. Tasks that start afterwards run on a new runtime.
With #6261 applied, one run truncated the same way. Another failed with task N was cancelled, which is the join error from a sort subtask. Which one you get depends on where the plan is when the runtime goes.
Expected behavior
A task whose producer was dropped fails, instead of reporting end of stream.
Additional context
The producer knows when the stream has ended, so the end could be made explicit: send a final message when stream.next() returns None, and treat a channel that closes without it as an error. #6261's BatchProducer keeps the task's JoinHandle, which would also let executePlan tell a completed producer from a cancelled one.
Separately, release_runtime's doc comment says the runtime is shut down in the background so that the calling JNI thread is not blocked, but shutdown_timeout(Duration::from_secs(3)) blocks the caller for up to 3 s.
Describe the bug
A native plan with no JVM input, such as a native Parquet scan feeding a native sort, runs on a Tokio task that sends its batches to the Spark task thread over a channel (
jni_api.rs:1152-1189).executePlantreats a closed channel as the end of the stream (None => Ok(-1),jni_api.rs:1214-1217), so a producer that was dropped looks the same as one that finished.The runtime drops it when it shuts down.
release_runtime(jni_api.rs:427-432) callsshutdown_timeout, which cancels every spawned task at its next yield. The producer's sender goes with it,blocking_recvon the task thread returnsNone, and the task finishes normally with whatever the plan had produced so far.CometExecutorPlugin.shutdowncallsNativeBase.releaseNative().Executor.stop()is also the executor's JVM shutdown hook, and in 3.4.3, 3.5.8, 4.0.1 and 4.1.3 it callsthreadPool.shutdown(), which neither waits for nor interrupts running tasks, then the plugins'shutdown(), and only thenenv.stop(). So when an executor gets a SIGTERM (YARN preemption, a Kubernetes eviction, a spot reclaim without decommissioning), a task on this path can report success with truncated output while the RpcEnv is still up.What gets truncated depends on the consumer. A result task returns partial rows, a Spark writer over a Comet child commits a partial file, and a JVM shuffle writer writes a partial map output, which outlives the executor when an external shuffle service serves it. The native local shuffle writer fails instead, because it checks that the plan was drained before it publishes offsets (
jni_api.rs:1413-1428). I haven't checked the Celeborn destination, which skips that check.Steps to reproduce
A throwaway suite on
mainat634e37d08(Spark 4.1,local[4], 2g off-heap):spark.read.parquet(path).sortWithinPartitions("s").rdd.count()in aFuture.NativeBase.releaseNative()1.5 s after four tasks have started.The count comes back as 1,200,000 instead of 4,800,000, with no exception and no warning. The four tasks that were running log
Finished task ... result sent to driverabout 200 ms after the release. Tasks that start afterwards run on a new runtime.With #6261 applied, one run truncated the same way. Another failed with
task N was cancelled, which is the join error from a sort subtask. Which one you get depends on where the plan is when the runtime goes.Expected behavior
A task whose producer was dropped fails, instead of reporting end of stream.
Additional context
The producer knows when the stream has ended, so the end could be made explicit: send a final message when
stream.next()returnsNone, and treat a channel that closes without it as an error. #6261'sBatchProducerkeeps the task'sJoinHandle, which would also letexecutePlantell a completed producer from a cancelled one.Separately,
release_runtime's doc comment says the runtime is shut down in the background so that the calling JNI thread is not blocked, butshutdown_timeout(Duration::from_secs(3))blocks the caller for up to 3 s.