Repository navigation
perf: charge the local shuffle writer's buffers to the memory pool - #6336
dwsmith1983 wants to merge 18 commits into
Conversation
The local shuffle writer's data file and spill file buffers, its spill copy buffer, the scratch it encodes blocks into and its zstd context were held outside any reservation, about 4 MiB per task that spilled with the 1 MiB default write buffer. The repartitioner now hands the writer a second reservation on the consumer it already registers, so no new consumer changes the fair share, and the writer keeps that reservation equal to what its buffers hold. A spill still frees only the main reservation, so memory_spilled_bytes and maxBufferBytes are unchanged. When the pool refuses a buffer's full size the writer uses an 8 KiB buffer and charges that. Closes apache#6196.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Multi-partition local shuffle retained file buffers, encoding scratch, and zstd workspace without charging them to the memory pool.
- Design approach: Give the writer a
reservation.new_empty()sibling on the existing consumer, separate from buffered input batches. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Reservation ownership, error cleanup, spill copying, and partition offsets remain consistent. Compared Spark memory acquisition and shuffle commit contracts across 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. Partition assignment and IPC format are unchanged.
- Key design decisions: Reusing the consumer preserves fair-pool shares. Separate reservations preserve
maxBufferBytesand spilled-memory metrics. Refused file/copy buffers fall back to at most 8 KiB. - Implementation sketch: One reservation hook and a shared allocation helper keep accounting with the buffer owner. Capacity-based synchronization includes retained zstd memory and handles failed writes.
- Behavioral changes worth calling out: Accounting for these buffers can trigger earlier spills. Small fallback buffers trade I/O throughput for lower memory usage. Unchanged reservation sizes avoid pool calls. Single-partition, empty-schema, and remote writers retain their existing behavior.
- Suggested improvements: None at the P1/P2 bar.
Reviewed the full eight-file diff from base 65a0cda1cf62877a37ec1a1f5f9ebc9dd1405ebd to head f0aa284d8bcdc7c75826c8c409f1bc1ee8329be7. Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-shuffle-pr. The PR is not a draft. The discussion snapshot contained no existing reviews, issue comments, inline comments, or threads.
Validation: All 177 shuffle unit tests passed. Two disposable fault-injection tests also passed, covering accounting after encoding failures and failed spill-file reads under ample and zero-capacity pools. The checkout was restored clean.
Exact-head CI: Comet CI and CodeQL report action_required. The label workflow remains queued. There is no completed CI test verdict. Local validation did not rerun JVM end-to-end suites, the Spark SQL matrix, production JNI-backed pool tests, or benchmarks.
|
This is a light fully automated review since there are so many PRs open.
When |
…ves but the counters miss The multi-partition local shuffle writer charges its zstd context to the pool, and libzstd allocates that workspace through libc malloc, so it is in `reserved` without ever being in `allocated`. Two passages said no C library memory reaches the pool; both now call this context out as the exception, and the overhead sizing step says the difference understates what has to fit by one context per such task.
Done in 0d23cb4. The "Non-Rust allocations" bullet in |
|
@sunchao CI passed on |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Multi-partition local shuffle retained file buffers, encoding scratch, and zstd workspace without charging them to the memory pool.
- Design approach: Give the writer a separate
reservation.new_empty()reservation on the existing consumer and synchronize it with retained buffer sizes. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Spill replay, partition offsets, reservation cleanup, and fallback-buffer output remain consistent. Checked Spark memory-acquisition and shuffle-commit contracts against upstream sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0.
- Key design decisions: Sharing the consumer preserves fair-pool shares. Keeping buffer accounting separate preserves
maxBufferBytesand spilled-memory metrics. Refused file and copy buffers fall back to at most 8 KiB. - Implementation sketch: A reservation hook and shared allocation helper keep accounting with the buffer owner. Capacity synchronization handles retained scratch and zstd memory without pool calls when the size is unchanged.
- Behavioral changes worth calling out: Compared with
branch-1.1, the intended changes reduce available pool headroom and can trigger earlier spills. Smaller fallback buffers trade throughput for lower memory use. The author’s pressure benchmark reports approximately 4% higher runtime with three additional spills. Single-partition, empty-schema, and remote writers retain their previous behavior. - Suggested improvements: None at the P1/P2 bar. The existing zstd accounting documentation concern is addressed.
Reviewed the full nine-file diff from base fef94f6cd78b18151dff57b7a936798385356de5 to head 8274c7a55b6bf44f7f7126941f26f34377685ef2, including existing reviews and discussions. The PR remains open and non-draft. Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-shuffle-pr.
Validation: All 177 shuffle unit tests passed. Two disposable fault-injection tests also passed, checking accounting after encoding failures and failed spill-file opens with ample and zero-capacity pools. The checkout was restored clean.
Exact-head CI: Comet CI remains in progress. At final inspection, 17 checks succeeded, 14 were skipped, and seven remained running, with no reported failures. Running checks include Rust tests and Spark 4.1 shuffle tests. CodeQL passed. Spark SQL and macOS jobs were skipped. Local validation did not run JVM end-to-end suites, the Spark SQL matrix, production JNI-backed pool tests, or benchmarks.
|
@andygrove when you get a chance, could you review this and approve the CI run? The branch is synced with |
andygrove
left a comment
There was a problem hiding this comment.
With maxBufferBytes at its default of 0, every spill is a pressure spill, so ensure_spill_file_created asks for the spill file buffer while the batch reservation that just failed to grow is still full. That try_grow is usually refused, and it always is when a batch is smaller than the write buffer, so the buffer drops to 8 KiB. It then stays at 8 KiB for the rest of the task, even though the spill frees the whole batch reservation a moment later. I checked this with a GreedyMemoryPool probe at 4, 16 and 64 MiB with no buffer limit. The spill buffer finished the task at 8 KiB every time, including the 16 MiB run that had 15 MiB free before the final write.
On TPC-H SF1 lineitem through shuffle_bench (200 partitions, lz4, 16 MiB pool, ci profile, 7 interleaved runs), shuffle write time went from 0.140s on main to 0.185s with this PR, and the total went up about 1%. At 32 and 64 MiB I couldn't see a difference.
Could spill() call write_burst_complete after self.reservation.free(), and have the local writer retry the full-size spill buffer there? PartitionedSpill can try_grow the difference, flush the small buffer, take the writer back with BufWriter::into_parts, and wrap it in a full-size one. I tried that locally. Write time went back to 0.142s and all 179 shuffle tests still pass. A test could spill under pool pressure with no buffer limit and check that the spill buffer is back to full size once the spill has freed its batches.
With no buffer limit every spill is a pressure spill, so the spill file buffer was requested while the spilled batches were still charged and fell back to 8 KiB for the rest of the task. spill() now calls write_burst_complete after freeing the batch reservation, and the local writer asks the pool for the full size again, moving any buffered bytes into the larger buffer without writing to the file.
Which issue does this PR close?
Closes #6196.
Rationale for this change
The local shuffle writer allocates buffers of
spark.comet.shuffle.native.writeBufferSizebytes that no reservation covers: the data fileBufWriter, the spill fileBufWriter, the bufferfinish_partitionreads spilled ranges through, and the recycled scratch that blocks are encoded into, plus the zstd context when zstd is configured. With the 1 MiB default, a multi-partition task that has spilled holds about 4 MiB that neither Comet's pool nor Spark'sTaskMemoryManagersees, so it has to fit inspark.executor.memoryOverhead.What changes are included in this PR?
MultiPartitionShuffleRepartitioner::try_newhands the writer a second reservation on the consumer it already registers,reservation.new_empty(), through a newPartitionWriter::attach_buffer_reservationwith a no-op default. No new consumer, sofair_unified's per-consumer share is unchanged, and the existing reservation,maxBufferBytesandmemory_spilled_bytesare untouched:spill()still frees only the main reservation.LocalPartitionWriterkeeps that reservation equal to the bytes its buffers hold: the data file buffer capacity, the spill file buffer capacity, the copy buffer length, the scratch capacity and the zstd context size. It syncs at the end ofwrite,finish_partition,write_burst_completeandfinish_all, on every exit.resizemakes no pool call when the size is unchanged, so steady state adds no pool traffic per batch; in practice it is one grow and one shrink per spill when zstd is on, plus the one-time buffer grows, plus onetry_growper spill while the spill file buffer is below full size.try_grow. When the pool refuses, the buffer is 8 KiB and that is what gets charged. The data file buffer is built before the consumer registers, so on a refusal it is rebuilt at the fallback size while still empty. The spill file and copy buffers are reserved after their file is open, so a failed open leaves no charge. Reads of a range larger than the copy buffer already go throughio::copy.spill()now callswrite_burst_completeafterreservation.free(), and the local writer then asks for the full size again (PartitionedSpill::restore_buffer):try_growof the difference, then a newBufWriterat full size that takes the old buffer's bytes. The bytes always fit the larger empty buffer, so the swap writes nothing to the file and cannot fail, which matters becausewrite_burst_completereturns nothing. A refusal leaves the buffer as it was, and the next spill asks again. Once grown, the spill buffer stays charged until the final write, so under steady pressure the batches have that much less room, as they would with a full-size buffer from the start; a pool smaller than about one write buffer keeps the 8 KiB buffer. The data file buffer is sized when the writer attaches, before this consumer holds anything, so pressure from the writer's own batches does not shrink it. The copy buffer is sized during the final write while the last batches are still charged and is not grown back; a range larger than it goes throughio::copyinstead of onepread.grow, which Comet's pools record as overcommit. The scratch's transient peak between syncs, the write buffer size plus one block, stays uncharged.fair_unified's shares. The tuning guide and the native shuffle contributor doc list them as untracked.Benchmark
shuffle_benchover TPC-Hlineitem(the first 6,001,215 rows of the SF10 file, the SF1 row count), 200 partitions, lz4, no buffer limit,ciprofile, 7 rounds per pool with the three builds' order rotated each round, in a Linux container. Medians of the shuffle write time, with spill counts:Charging the buffers takes about 2 MiB of the pool, which is the few extra spills. Before the restore, the spill buffer stayed at 8 KiB after a pressure spill for the rest of the task; with it, the buffer is back to full size after each spill and write time is back at main's.
How are these changes tested?
Shuffle crate tests with a
GreedyMemoryPool: the buffers are charged after attach, after a spill and after the write, the main reservation is empty after a spill, and the pool reads zero after drop; with pools of 1.5 MiB and 512 KiB the refused buffers fall back to 8 KiB, the reservation equals the bytes held after every step, and the output is byte-identical to a 1 GiB pool run; a refused copy buffer copies a range larger than it throughio::copyand a smaller one throughpread; withmaxBufferBytesset,spill_countandmemory_spilled_bytesmatch a run without buffer charging; with zstd, the reservation grows by the context size after a write and shrinks afterwrite_burst_completeandfinish_all; growing only the encode scratch moves the reservation by exactly the scratch's growth; with no buffer limit and pools of 1, 4 and 16 MiB, the spill buffer is back to full size after every pressure spill, the reservation equals the bytes held, and every row is written once; at the spill file level, a buffer the full pool cut to 8 KiB stays there while the pool is full, then grows by exactly the difference with its buffered bytes moved and nothing written until the flush, and a spill file that failed a write is not grown back.The shuffle crate, the memory pool tests, clippy, fmt and
cargo bench --no-runpass;CometNativeShuffleSuiteandCometShuffleSuitepass on Spark 3.5 andCometNativeShuffleSuiteon 4.1.