Skip to content

fix: let a spilled final aggregate read its spill files back past its memory share - #6544

Draft
andygrove wants to merge 4 commits into
apache:mainfrom
andygrove:fix/spill-replay-overcommit-6254
Draft

andygrove wants to merge 4 commits into
apache:mainfrom
andygrove:fix/spill-replay-overcommit-6254

Conversation

@andygrove

@andygrove andygrove commented Oct 2, 2026 •

Copy link
Copy Markdown
Member

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.

A final aggregate whose input is sorted on some of its grouping keys, such as one over a sort in the same native plan, runs as an OrderedFinalAggregateStream instead. It spills and replays its spill files the same way, and fails the same way.

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 instead: the consumer is a FinalHashAggregateStream or an OrderedFinalAggregateStream, 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 take no lock unless Spark refuses them.
  • The memory management guide describes the exception.

The replay asks for memory only after it has aggregated a batch, so like a grow, the request is for memory that already exists. Recording it uses the overcommit that 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. #6583 tracks that.

The OrderedFinalAggregateStream case and the guide update came out of a self-review with the review-comet-pr and review-comet-memory-pr skills.

How are these changes tested?

  • Six new Rust tests in fair_pool.rs, unified_pool.rs and spill_replay.rs cover the replay, the consumers that count as final aggregates, 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.
  • Two new CometAggregateSuite tests check 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 for the hash aggregate: the 2 that fail on main pass, and the other 16 spill the same number of times as on main. A sweep of 24 for the ordered aggregate: the 5 that fail without OrderedFinalAggregateStream in the check pass, and the other 19 spill the same number of times.
  • CometAggregateSuite, CometTaskMetricsSuite and CometExecIteratorLifecycleSuite on the default Spark 4.1 profile: 179 passed. cargo test -p datafusion-comet --lib: 565 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.

… memory share

DataFusion 55's FinalHashAggregateStream replays its merged spill files through a stream that
can't spill, so a refused memory request there fails the task. The merge's read buffers are a
sibling reservation of the same consumer and take as many spill files as fit, so the replay often
finds the consumer's share already taken (apache#6254).

Both Comet pools now record that request as overcommit instead of refusing it. They recognize it
as a request from a FinalHashAggregateStream consumer while another of its reservations holds
memory, which in DataFusion 55.1 happens only during the replay. The replay emits its finished
groups after every batch, so the overcommit stays around one batch of groups, and releases repay
it first. Refusals while the aggregate reads its input, and while the merge picks its files, are
unchanged.

Remove this once Comet's DataFusion includes apache/datafusion#25383.
@andygrove andygrove added bug Something isn't working area:aggregation Hash aggregates, aggregate expressions area:memory Memory pools, reservations, OOM handling regression A bug that did not affect the most recent Comet release labels Oct 2, 2026
DataFusion runs a final aggregate whose input is sorted on some of its
grouping keys as an OrderedFinalAggregateStream. It merges and replays its
spill files the same way FinalHashAggregateStream does, so its replay
failed the task the same way. Treat its consumer as a final aggregate too.

Explain in spill_replay.rs why recording the replay's request is safe: the
replay asks for memory only after it has aggregated a batch, so the memory
already exists, as it does for a grow. Describe the exception in the memory
management guide, which said the fair pool always refuses a request that
fails its local checks.
The fair pool is the only one with a share and a pool total to skip, so say
so. The greedy pool takes its tracking lock for other consumers when Spark
refuses them, not never. Point whoever removes the workaround at the tests
that show whether the replay still needs it.
fair_pool.rs and unified_pool.rs each had a copy of the same two spill
replay scenarios. Run them once, in spill_replay.rs, against both pool
types as createPlan builds them, which also checks that the wrappers pass
register through to the greedy pool's tracking. Keep a greedy pool test for
its per-consumer map, and drop a test import that the pool's own imports
now cover.

Share the final aggregate spill checks between the two apache#6254 tests in
CometAggregateSuite, and use checkCometAnswer, which collects once and
labels the Comet answer as Comet's.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:aggregation Hash aggregates, aggregate expressions area:memory Memory pools, reservations, OOM handling bug Something isn't working regression A bug that did not affect the most recent Comet release

Projects

None yet

Development

Successfully merging this pull request may close these issues.

A native final aggregate that has spilled can fail the task during its replay

1 participant