Skip to content

A native final aggregate that has spilled can fail the task during its replay #6254

Description

@andygrove

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:

org.apache.comet.CometNativeException: Additional allocation failed for FinalHashAggregateStream[0] with top memory consumers (across reservations) as:
  FinalHashAggregateStream[0]#28(can spill: true) consumed 23.1 MB, peak 23.1 MB.
Error: Failed to acquire 1828080 bytes plus 0 bytes overcommitted, only got 917424 bytes. Reserved: 24248400 bytes

After spilling, DataFusion 55's FinalHashAggregateStream merges the sorted spill runs and replays them through an OrderedFinalAggregateStream in Sorted mode (into_replay_stream in aggregates/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 as try_grow allows, 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 its DiskManagerBuilder without 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:

spark.range(0, 2000000, 1, 4)
  .selectExpr("id", "concat(sha2(cast(id as string), 256), sha2(cast(id * 7 as string), 256)) as s")
  .write.parquet(path)

spark.read.parquet(path)
  .groupBy("s").agg(count(lit(1)), max("id"))
  .write.format("noop").mode("overwrite").save()

The task fails with the error above. The refusal comes from Spark's per-task share, 24 MB here, not from the fair_unified limit. 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_in in prepare_datafusion_session_context would 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.

Activity

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

Metadata

Metadata

Assignees

Labels

area:aggregationHash aggregates, aggregate expressionsarea:memoryMemory pools, reservations, OOM handlingbugSomething isn't workingrequires-triage

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions