dwsmith1983 opened a new pull request, #6336: URL: https://github.com/apache/datafusion-comet/pull/6336
## 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. - 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`. - 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 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. | | spill count | memory_spilled_bytes | spilled bytes | insert (spills) | final write | total | |---|---|---|---|---|---|---| | main | 36 | 607,136,768 | 153,994,145 | 692 ms | 25 ms | 717 ms | | this PR | 39 | 613,866,496 | 153,571,301 | 709 ms | 36 ms | 745 ms | 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 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. 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. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
