andygrove opened a new issue, #6254:
URL: https://github.com/apache/datafusion-comet/issues/6254

   ### 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](https://github.com/apache/datafusion-comet/blob/bc4be39964cbe9cdb5f2a949740a8164e6b5755b/native/core/src/execution/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:
   
   ```scala
   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.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to