Repository navigation
A native final aggregate that has spilled can fail the task during its replay #6254
Copy link
Copy link
Closed
Labels
area:aggregationHash aggregates, aggregate expressionsHash aggregates, aggregate expressionsarea:memoryMemory pools, reservations, OOM handlingMemory pools, reservations, OOM handlingbugSomething isn't workingSomething isn't workingpriority:mediumFunctional bugs, performance regressions, broken featuresFunctional bugs, performance regressions, broken featuresregressionA bug that did not affect the most recent Comet releaseA bug that did not affect the most recent Comet release
Description
Activity
- addedarea:aggregationHash aggregates, aggregate expressionsHash aggregates, aggregate expressionsarea:memoryMemory pools, reservations, OOM handlingMemory pools, reservations, OOM handling
on Sep 26, 2026 - addedpriority:mediumFunctional bugs, performance regressions, broken featuresFunctional bugs, performance regressions, broken featuresregressionA bug that did not affect the most recent Comet releaseA bug that did not affect the most recent Comet releaseand removed
on Sep 28, 2026 This is a regression in 1.1.0. The 1.1.0 regression audit (#6399) ran this issue's reproducer on both releases with the Spark 4.1 profile. 1.1.0-rc1 failed all 6 runs at 96m of off-heap memory with
Additional allocation failed for FinalHashAggregateStream[0] ... Failed to acquire 1828080 bytes, and also failed at 80m and 88m. 1.0.0 passed all 16 runs between 64m and 128m, spilling about 40 files and returning the right answer. It's tracked in #6402.- added 5 commits that reference this issue
on Oct 1, 2026 - added a commit that references this issue
on Oct 6, 2026 - added a commit that references this issue
on Oct 6, 2026
Metadata
Metadata
Assignees
Labels
area:aggregationHash aggregates, aggregate expressionsHash aggregates, aggregate expressionsarea:memoryMemory pools, reservations, OOM handlingMemory pools, reservations, OOM handlingbugSomething isn't workingSomething isn't workingpriority:mediumFunctional bugs, performance regressions, broken featuresFunctional bugs, performance regressions, broken featuresregressionA bug that did not affect the most recent Comet releaseA bug that did not affect the most recent Comet release
Describe the bug
Under memory pressure, a native final hash aggregate that has already spilled can fail the task instead of spilling again, even though its consumer is marked as spillable:
After spilling, DataFusion 55's
FinalHashAggregateStreammerges the sorted spill runs and replays them through anOrderedFinalAggregateStreaminSortedmode (into_replay_streaminaggregates/hash_stream.rs, datafusion-physical-plan 55.1.0). That stream is built without a spill context, so a refused resize of its table comes back as an error. The merge feeding it takes as many runs astry_growallows, because the spill merge fan-in defaults to unlimited (DEFAULT_MAX_SPILL_MERGE_FAN_IN = 0). The merge's reservation belongs to the same consumer, so it grows until the pool refuses and leaves nothing for the replay.apache/datafusion#25423 describes the same failure. apache/datafusion#25424 closed it by bounding the fan-in in the test only (
datafusion.runtime.max_spill_merge_fan_in = 2), so the behavior is unchanged in 55.1. Comet builds itsDiskManagerBuilderwithout a fan-in (jni_api.rs#L834-L836), so it runs with the unlimited default.Steps to reproduce
Use
spark.memory.offHeap.size=96m,local[4]and 4 shuffle partitions, and run a grouped aggregate over 2M distinct 128-character strings:The task fails with the error above. The refusal comes from Spark's per-task share, 24 MB here, not from the
fair_unifiedlimit. I haven't run the same query with Comet disabled at this budget.Expected behavior
A final aggregate that has spilled completes, spilling again or merging in more passes if it has to.
Additional context
Setting a finite fan-in with
DiskManagerBuilder::with_max_spill_merge_fan_ininprepare_datafusion_session_contextwould leave room for the replay, at the cost of more merge passes. The real fix belongs upstream: either the replay stream can spill, or the merge leaves headroom for it.