You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
#6544 works around #6254 in Comet's memory pools. Once a final aggregate has spilled, DataFusion 55 merges its spill files and replays them through an aggregate that can't spill. This applies to FinalHashAggregateStream, and to OrderedFinalAggregateStream when the input is sorted on some of the grouping keys. The merge reserves read buffers for as many spill files as fit, in a sibling reservation of the same consumer, so the replay's next request is often refused and the task fails.
The workaround, native/core/src/execution/memory_pools/spill_replay.rs, records such a refused request the way grow does, and carries what Spark doesn't grant as overcommit. It depends on DataFusion's consumer names and on when those aggregates hold sibling reservations, which can change in any DataFusion release. The overcommit is also memory that Spark doesn't know about.
apache/datafusion#25383 fixes the cause upstream: when the merge picks its spill files, it leaves room for the replay. It merged into DataFusion main on 2026-09-17. Its branch-55 backport, apache/datafusion#25814, was closed without merging, so the fix first ships in DataFusion 56 (#6410).
Delete spill_replay.rs and its call sites. In fair_pool.rs that is refuse and the three calls to it in try_grow. In unified_pool.rs it is the final_aggregates map, the register, unregister and track updates to it, and the replay branch in try_grow. Delete the Rust tests that cover them.
Remove the "Final aggregates reading their spill files back" section from docs/source/contributor-guide/memory_management.md, and the link to it in the fair pool section.
Keep the two A native final aggregate that has spilled can fail the task during its replay #6254 tests in CometAggregateSuite, "final aggregate that has spilled reads its spill files back" and "ordered final aggregate that has spilled reads its spill files back". On DataFusion 55.1 both fail without the workaround, so they show whether the upstream fix is enough on its own.
Additional context
apache/datafusion#25383 reserves room for the replay equal to the merge's buffers. Its description calls that a budgeting policy rather than a guarantee that every aggregation fits. So the upgrade alone doesn't prove the workaround can go. If either test above fails without it, keep the workaround and find out why. #6254 has a larger reproducer that is worth running too.
What is the problem the feature request solves?
#6544 works around #6254 in Comet's memory pools. Once a final aggregate has spilled, DataFusion 55 merges its spill files and replays them through an aggregate that can't spill. This applies to
FinalHashAggregateStream, and toOrderedFinalAggregateStreamwhen the input is sorted on some of the grouping keys. The merge reserves read buffers for as many spill files as fit, in a sibling reservation of the same consumer, so the replay's next request is often refused and the task fails.The workaround,
native/core/src/execution/memory_pools/spill_replay.rs, records such a refused request the waygrowdoes, and carries what Spark doesn't grant as overcommit. It depends on DataFusion's consumer names and on when those aggregates hold sibling reservations, which can change in any DataFusion release. The overcommit is also memory that Spark doesn't know about.apache/datafusion#25383 fixes the cause upstream: when the merge picks its spill files, it leaves room for the replay. It merged into DataFusion
mainon 2026-09-17. Itsbranch-55backport, apache/datafusion#25814, was closed without merging, so the fix first ships in DataFusion 56 (#6410).Describe the potential solution
After #6544 merges and Comet upgrades to a DataFusion release that includes apache/datafusion#25383:
spill_replay.rsand its call sites. Infair_pool.rsthat isrefuseand the three calls to it intry_grow. Inunified_pool.rsit is thefinal_aggregatesmap, theregister,unregisterandtrackupdates to it, and the replay branch intry_grow. Delete the Rust tests that cover them.docs/source/contributor-guide/memory_management.md, and the link to it in the fair pool section.CometAggregateSuite, "final aggregate that has spilled reads its spill files back" and "ordered final aggregate that has spilled reads its spill files back". On DataFusion 55.1 both fail without the workaround, so they show whether the upstream fix is enough on its own.Additional context
apache/datafusion#25383 reserves room for the replay equal to the merge's buffers. Its description calls that a budgeting policy rather than a guarantee that every aggregation fits. So the upgrade alone doesn't prove the workaround can go. If either test above fails without it, keep the workaround and find out why. #6254 has a larger reproducer that is worth running too.