Repository navigation
Conversation
…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.
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
left a comment
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
The kill check in readAsRawStream() fixes the issue behind this review, and the new test passes in CI.
sunchao
left a comment
There was a problem hiding this comment.
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 callsavailable(). This adds stream calls on large shuffle reads. - Design approach: Replace the channel with shared
readFullylogic, use a reusable heap array, and exposespark.comet.shuffle.readBufferSizewith 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
InterruptibleIteratorandTaskContextImplsources 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.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1. 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 bypassInterruptibleIterator. Celeborn retains its existing stream-level kill checks. - Implementation sketch:
CometShuffleBlockIterator.readFullyfills 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 ofavailable()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
skipKillCheckcomparison 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
left a comment
There was a problem hiding this comment.
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.
Which issue does this PR close?
Part of #6528. It also covers the
Channels.newChannelpart 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 callsin.available()before every read after the first. A local shuffle block arrives as Spark's stream wrappers over aFileInputStream, so every 8 KiB costs areadsyscall and anavailable()call, which is anfstatand anlseek, 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.availablealone took about 6% of that task's samples, on both paths.What changes are included in this PR?
CometShuffleBlockIterator.readFully, instead of wrapping the stream in a channel. It fills the existing buffer frominin 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 callsavailable()any more.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, likespark.comet.tracing.enabled.CometBlockStoreShuffleReader.readAsRawStream()returns now callscontext.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 noInterruptibleIteratorin 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.NativeBatchDecoderIteratorstill stops a killed task at its next batch, through Spark'sInterruptibleIterator.close()already closesindirectly. Task completion still closesin, which unblocks a stalled remote read.Larger reads are not free:
FileInputStream.readallocates a native buffer of the requested size for every read above 8 KiB. On macOS, freeing that buffer shows up asmadviseeven 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 countsavailable()calls and records the largest read, and asserts that there are noavailable()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 throughCometShuffledBatchRDD.computeAsShuffleBlockIterator, the entry point of the direct-read path. It marks the task context killed after the first batch and checks that the nexthasNext()throwsTaskKilledException.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:
CometShuffleBlockIterator)NativeBatchDecoderIterator)In profiles of that task,
available()disappears, and the time spent inreadsyscalls 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).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:mainThe 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.