andygrove opened a new issue, #6295: URL: https://github.com/apache/datafusion-comet/issues/6295
### 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. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
