Skip to content

perf: read shuffle blocks in 64 KiB pieces instead of through Channels.newChannel - #6805

Queued
comphead wants to merge 5 commits into
apache:mainfrom
comphead:perf-shuffle-reader-chunked-reads
Queued

comphead wants to merge 5 commits into
apache:mainfrom
comphead:perf-shuffle-reader-chunked-reads

Conversation

@comphead

@comphead comphead commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Part of #6528. It also covers the Channels.newChannel part of item R6 in the shuffle performance review, #5905.

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?

  • Both readers call one new helper, CometShuffleBlockIterator.readFully, instead of wrapping the stream in a channel. It fills the existing buffer from in in pieces of up to the read buffer size, through a heap array reused per thread, and 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.
  • A new config, spark.comet.shuffle.readBufferSize, sets the size of each read. It defaults to 64 KiB and is read on the executor when a reader is created, like spark.comet.tracing.enabled.
  • The direct-read stream that CometBlockStoreShuffleReader.readAsRawStream() returns now calls context.killTaskIfInterrupted() before every read, as the Celeborn reader's stream already does. The JDK channel checked the thread's interrupt status on every read, and on this path it was the only check inside a shuffle block, because native code pulls the stream on the task thread with no InterruptibleIterator in between. A killed task on this path now stops at its next read, also when the kill does not interrupt its thread, which this path used to ignore. NativeBatchDecoderIterator still stops a killed task at its next batch, through Spark's InterruptibleIterator.
  • Cleanup is unchanged, since close() already closes in directly. Task completion still closes in, which unblocks a stalled remote read.

Larger reads are not free: FileInputStream.read allocates a native buffer of the requested size for every read above 8 KiB. On macOS, freeing that buffer shows up as madvise even at 64 KiB, at about 3% of the hot task above. The sweep below found 32 KiB to 1 MiB within noise of each other, and 8 KiB and 4 MiB slower.

How are these changes tested?

  • A new check, run by "shuffle readers read in pieces of the read buffer size without asking for available" in CometCelebornShuffleReaderSuite, reads a 1 MiB block through each reader with a 16 KiB read buffer size. It uses a stream that counts available() calls and records the largest read, and asserts that there are no available() calls and that the largest read is exactly 16 KiB.

  • A new test, "a killed task stops at the next read of a local shuffle's direct-read stream" in CometCelebornShuffleReaderSuite, reads a local shuffle through CometShuffledBatchRDD.computeAsShuffleBlockIterator, the entry point of the direct-read path. It marks the task context killed after the first batch and checks that the next hasNext() throws TaskKilledException.

  • 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. With AQE on, skew join handling was off, so that the hot partition stays whole and the join stays in Comet (AQE skew join split makes the join fall back to Spark with no fallback reason #6530). Each table comes from one run with its arms interleaved. Absolute times differ by a few percent between runs on this machine.

    The hot reduce task of Comet, median of 5 runs, before and after this change at 64 KiB:

    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

    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.

    Sweeping spark.comet.shuffle.readBufferSize, median of 5 runs. With AQE off, the shuffle read operator time isolates the read path. With AQE on, only the whole hot task is available, because the direct-read path reports no read time (ShuffleScan metrics are dropped under AQE: shuffle read and decode time missing on the direct-read path #6804).

    Read buffer size AQE on: hot task AQE off: shuffle read operator
    8 KiB 1,413 ms 995 ms
    16 KiB 1,346 ms 959 ms
    32 KiB 1,305 ms 895 ms
    64 KiB (default) 1,281 ms 920 ms
    128 KiB 1,365 ms 926 ms
    256 KiB 1,281 ms 927 ms
    1 MiB 1,298 ms 932 ms
    4 MiB 1,358 ms 1,053 ms

    How long a killed reduce task keeps running, from a second local benchmark that is not part of this PR either. One map task writes 67,108,864 rows as one shuffle block of 1.4 GiB, and one reduce task reads it with a native project. The job is cancelled 1 s into the read, with and without spark.job.interruptOnCancel. Time from the cancel to the end of the task, median of 5:

    Read path Kill main This PR
    AQE on, direct read interrupts the thread 1 ms 1 ms
    AQE on, direct read does not interrupt the thread 3,365 ms 1 ms
    AQE off, JVM decode either 1 ms 1 ms

    The check costs nothing measurable. The hot task of the skewed join above, with AQE on, took 1,369 ms with the check and 1,368 ms with it switched off in a local build, median of 9 runs interleaved in one JVM.

…ls.newChannel`

Both shuffle readers, `NativeBatchDecoderIterator` and `CometShuffleBlockIterator`,
read blocks through `Channels.newChannel(in)`. For a stream that is not a plain
`FileInputStream`, the JDK channel copies at most 8 KiB per read and calls
`in.available()` before every read after the first, which costs an fstat and an
lseek on a local shuffle file. Read `in` directly in pieces of up to 64 KiB instead.
@github-actions github-actions Bot added enhancement New feature or request performance area:shuffle Shuffle (JVM and native) labels Oct 8, 2026
Add `spark.comet.shuffle.readBufferSize`, 64 KiB by default. Both shuffle readers now
call one `CometShuffleBlockIterator.readFully`, which owns a single buffer per thread,
and one test covers both readers.
The strict-warnings build rejects the implicit Int to Long widening in
`createWithDefault`, so convert the default explicitly.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for this. Reading in larger pieces and dropping the available() call per piece is a nice contained win. I found one behavior change against 1.1 that I'd like to see addressed before this merges, so I'm requesting changes. A killed task on the direct read path no longer stops until its current shuffle block runs out. Details and a fix I tried locally are inline.

return -1;
}
throw new EOFException("Data corrupt: unexpected EOF while reading batch header");
readFully(inputStream, headerBuf, readBufferSize);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Dropping Channels.newChannel also drops the only thing that stopped a killed task on the direct read path. The JDK channel checks the thread's interrupt status on every read and throws ClosedByInterruptException, and native calls hasNext() on the task thread. readAsRawStream() in CometBlockStoreShuffleReader returns a bare SequenceInputStream, so this path has no InterruptibleIterator. With this change a killed task keeps reading until the shuffle block it is on runs out.

I measured it locally with a 20M row shuffle into a single block, read by a native project, and killed the reduce task 1 s into its 5.6 s read. With an interrupting kill, which is what speculation and Spark Connect and Thrift server cancellation use, the task stopped 11 ms after the kill on main and kept going for 4.6 s on this branch. The skewed partition this PR is aimed at is also the task most likely to get a speculative copy, and the losing copy now holds its core and memory until it finishes reading its current block.

Could readAsRawStream() wrap the stream so that each read calls context.killTaskIfInterrupted(), the way the Celeborn reader's stream already does? I tried that on top of this branch. Kills with and without an interrupt both stopped within about 10 ms, and a task that isn't killed took the same time. A test that marks the task context interrupted and checks that the next read throws TaskKilledException would keep this from regressing again.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 5c025a8acb as you suggested. I measured it locally with one 1.4 GiB shuffle block read by a native project and cancelled 1 s into the read. The direct-read task now stops 1 ms after the cancel, with or without an interrupt. Before the fix it kept reading for about 3.4 s either way, and on main it did so for kills without an interrupt. The table against main is in the description.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for turning this around so quickly. The kill check in readAsRawStream() is what I had in mind. It matches what the Celeborn stream already does, and the new test would fail without it. I have two asks before I approve. The first is inline: moving the two benchmarks out of this PR, which also clears the lint failure.

The second is one more pass over the description. The paragraph about interrupts says a killed task stops at the next batch through InterruptibleIterator. That is still true for NativeBatchDecoderIterator, but the direct-read path now stops at its next read through the check in readAsRawStream(). That check also catches kills without an interrupt, which that path used to ignore. The "What changes" list could mention it too, and the testing section still says the benchmark is not part of this PR. The description becomes the commit message, so it is worth getting right.

* variant name only labels the classes on the classpath. `paths=aqe,noaqe` selects the read
* paths, `rows=N` the rows of the shuffle and `trials=N` the trials.
*/
object CometShuffleReadKillBenchmark extends CometBenchmarkBase {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could both new benchmarks come out of this PR? They read as tooling for this review, and both scaladocs say so. Their kill-check A/B sets spark.comet.benchmark.skipKillCheck, which only a temporary build reads. Nothing in the tree reads it, so the default fix and head variants here run the same code, and so do the two killcheck arms of CometHashJoinTaskTimeBenchmark.

CometHashJoinTaskTimeBenchmark looks like the skew-join reproduction that #6528 says is not in a PR yet. I think it deserves its own PR against that issue, without the kill-check mode, where it can be reviewed on its own. At 1,084 lines it would be the largest benchmark in the tree. Dropping it here also fixes the red lint jobs, which all come from two s interpolators with nothing to interpolate at lines 633 and 658.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed, removed both in 1346697f8c, which also clears the lint failures. They were local tools for this review and came in by accident with the fix. I'll keep CometHashJoinTaskTimeBenchmark for #6528 and open it as its own PR there, without the kill-check mode.

@andygrove
andygrove dismissed their stale review October 9, 2026 21:48

The kill check in readAsRawStream() fixes the issue behind this review, and the new test passes in CI.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed the full 13-file diff from 412468e2817207eff0c1907d1c4aac8c3f06680d to 5c025a8acb2278a576895cbff6c781f930267969, including both benchmarks and the documentation changes. The PR is not a draft. Read all supplied reviews, comments, and threads before forming findings.

No additional introduced P1/P2 issues found within this review. The previously reported cancellation issue is addressed. The already-reported benchmark/lint blocker remains unresolved.

Summary

  • Prior state and problem: Both shuffle readers used Channels.newChannel, which splits wrapped-stream reads into 8 KiB pieces and calls available(). This adds stream calls on large shuffle reads.
  • Design approach: Replace the channel with shared readFully logic, use a reusable heap array, and expose spark.comet.shuffle.readBufferSize with a 64 KiB default. The local direct-read stream also gains explicit task-kill checks.
  • Correctness: Short reads, payload preservation, consecutive frames, clean EOF, truncated headers and bodies, invalid lengths, and idempotent close passed the local comparison. The framing and native decode interfaces remain unchanged. Spark's InterruptibleIterator and TaskContextImpl sources support the cancellation approach. No additional correctness issue met the P1/P2 bar.
  • Compatibility analysis: Checked relevant Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Their task-kill contract is identical. Compared the touched reader paths with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1. The intended differences are read granularity, configuration, and direct-read cancellation. No shuffle-format or SQL-result change was identified.
  • Key design decisions: Resolve the buffer size when constructing readers, validate its range before narrowing to Int, and keep cancellation checks at the stream boundary where direct native reads bypass InterruptibleIterator. Celeborn retains its existing stream-level kill checks.
  • Implementation sketch: CometShuffleBlockIterator.readFully fills the existing destination buffers. Both readers call it for headers and bodies. Local and Celeborn construction paths pass the configured size. Added tests cover read sizes, absence of available() calls, and cancellation after a batch.
  • Performance: The local 1 MiB frame probe reduced stream reads from 129 to 17 and available() calls from 127 to zero. This validates the mechanism, not end-to-end speedup. The default retains one 64 KiB JVM-heap scratch array per task thread, outside the native memory pool. The reported task-time improvements remain the author's measurements.
  • Design: The production change is contained and preserves existing decoding and ownership boundaries. The full diff also adds substantial benchmark tooling and removes or rewrites FAQ/comparison material. Those changes were included in the review; no additional P1/P2 finding arose from them.
  • Abstraction & complexity: One shared read loop avoids diverging implementations without introducing another reader hierarchy. Thread-local scratch storage follows the existing decoder-buffer pattern. The benchmarks' inert skipKillCheck comparison is already covered by an existing review thread.
  • Behavioral changes worth calling out: Direct reads now detect a task marked killed even without a thread interrupt. JVM decoding retains batch-boundary checks through InterruptibleIterator. Removing the channel also removes its automatic interrupt-driven closure during a read. The existing description-review comment already requests clarification of these distinctions.
  • Suggested improvements: Resolve the existing benchmark/lint blocker. Its A/B variants currently select the same production implementation, and redundant interpolators fail required lint checks. No additional P1/P2 recommendations are raised.

Routed skills: review-comet-pr, review-comet-shuffle-pr, review-comet-memory-pr, and review-comet-ffi-pr, with the relevant contributor documentation.

Exact-head CI: Checks associated with this SHA show successful Linux Spark 4.1 runtime suites, TPC-H/TPC-DS verification, native build/Rust tests, and Celeborn reflection compatibility. The shuffle job passed all 555 tests, including both added checks. Six Java-lint matrix jobs and Scala syntactic lint failed, leaving Required Checks red. Logs identify the already-discussed interpolators in CometHashJoinTaskTimeBenchmark.scala at lines 633 and 658. The workflow tested GitHub's merge ref 83b7594, containing this head, rather than the standalone checkout. See CI run.

Validation limits: A disposable Java 21 harness compiled the actual base and head reader sources and passed 307 base assertions and 355 head assertions. No full local Maven/native build or end-to-end benchmark was run. Spark's own SQL suites, other Spark runtime profiles, macOS, Iceberg suites, and benchmark checks were not exercised by this CI run. No real remote-service cancellation reproduction was performed. The project checkout remains unchanged.

`CometHashJoinTaskTimeBenchmark` and `CometShuffleReadKillBenchmark` were local tools for
measuring this change and its review, and came in with the previous commit by accident. They
stay local, and the skew-join benchmark will get its own PR against apache#6528.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for moving the benchmarks out and for the description pass. Both asks from my last review are covered. The reader code hasn't changed since the round where the kill check went in, and CI is green at 1346697f8c, including the two new checks in the shuffle job, so I'm approving.

While re-reading the direct-read close path I noticed that closing the stream early still fetches every remaining block. That predates this PR, so I filed #6865 for it.

@andygrove
andygrove added this pull request to the merge queue Oct 11, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:shuffle Shuffle (JVM and native) enhancement New feature or request performance

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants