Conversation
… 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.
This was referenced 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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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 anOrderedFinalAggregateStreamthat 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'sGroupedHashAggregateStreamignored 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
OrderedFinalAggregateStreaminstead. 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?
memory_pools/spill_replay.rsdecides when a refused request is recorded instead: the consumer is aFinalHashAggregateStreamor anOrderedFinalAggregateStream, 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.growdoes.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 thatgrowalready uses. Spark grants what it can, the rest is carried as debt, releases repay the debt first, and while any is outstanding every othertry_growin 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
OrderedFinalAggregateStreamcase and the guide update came out of a self-review with thereview-comet-prandreview-comet-memory-prskills.How are these changes tested?
fair_pool.rs,unified_pool.rsandspill_replay.rscover 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.CometAggregateSuitetests check the answer and that the final aggregate spilled:mainit fails withFailed to acquire 115952 bytes where this consumer already holds 3031120 bytes and the fair limit is 3145728 bytes.OrderedFinalAggregateStream, and gives it a 2.5 MiB pool. WithoutOrderedFinalAggregateStreamin the check it fails withFailed to acquire 112528 bytes where this consumer already holds 2555968 bytes and the fair limit is 2621440 bytes.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. Withgreedy_unifiedit passed 3 of 3, and failed 2 of 2 with the check disabled.mainpass, and the other 16 spill the same number of times as onmain. A sweep of 24 for the ordered aggregate: the 5 that fail withoutOrderedFinalAggregateStreamin the check pass, and the other 19 spill the same number of times.CometAggregateSuite,CometTaskMetricsSuiteandCometExecIteratorLifecycleSuiteon 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.