andygrove opened a new pull request, #6544:
URL: https://github.com/apache/datafusion-comet/pull/6544

   ## Which issue does this PR close?
   
   Closes #6254.
   
   ## Rationale for this change
   
   DataFusion 55 moved Comet's final hash aggregates onto 
`FinalHashAggregateStream` (apache/datafusion#24061). Once it has spilled, it 
merges its sorted spill files and replays them through an 
`OrderedFinalAggregateStream` that has no way to spill, so a refused memory 
request during the replay fails the task. The merge reserves read buffers for 
as many spill files as fit, in a sibling reservation of the same consumer, so 
the replay often finds the consumer's share already taken. 1.0.0 didn't fail 
here because DataFusion 54's `GroupedHashAggregateStream` ignored a refused 
reservation during the replay.
   
   The upstream fix, apache/datafusion#25383, leaves the replay room when the 
merge picks its files, but it is only on DataFusion `main`. This works around 
the failure in Comet's memory pools until we upgrade.
   
   ## What changes are included in this PR?
   
   - A new `memory_pools/spill_replay.rs` decides when a refused request is 
recorded as overcommit instead: the consumer is a `FinalHashAggregateStream`, 
and another of its reservations holds memory. In DataFusion 55.1 that only 
happens while the replay grows and the merge holds its read buffers. While the 
aggregate reads its input, its table is its only reservation holding memory, so 
a refusal still makes it spill. The merge picks its files while nothing else is 
held, so a refusal still limits how many it opens.
   - The fair pool sends all three of its refusals (fair limit, pool limit and 
Spark) through that check, and records a matching request the way `grow` does.
   - The greedy pool now tracks what each final aggregate's consumer holds 
across its reservations, so it can make the same check. Other consumers aren't 
tracked and still never take a lock.
   
   The overcommit is the mechanism `grow` already uses. Spark grants what it 
can, the rest is carried as debt, releases repay the debt first, and while any 
is outstanding every other `try_grow` in the task is refused. The replay emits 
every finished group after each batch, so what it holds stays around one batch 
of groups. In the issue's reproducer it asked for 1.8 MB on top of a 24 MB 
share.
   
   This should be removed once Comet's DataFusion includes 
apache/datafusion#25383.
   
   ## How are these changes tested?
   
   - Five new Rust tests in `fair_pool.rs` and `unified_pool.rs` cover the 
replay and the refusals that must stay: the aggregate reading its input, the 
merge picking its files, and another operator with a sibling reservation. 
Dropping the check fails the replay tests, and dropping the sibling condition 
fails the others.
   - A new `CometAggregateSuite` test, "final aggregate that has spilled reads 
its spill files back (issue #6254)", gives the final aggregate a 3 MiB pool. On 
`main` it fails with `Failed to acquire 115952 bytes where this consumer 
already holds 3031120 bytes and the fair limit is 3145728 bytes`. It checks the 
answer and that the final aggregate spilled.
   - The issue's reproducer (96m off-heap, `local[4]`, 4 shuffle partitions) 
passed 3 of 3 runs at 96m, 80m and 88m with the right answer. With the check 
disabled it failed 3 of 3 at 96m with the issue's error. With `greedy_unified` 
it passed 3 of 3, and failed 2 of 2 with the check disabled.
   - A sweep of 18 data sizes and pool fractions: the 2 that fail on `main` 
pass, and the other 16 spill the same number of times as on `main`.
   - `CometAggregateSuite`, `CometTaskMetricsSuite` and 
`CometExecIteratorLifecycleSuite` on the default Spark 4.1 profile: 178 passed. 
`cargo test -p datafusion-comet --lib`: 564 passed. Clippy (`--all-targets -D 
warnings`) and rustfmt are clean, and the tests compile on Spark 3.5 / Scala 
2.12.
   
   The Spark SQL tests run Comet in on-heap mode, which uses an unbounded pool, 
so they don't exercise this change.
   


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