Skip to content

perf: charge the local shuffle writer's buffers to the memory pool - #6336

Open
dwsmith1983 wants to merge 18 commits into
apache:mainfrom
dwsmith1983:perf/6196-shuffle-buffer-reservation
Open

dwsmith1983 wants to merge 18 commits into
apache:mainfrom
dwsmith1983:perf/6196-shuffle-buffer-reservation

Conversation

@dwsmith1983

@dwsmith1983 dwsmith1983 commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6196.

Rationale for this change

The local shuffle writer allocates buffers of spark.comet.shuffle.native.writeBufferSize bytes that no reservation covers: the data file BufWriter, the spill file BufWriter, the buffer finish_partition reads 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's TaskMemoryManager sees, so it has to fit in spark.executor.memoryOverhead.

What changes are included in this PR?

  • MultiPartitionShuffleRepartitioner::try_new hands the writer a second reservation on the consumer it already registers, reservation.new_empty(), through a new PartitionWriter::attach_buffer_reservation with a no-op default. No new consumer, so fair_unified's per-consumer share is unchanged, and the existing reservation, maxBufferBytes and memory_spilled_bytes are untouched: spill() still frees only the main reservation.
  • LocalPartitionWriter keeps 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 of write, finish_partition, write_burst_complete and finish_all, on every exit. resize makes 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 one try_grow per spill while the spill file buffer is below full size.
  • Buffers allocated up front ask for the full write buffer size with 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 through io::copy.
  • A spill under pool pressure opens the spill file while its batches are still charged, so its buffer usually starts at 8 KiB. spill() now calls write_burst_complete after reservation.free(), and the local writer then asks for the full size again (PartitionedSpill::restore_buffer): try_grow of the difference, then a new BufWriter at 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 because write_burst_complete returns 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 through io::copy instead of one pread.
  • Memory that already exists when it is charged (the scratch, the zstd context) uses grow, which Comet's pools record as overcommit. The scratch's transient peak between syncs, the write buffer size plus one block, stays uncharged.
  • The single-partition, empty-schema and remote shuffle writers are unchanged; they have no reservation to charge, and registering one would change fair_unified's shares. The tuning guide and the native shuffle contributor doc list them as untracked.

Benchmark

shuffle_bench over TPC-H lineitem (the first 6,001,215 rows of the SF10 file, the SF1 row count), 200 partitions, lz4, no buffer limit, ci profile, 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:

pool main this PR before the restore this PR
16 MiB 0.128 s (67 spills) 0.130 s (73 to 77) 0.130 s (76 to 77)
32 MiB 0.127 s (34 to 35) 0.146 s (35 to 37) 0.129 s (36 to 37)
64 MiB 0.126 s (17) 0.130 s (18) 0.120 s (18)

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 through io::copy and a smaller one through pread; with maxBufferBytes set, spill_count and memory_spilled_bytes match a run without buffer charging; with zstd, the reservation grows by the context size after a write and shrinks after write_burst_complete and finish_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-run pass; CometNativeShuffleSuite and CometShuffleSuite pass on Spark 3.5 and CometNativeShuffleSuite on 4.1.

dwsmith1983 and others added 2 commits September 29, 2026 00:07
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.
@github-actions github-actions Bot added enhancement New feature or request performance area:shuffle Shuffle (JVM and native) labels Sep 28, 2026

@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.

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 maxBufferBytes and 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.

@andygrove

Copy link
Copy Markdown
Member

This is a light fully automated review since there are so many PRs open.

buffer_bytes_held at native/shuffle/src/writers/local/local_partition_writer.rs:195 charges codec_context.retained_bytes(), and that workspace is allocated by libzstd through libc malloc (zstd-sys only swaps the allocator on wasm). So with this PR the pool tracks C library memory that the AccountingAllocator behind native_allocated never sees, and two docs now say the opposite. The "Non-Rust allocations" bullet at docs/source/contributor-guide/memory_management.md:406 says neither the memory pool nor the allocation counters see such memory. Step 3 of "Sizing the Overhead from the Memory Usage Log" at docs/source/user-guide/latest/tuning/memory.md:175 says neither allocated nor reserved includes memory from native C libraries such as zstd.

When spark.comet.shuffle.compression.codec is zstd, reserved now holds one context per multi-partition shuffle task that is mid-spill or in its final write (about 1.3 MiB at the default level and close to 8 MiB at levels 7 and 8), and that also comes straight off the allocated - reserved difference the tuning page sizes the overhead from. Charging it still seems right to me, since it is real memory Spark should budget for. Could those two passages call out the shuffle writer's zstd context as the exception, alongside the list this PR already updates higher up in the same tuning page?

…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.
@dwsmith1983

Copy link
Copy Markdown
Contributor Author

Could those two passages call out the shuffle writer's zstd context as the exception, alongside the list this PR already updates higher up in the same tuning page?

Done in 0d23cb4. The "Non-Rust allocations" bullet in memory_management.md now says the pool tracks such memory only where an operator charges it by hand, and names the multi-partition shuffle writer's zstd context as the one case today, in reserved and never in allocated. Step 3 of the sizing section in tuning/memory.md says the same, and that the allocated - reserved difference understates what has to fit in the overhead by one context per such task, with the two sizes you gave.

@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@sunchao CI passed on a888fbeb6 and it has your approval. Could it go into the merge queue?

@dwsmith1983
dwsmith1983 requested a review from sunchao October 4, 2026 15:04

@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.

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 maxBufferBytes and 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.

@dwsmith1983

Copy link
Copy Markdown
Contributor Author

@andygrove when you get a chance, could you review this and approve the CI run? The branch is synced with main.

@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.

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.
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.

Charge the native shuffle writer's write buffers to the memory pool

3 participants