andygrove opened a new issue, #6294:
URL: https://github.com/apache/datafusion-comet/issues/6294

   ### 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):
   
   1. Write 4.8M rows to 16 Parquet files.
   2. Run `spark.read.parquet(path).sortWithinPartitions("s").rdd.count()` in a 
`Future`.
   3. 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.
   


-- 
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]

Reply via email to