andygrove opened a new pull request, #6362:
URL: https://github.com/apache/datafusion-comet/pull/6362

   ## Which issue does this PR close?
   
   Closes #6294.
   
   ## Rationale for this change
   
   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. `executePlan` in `native/core/src/execution/jni_api.rs` 
took a closed channel for the end of the stream. The channel also closes when 
the task is cancelled, and a runtime cancels every task it has when it shuts 
down, which `release_runtime` does when the executor plugin stops. A Spark task 
whose producer was cancelled therefore ended successfully with only the batches 
it had read so far: a result task returned partial rows, a writer over a Comet 
child could commit a partial file, and a JVM shuffle writer could write partial 
map output. Only the local native shuffle writer failed, because it checks 
afterwards that the plan published its partition offsets.
   
   ## What changes are included in this PR?
   
   Stacked on #6261, which restructures this code. Until it lands, the diff 
here also includes its commits. Only the commit after 29ce01d16 is new:
   
   - 54069eab3 fix: fail a native plan whose producer task is cancelled before 
its stream ends
   
   Related to #6261.
   
   - `BatchProducer`'s task sets a `stream_ended` flag once it has sent the 
stream's last batch, before its sender is dropped. A new 
`BatchProducer::next_batch` returns `None` for a closed channel only when the 
flag is set. Otherwise it fails with `CometNativeException: Comet Internal 
Error: The Tokio task running the native plan was cancelled before the plan 
produced all of its output, for instance because Comet's Tokio runtime was shut 
down`. Every consumer of the plan's output sees that error from `executePlan`, 
including a remote shuffle destination, which does not check that the plan was 
drained. The local native shuffle writer now fails at the same point, with this 
message rather than the later "the plan was not drained to completion".
   - `executePlan` reads batches through `next_batch`, and so do #6261's 
producer tests, so nothing reads the channel directly.
   - The end is a flag rather than a final message on the channel, because a 
final message needs room in the channel, which holds up to two batches. A task 
cancelled while waiting to send it would fail a Spark task whose output was 
complete. The flag is stored without waiting, before the channel closes.
   - `releasePlan` is unchanged. `BatchProducer::stop` does not look at the 
flag, so releasing a plan that finished, or whose consumer stopped early, for 
example for a LIMIT, is not an error.
   - Two sentences on this in the development guide's description of the async 
I/O path.
   
   ## How are these changes tested?
   
   - Two native unit tests next to #6261's `BatchProducer` tests in 
`jni_api.rs`:
     - `a_batch_producer_cancelled_mid_stream_is_an_error_not_the_end`: the 
consumer reads the first batch of a stream that then waits for input, the 
runtime is shut down with `shutdown_timeout` while the consumer waits for the 
next batch, and `next_batch` must fail. With the old end-of-stream check it 
failed with `expected an error, got Ok(None)`.
     - `a_batch_producer_ends_cleanly_once_its_stream_has_ended`: the task of a 
two-batch stream finishes before the consumer reads anything, and the runtime 
is shut down. The consumer still gets both batches and then the end, and `stop` 
succeeds.
   - `cargo test -p datafusion-comet`: 551 passed, 5 ignored. The two new tests 
passed 40 repeated runs. `cargo clippy --all-targets --workspace -- -D 
warnings` is clean.
   - End to end, with a throwaway suite that is not in this PR: a 
single-partition native Parquet scan of 1,000,000 rows at batch size 1024, read 
by a task that calls `NativeBase.releaseNative()` on its own thread after its 
first row and then drains the iterator. At 29ce01d16, #6261's head, the job 
returned 3,072 rows without an error: the batch the task had read and the two 
in the channel. With this change the task fails with the `CometNativeException` 
above. The test is not included because it shuts down the executor-wide runtime 
that the rest of a test JVM shares, and `releaseNative()` only takes effect 
once per JVM.
   - `CometExecIteratorLifecycleSuite`, `CometExecSuite`, 
`CometNativeShuffleSuite` and `CometTaskMetricsSuite` on the default Spark 4.1 
profile: 243 tests passed, and the new error does not appear in their log.
   


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