Skip to content

Remove the spill replay workaround from the memory pools once DataFusion includes apache/datafusion#25383 #6583

Description

@andygrove

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

Describe the potential solution

After #6544 merges and Comet upgrades to a DataFusion release that includes apache/datafusion#25383:

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

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions