comphead opened a new pull request, #6805:
URL: https://github.com/apache/datafusion-comet/pull/6805
## Which issue does this PR close?
Part of #6528.
## Rationale for this change
Both Comet shuffle readers read blocks through `Channels.newChannel(in)`:
- `NativeBatchDecoderIterator`, the JVM decode path, which a native operator
reads through when AQE is off.
- `CometShuffleBlockIterator`, the direct-read path behind AQE query stages.
For a stream that is not a plain `FileInputStream`, the JDK returns a
channel that copies at most 8 KiB per read and calls `in.available()` before
every read after the first. A local shuffle block arrives as Spark's stream
wrappers over a `FileInputStream`, so every 8 KiB costs a `read` syscall and an
`available()` call, which is an `fstat` and an `lseek`, each through its own
JNI call.
In a skewed shuffled hash join whose hot reduce task reads 1.5 GiB of
shuffle data (#6528), `FileInputStream.available` alone took about 6% of that
task's samples, on both paths.
## What changes are included in this PR?
- `NativeBatchDecoderIterator` and `CometShuffleBlockIterator` read the
stream directly. A `readFully` helper fills the existing buffer from `in` in
pieces of up to 64 KiB, through a thread-local heap array. It stops short only
at the end of the stream, so the end-of-stream and corruption checks are
unchanged.
- Neither reader calls `available()` any more, and reads are no longer
limited to 8 KiB.
- Cleanup is unchanged, since `close()` already closes `in` directly. One
difference: the JDK channel turned a thread interrupt during a read into
`ClosedByInterruptException`. Now a killed task finishes the read in progress
and stops at the next batch, through Spark's `InterruptibleIterator`. Task
completion still closes `in`, which unblocks a stalled remote read.
Pieces are capped at 64 KiB because `FileInputStream.read` allocates a
native buffer of the requested size for every read above 8 KiB. On macOS,
freeing that buffer still shows up as `madvise` at 64 KiB, at about 3% of the
hot task above.
## How are these changes tested?
- Two new tests, one per reader:
- "shuffle block iterator reads large blocks in big pieces without asking
for available" in `CometNativeShuffleSuite`.
- "decoder reads large blocks in big pieces without asking for available
bytes" in `CometCelebornShuffleReaderSuite`, which runs a new check in
`NativeBatchDecoderIteratorLifecycleChecks`.
Each reads a 1 MiB block through a stream that counts `available()` calls
and records the largest read, and asserts that there are no `available()` calls
and that a read is larger than 8 KiB.
- The existing lifecycle and concurrency checks of both readers cover end of
stream, corruption reporting and cleanup.
- Measured with a local benchmark that is not part of this PR: Spark 4.1.3,
`local[8]`, a release build on an Apple M3 Max, a shuffled hash join over
2,097,152 rows with 3 nested and 100 flat columns, 95% of them on one key. The
hot reduce task of Comet, median of 5 runs:
| Shuffle read path | Before | After |
|---|---:|---:|
| AQE on, direct read (`CometShuffleBlockIterator`) | 1,411 ms | 1,227 ms |
| AQE off (`NativeBatchDecoderIterator`) | 1,795 ms | 1,478 ms |
With AQE on, skew join handling was off, so that the hot partition stays
whole and the join stays in Comet (#6530). In profiles of that task,
`available()` disappears, and the time spent in `read` syscalls and in the W^X
toggles that macOS on Apple Silicon makes at every JNI call drops.
--
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]