Skip to content

fix: [branch-1.1] let a spilled final aggregate read its spill files back past its memory share (#6544) - #6713

Merged
andygrove merged 1 commit into
apache:branch-1.1from
andygrove:backport-6544-branch-1.1
Oct 6, 2026
Merged

andygrove merged 1 commit into
apache:branch-1.1from
andygrove:backport-6544-branch-1.1

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Closes #6254 on branch-1.1. #6544 already closed it on main.

Rationale for this change

This is the branch-1.1 backport of #6544, which has the backport-1.1 label. #6254 is a 1.1.0 regression (it carries the regression label). DataFusion 55 moved Comet's final hash aggregates onto FinalHashAggregateStream. Once one 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. A final aggregate over input sorted on some of its grouping keys runs as an OrderedFinalAggregateStream and fails the same way.

The upstream fix, apache/datafusion#25383, is only on DataFusion main, and this branch is on DataFusion 55.1. #6544 works around the failure in Comet's memory pools, and this backport takes that workaround as it is.

Not proposed for branch-1.0: it is on DataFusion 54, which doesn't fail here, and #6544 has no backport-1.0 label.

What changes are included in this PR?

A cherry-pick (-x) of #6544's squash commit; see #6544 for what the change does. In short, a new SpillReplayPool wraps both Comet pools, between TaskSharedMemoryPool and TrackConsumersPool. When the pool refuses a try_grow from a final aggregate reading its spill files back, the wrapper records the request with the pool's grow instead. Every other call and every other refusal passes through unchanged. The change comes with seven Rust tests, two CometAggregateSuite tests and a section in the memory management guide.

The wrapper's logic in spill_replay.rs, its wiring in create_pool, both CometAggregateSuite test bodies and the guide's new text apply as on main. Three files conflicted (mod.rs, memory_management.md and CometAggregateSuite.scala), and the two pool files needed a small port. The adaptations:

How are these changes tested?

Run locally on this branch with JDK 17:

  • Rust: cargo test -p datafusion-comet --lib passes (580 tests), including the seven in spill_replay.rs. cargo fmt --check and workspace clippy with --all-targets -D warnings are clean. My local toolchain is Rust 1.97, and CI uses stable.
  • CometAggregateSuite, CometTaskMetricsSuite and CometExecIteratorLifecycleSuite on the default Spark 4.1 profile: 153 pass, with 2 that CometAggregateSuite already ignores on this branch. That includes the two new A native final aggregate that has spilled can fail the task during its replay #6254 tests, which each check that a native final aggregate ran and spilled.
  • On Spark 3.5 and Scala 2.12 the sources compile, scalastyle, spotless and scalafix (RemoveUnused, in CHECK mode) pass, and the two new tests pass.
  • Without the fix, both new tests fail. I took the SpillReplayPool wrapper out of create_pool, which leaves this branch's own pool stack, rebuilt the native library and ran the two tests. They fail with Failed to acquire 115952 bytes where this consumer already holds 3031120 bytes and the fair limit is 3145728 bytes and Failed to acquire 112528 bytes where this consumer already holds 2555968 bytes and the fair limit is 2621440 bytes, the same errors that fix: let a spilled final aggregate read its spill files back past its memory share #6544 reports. So this branch has the bug and the tests catch it. The fix changes no JVM code, so the JVM side was the same in both runs.
  • prettier --check on the guide passes.

By CI's path routing, which I checked with dev/ci/compute-changes.py, this PR runs the all-profiles Linux build, macOS, the Spark SQL jobs for 3.5, 4.0 and 4.1, and the four Iceberg versions. The Spark SQL tests run Comet on-heap, which uses an unbounded pool, so they don't exercise this change, and no run-* label is needed.

… memory share (apache#6544)

* 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 (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.

* 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 apache#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 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.

(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 apache#6271 added to memory_pools. Only that seam is ported, not apache#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 apache#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.
@andygrove andygrove added bug Something isn't working area:aggregation Hash aggregates, aggregate expressions area:memory Memory pools, reservations, OOM handling labels Oct 6, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: DataFusion 55.1 can fill a final aggregate’s memory share with spill-merge buffers, leaving its non-spillable replay unable to proceed.
  • Design approach: SpillReplayPool identifies replay through the consumer name and an occupied sibling reservation, then records refused growth through the existing grow mechanism.
  • Correctness / compatibility analysis: Verified DataFusion’s reservation lifetimes, merge admission, replay and cleanup paths. Ordinary refusals remain intact, and releases repay overcommit 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 wrapper isolates the DataFusion-specific workaround from both pool implementations. Its assumptions are documented, and issue #6583 tracks removal. The backport preserves the upstream wrapper logic.
  • Implementation sketch: Both off-heap pools receive the wrapper through create_pool. Seven Rust tests cover replay, rejection and accounting. Two Scala regression tests check results and native spilling. The memory guide describes the exception.
  • Behavioral changes worth calling out: Replay may exceed the fair share, pool limit or Spark grant while remaining accounted for. The wrapper adds dispatch and final-aggregate map/mutex bookkeeping. Other consumers avoid its tracking mutex. Performance assessment was limited to code paths.
  • Suggested improvements: None meeting the reproducible P1/P2 threshold.

Reviewed all six changed files at 8bcacb60688e0b3d03e7551fd24cd2eca7f2d56a against 7b7eec69282a7abc33a0e8e687ca594c72f140ec. The PR remains non-draft. The snapshot and refreshed discussion contained no reviews, issue comments or review threads.

Routed skills: review-comet-pr and review-comet-memory-pr.

Exact-head CI: The primary run has 7 successful checks, 18 in progress and 4 skipped, with no failures reported. Rust, JVM, macOS, Spark SQL and Iceberg validation remains pending. Completed label runs skipped substantive suites and do not establish integration success.

Local validation: cargo test -p datafusion-comet --lib memory_pools --offline --no-default-features --locked passed all 26 selected tests, including all seven new tests. No local JVM, Spark SQL or Iceberg suites were run. Default features and throughput/RSS measurements were not validated locally. Project source remains unchanged.

@andygrove
andygrove merged commit e9efd9f into apache:branch-1.1 Oct 6, 2026
163 checks passed
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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants