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]