Repository navigation
perf: charge the local shuffle writer's buffers to the memory pool - #6336
dwsmith1983 wants to merge 14 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 |
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.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.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
A pressure-triggered spill in the shuffle crate: 200 hash partitions, a 16 MiB
GreedyMemoryPool, 1 MiB write buffer, Lz4, 256 MiB of Int64 input in 8192-row batches, release build, medians of three runs.The buffers take about 2 MiB of the 16 MiB pool, which is the three extra spills. The spill file buffer's grow runs during a pressure spill while the main reservation is still full, so it gets the 8 KiB fallback for the rest of the task, which is the slower final merge. Reserved bytes after the write: data buffer 1 MiB, spill buffer 8 KiB, copy buffer 1 MiB, scratch 38 KiB. Upgrading that buffer to full size right after a spill frees the main reservation, when the pool has room again, is a possible follow-up.
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.The shuffle crate, the memory pool tests, clippy, fmt and
cargo bench --no-runpass;CometNativeShuffleSuiteandCometShuffleSuitepass on Spark 3.5 andCometNativeShuffleSuiteon 4.1.