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]

Reply via email to