Describe the bug
When Spark kills a task it interrupts the task thread, and the task stops at its next interruption check. A Comet task waiting for its native plan's next batch is parked in native code instead, in blocking_recv for a plan with no JVM input or in Handle::block_on otherwise (jni_api.rs:1200 and 1226). Thread.interrupt wakes neither. The task stays there until the plan produces a batch or finishes, and it keeps its core and its memory until then. A plan with no JVM input also keeps running on its Tokio task after that, until it next tries to send. That is the cause of #2453, which #6261 fixes.
Measured on main at 634e37d08 (local[4], one 2.4M-row file, setJobGroup(..., interruptOnCancel = true)):
- Native scan feeding a native sort, read with
foreachPartition(_.hasNext), cancelled 2 s after the task started: the job failed 49 ms after the cancel, but the task kept running for another 13.1 s.
- The same query with Comet disabled: the task ended 324 ms after the cancel.
- Native scan feeding a native shuffle write, cancelled after 1.5 s: the map task ran to completion, 16.0 s after the cancel.
Anything that relies on killing tasks to free their slots waits on this: job cancellation, speculative copies that lose, and stages that are no longer needed.
Expected behavior
A killed task stops its native plan in roughly the time Spark takes to stop its own operators.
Additional context
CelebornNativeShuffleDestination.watchForCancellation (CometNativeShuffleWriter.scala:564-580) already polls TaskContext.isInterrupted() on a scheduled executor and aborts the pusher. The same pattern could call a native entry point that cancels the plan. For a plan with no JVM input, that entry point could stop the producer and drop its stream, as #6261's BatchProducer::stop does. For the block_on path, it could set a flag that makes next_batch return an error and wake it. #4175 covers the narrower case of JVM UDF dispatch.
Describe the bug
When Spark kills a task it interrupts the task thread, and the task stops at its next interruption check. A Comet task waiting for its native plan's next batch is parked in native code instead, in
blocking_recvfor a plan with no JVM input or inHandle::block_onotherwise (jni_api.rs:1200and1226).Thread.interruptwakes neither. The task stays there until the plan produces a batch or finishes, and it keeps its core and its memory until then. A plan with no JVM input also keeps running on its Tokio task after that, until it next tries to send. That is the cause of #2453, which #6261 fixes.Measured on
mainat634e37d08(local[4], one 2.4M-row file,setJobGroup(..., interruptOnCancel = true)):foreachPartition(_.hasNext), cancelled 2 s after the task started: the job failed 49 ms after the cancel, but the task kept running for another 13.1 s.Anything that relies on killing tasks to free their slots waits on this: job cancellation, speculative copies that lose, and stages that are no longer needed.
Expected behavior
A killed task stops its native plan in roughly the time Spark takes to stop its own operators.
Additional context
CelebornNativeShuffleDestination.watchForCancellation(CometNativeShuffleWriter.scala:564-580) already pollsTaskContext.isInterrupted()on a scheduled executor and aborts the pusher. The same pattern could call a native entry point that cancels the plan. For a plan with no JVM input, that entry point could stop the producer and drop its stream, as #6261'sBatchProducer::stopdoes. For theblock_onpath, it could set a flag that makesnext_batchreturn an error and wake it. #4175 covers the narrower case of JVM UDF dispatch.