Repository navigation
fix: let a spilled final aggregate read its spill files back past its memory share - #6544
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.
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.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: DataFusion 55.1 can fill a final aggregate’s memory share with spill-merge buffers, causing its non-spillable replay to fail.
- Design approach: Identify replay through the final-aggregate consumer name and a sibling reservation holding memory, then account for refused growth through the existing overcommit mechanism.
- Correctness / compatibility analysis: Traced DataFusion’s table, merge, replay and reservation lifetimes. Normal table growth and merge admission retain their refusal paths. Replay growth remains accounted for, and releases repay debt before returning bytes to Spark. Compared relevant Spark aggregation and memory-manager sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No introduced P1/P2 issues found within this review.
- Key design decisions: The DataFusion-specific heuristic is isolated in
spill_replay.rs, with its removal tied to upstream #25383. It introduces no configuration or dependency changes. - Implementation sketch: The fair pool routes its refusals through the helper. The greedy pool tracks final-aggregate consumer totals across reservations. Shared Rust tests, two aggregate regression tests and the memory guide cover the change.
- Behavioral changes worth calling out: Spill replay can continue beyond its share or Spark’s grant, carrying the shortfall as debt. Greedy final-aggregate accounting adds map/mutex operations. Successful requests from other consumers avoid the new mutex. Performance was assessed from these call paths, without benchmark measurements.
- Suggested improvements: None meeting the reproducible P1/P2 threshold.
Reviewed full SHA 52bae81391beb4f37ce2163c1f973e6901bf9121 against ca1c143affc051998843e16cc85c3f9c76769bb5. Covered all 11 files in that comparison, including planner/cache differences originating from two base-only commits. The PR is not a draft. The snapshot and refreshed discussion contained no existing reviews or comments.
Routed skills: review-comet-pr, review-comet-memory-pr, and review-comet-shuffle-pr.
Exact-head CI: 25 checks succeeded and 15 were skipped, with no failed checks. Rust CI passed 2,011 tests. The Spark 4.1 execution shard passed 1,233 tests, including both new spill-replay regressions. Spark SQL, Iceberg and macOS jobs were skipped.
Local validation: cargo test -p datafusion-comet --lib memory_pools --offline --no-default-features passed all 30 selected tests. No local JVM or full Spark SQL suite was run, and no independent throughput or RSS measurements were taken. Project source remains unchanged.
comphead
left a comment
There was a problem hiding this comment.
Thanks for the thorough write-up and for covering OrderedFinalAggregateStream. I read the replay path in DataFusion 55.1 (hash_stream.rs, ordered_final_stream.rs and multi_level_merge.rs) against the sibling-reservation rule here, and it matches what those streams do. I did not find a correctness problem in this read. A few small points inline.
| self.spark.overcommit(), | ||
| self.pool_size | ||
| ); | ||
| return self.refuse(&mut state, reservation, additional, err); |
There was a problem hiding this comment.
As far as I can tell, this is the one refuse call site that none of the new tests reach. The fair pool test trips the fair limit first, and the shared tests keep pool_size at 1000 so only Spark refuses. Putting return Err(err) back here would not fail any of them. A sketch that should reach it (not run): a pool of 100 over a Spark that grants everything, other registers and grows to 70, then FinalHashAggregateStream[0] registers, grows merge by 25 and grows merge.new_empty() by 20. The share is 50, so 25 + 20 passes the fair check, and 70 + 25 + 20 is over the pool size.
There was a problem hiding this comment.
Added your scenario as a_final_aggregate_reading_its_spill_files_back_may_pass_the_pool_limit. With the wrapper all three refusals take the same path now, but I checked that the test does reach the pool limit: with the wrapper taken out it fails with Failed to acquire 20 bytes where 95 bytes already reserved (0 bytes overcommitted) and the pool limit is 100 bytes.
|
|
||
| /// Records `additional` bytes for `reservation` whatever the fair and pool limits, and carries | ||
| /// what Spark doesn't grant as overcommit. See [`SparkMemory`]. | ||
| fn record( |
There was a problem hiding this comment.
#5613 is open and reworks these same functions. It charges under the state lock, calls Spark without it, settles afterwards, and consumer_used there takes the consumer id. The heads conflict in fair_pool.rs and the memory guide (checked with git merge-tree). record takes the locked state and calls Spark while it is held, and all three refusals go through refuse. In #5613 the Spark refusal is handled after the lock is dropped. I expect whichever PR lands second to need more than a textual rebase.
That, and the fact that #6583 will have to undo edits in both pools and their tests, makes me wonder about a MemoryPool wrapper in spill_replay.rs, like PlanMemoryPool and TaskSharedMemoryPool. On a ResourcesExhausted from a replaying final aggregate it would call the inner grow, which already skips the limits and carries the shortfall as overcommit, so neither pool would change. The cost is a per-consumer map of its own for fair_unified and one more layer for overcommit() to look through. I have not compiled this, so please read it as an option.
There was a problem hiding this comment.
Thanks, I went with the wrapper in cc60fe7. SpillReplayPool sits between TaskSharedMemoryPool and TrackConsumersPool, and on a ResourcesExhausted from a replaying final aggregate it calls the inner grow. Both pools are back to main's versions, and git merge-tree against #5613 now only conflicts on the pool stack diagram in the guide, where both PRs add a line. The costs are the ones you listed: its own map, which only holds final aggregates, and one more line in overcommit(). It only matches ResourcesExhausted, so a failed JNI call during the replay still goes back to the aggregate, and there is a test for that.
| //! Once one of DataFusion 55's final aggregates 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 there fails the task (#6254). `FinalHashAggregateStream` does this, and so does | ||
| //! `OrderedFinalAggregateStream` itself, which DataFusion uses when the input is sorted on some of | ||
| //! the grouping keys. The merge reserves read buffers for as many spill files as fit, and those | ||
| //! buffers belong to the same consumer, so the replay often finds the consumer's share already | ||
| //! taken. The replay only asks for memory once it has aggregated a batch, so like a `grow`, the | ||
| //! request is for memory that already exists. The pools record it the way they record a `grow`, | ||
| //! carrying what Spark doesn't grant as overcommit. The replay emits every finished group after | ||
| //! each batch, so it holds about one batch of groups, and releasing memory repays the overcommit | ||
| //! first. |
There was a problem hiding this comment.
The module doc and the new guide section repeat each other closely (the replay can't spill, the merge buffers take the share, the request is recorded like a grow), and the doc on is_spill_replay repeats the guide's second paragraph about when a sibling holds memory. For overcommit, the guide keeps a short paragraph and leaves the detail to spark_memory.rs. It might be worth doing the same here. The DataFusion 55.1 facts that an upgrader has to re-check could stay in this module, and the guide could say what the pools do and link to it. The module doc could also name #6583 next to the removal condition so that the removal is easy to find.
There was a problem hiding this comment.
Agreed. The guide section is now one paragraph on what SpillReplayPool does, and it points to spill_replay.rs for the rest. The module doc keeps the DataFusion 55.1 facts an upgrade has to re-check as a list, and names #6583 next to the removal condition. I also moved the guide section to after the task-shared pools section, because #5613 adds its anchor text where it used to be.
|
|
||
| /// Refuses a `try_grow` with `err`, unless it comes from a final aggregate reading its spill | ||
| /// files back, which can't spill. That request is recorded instead; see [`spill_replay`]. | ||
| fn refuse( |
There was a problem hiding this comment.
refuse returns Ok(()) after recording the request, so return self.refuse(...) at the call sites reads as the opposite of what it can do. A name like refuse_unless_replay would say it.
There was a problem hiding this comment.
refuse is gone with the wrapper. Its try_grow matches the refusal and only calls grow for the replay, so there is nothing left to rename.
…commit-6254 # Conflicts: # spark/src/test/scala/org/apache/comet/exec/CometAggregateSuite.scala
The workaround for apache#6254 lived in both Comet pools. The fair pool sent its three refusals through a check, and the greedy pool tracked each final aggregate across its reservations. apache#5613 reworks the fair pool's try_grow, and apache#6583 has to take the workaround out again, so both would have had to rework the pools. SpillReplayPool, in spill_replay.rs, now wraps the tracked Comet pool instead. It keeps a total for each final aggregate's consumer. When the pool refuses a request from one while another of its reservations holds memory, it calls the pool's grow, which skips the fair pool's limits and carries what Spark doesn't grant as overcommit. fair_pool.rs and unified_pool.rs are back to main's versions, and overcommit() looks through the wrapper. The Rust tests run against both pool types as createPlan builds them. They now reach the fair pool's pool-limit refusal too, and check that a failed Spark call during the replay is returned rather than recorded. The guide adds the wrapper to the pool stack and keeps a short section. The DataFusion 55.1 facts that an upgrade has to re-check stay in spill_replay.rs, which names apache#6583 next to the removal condition.
… memory share (#6544) (#6713) * fix: let a spilled final aggregate read its spill files back past its 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 (#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. * fix: let an ordered final aggregate read its spill files back too 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. * docs: tighten the spill replay comments and guide 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. * test: share the spill replay tests between the pools 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 #6254 tests in CometAggregateSuite, and use checkCometAnswer, which collects once and labels the Comet answer as Comet's. * refactor: move the spill replay workaround into a pool wrapper The workaround for #6254 lived in both Comet pools. The fair pool sent its three refusals through a check, and the greedy pool tracked each final aggregate across its reservations. #5613 reworks the fair pool's try_grow, and #6583 has to take the workaround out again, so both would have had to rework the pools. SpillReplayPool, in spill_replay.rs, now wraps the tracked Comet pool instead. It keeps a total for each final aggregate's consumer. When the pool refuses a request from one while another of its reservations holds memory, it calls the pool's grow, which skips the fair pool's limits and carries what Spark doesn't grant as overcommit. fair_pool.rs and unified_pool.rs are back to main's versions, and overcommit() looks through the wrapper. The Rust tests run against both pool types as createPlan builds them. They now reach the fair pool's pool-limit refusal too, and check that a failed Spark call during the replay is returned rather than recorded. The guide adds the wrapper to the pool stack and keeps a short section. The DataFusion 55.1 facts that an upgrade has to re-check stay in spill_replay.rs, which names #6583 next to the removal condition. (cherry picked from commit ba7c892) Adapted for branch-1.1: - The tests build the pools over a fake Spark through create_pool, a seam that #6271 added to memory_pools. Only that seam is ported, not #6271's memory usage log change, which is not on branch-1.1: create_memory_pool builds the pools through create_pool, with_spark is pub(super), and the pools' new constructors, which nothing calls any more, are removed. - overcommit() and unwrap_task_shared(), which only that log reads, are not ported, so nothing looks through the wrapper. spill_replay.rs drops unwrap_spill_replay, and its tests drop their overcommit(&pool) assertions, which restate pool.reserved() minus what the fake Spark holds. The test that shrinks the replay checks pool.reserved() instead. - memory_management.md: the paragraph about the task's shared CometTaskMemoryManager, from #6261, is not on branch-1.1, so the new section follows the task-shared pools section directly. - CometAggregateSuite: only its imports differ. This branch does not import Column, CometBaseAggregateExec or CometSortAggregateExec.
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?
SpillReplayPoolinmemory_pools/spill_replay.rswraps both Comet pools, betweenTaskSharedMemoryPoolandTrackConsumersPool. When the pool refuses atry_growfrom a final aggregate reading its spill files back, the wrapper records the request with the pool'sgrowinstead. Every other call and every other refusal passes through unchanged.CometFairMemoryPoolandCometUnifiedMemoryPoolthemselves don't change.FinalHashAggregateStreamor anOrderedFinalAggregateStreamand another of the consumer's 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.ResourcesExhausted, so a failed JNI call during the replay still goes back to the aggregate.overcommit()looks through the wrapper. The memory management guide adds the wrapper to the pool stack and gives it a short section, and leaves the DataFusion details that an upgrade has to re-check tospill_replay.rs.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.
An earlier revision changed both pools directly. Review suggested the wrapper, which leaves the pools alone for #5613 and keeps the removal to one module. The
OrderedFinalAggregateStreamcase and the guide update came out of a self-review with thereview-comet-prandreview-comet-memory-prskills.How are these changes tested?
spill_replay.rs. Five of them build the pools the waycreatePlandoes, and three of those run against both pool types. They cover the replay going past Spark's grant, past the fair pool's share and past its pool limit, and the refusals that must stay: the aggregate reading its input, the merge picking its files, and another operator with a sibling reservation. They also cover a failed Spark call during the replay, the per-consumer tracking, and which consumers count as final aggregates. Five mutations each fail at least one of them: never recording, recording any error rather than only a refusal, dropping the sibling condition, not tracking shrinks, andovercommit()not looking through the wrapper.CometAggregateSuitetests check the answer and that the final aggregate spilled:Failed 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. Without the wrapper 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, and 3 of 3 withgreedy_unifiedat 96m. Without the wrapper it failed 3 of 3 at 96m with the issue's error, with either pool.mainpass, and the other 16 spill the same number of times as onmain. A sweep of 24 for the ordered aggregate, run on the earlier revision: the 5 that fail withoutOrderedFinalAggregateStreamin the check pass, and the other 19 spill the same number of times. The wrapper makes the same decisions as that revision, and the reproducer and the hash aggregate sweep gave the same spill counts on both.CometAggregateSuite,CometTaskMetricsSuiteandCometExecIteratorLifecycleSuiteon the default Spark 4.1 profile: 189 passed.cargo test -p datafusion-comet --lib: 645 passed. Clippy (--all-targets --workspace -D warnings) and rustfmt are clean.The Spark SQL tests run Comet in on-heap mode, which uses an unbounded pool, so they don't exercise this change.