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

   ### Describe the bug
   
   A native grouped aggregate whose groups carry large state, such as 
`collect_list` over a low-cardinality key, fails the task once it spills, while 
Spark completes the same query. The error is `CometNativeException: Additional 
allocation failed for FinalHashAggregateStream`. Spilling can't help: when it 
fails, the aggregate holds nothing, and a single request is larger than its 
`fair_unified` share. The same data spread over many small groups spills and 
succeeds with the same off-heap size.
   
   The mechanism is in DataFusion. Spilled aggregate state is written in 
batches of up to `batch_size` rows whatever their size in bytes. With few 
groups, one batch holds the whole table, and reading it back needs more memory 
than the budget that forced the spill. I filed that upstream as 
apache/datafusion#25851, with a plain SQL repro on DataFusion main. On 55.1, 
which Comet uses, the failing request can also be the spill merge's read-back 
reservation. That is twice the largest spilled batch, and 55.1 returns the 
error when even one spill file doesn't fit.
   
   Two things on the Comet side make it harder to live with:
   
   1. The obvious mitigation, a smaller `spark.comet.batchSize`, would cap the 
rows in each spilled batch. But below 8192 it stops `CometConf` from 
initializing on executors (#6286).
   2. The error doesn't say that spilling cannot satisfy the request, or which 
settings could. In the repro below, the share is a third of the pool because 
two non-spillable global aggregates are registered alongside the grouped one. 
The request is larger than the whole pool, so #5465 (splitting only among 
spillable consumers) would not rescue this case.
   
   ### Steps to reproduce
   
   On main at 65b334bfd with the default Spark 4.1.3 profile, `local[1]`, 
`spark.memory.offHeap.enabled=true`, `spark.memory.offHeap.size=128m`, Comet 
exec and shuffle enabled, `spark.sql.adaptive.enabled=false` and 
`spark.sql.shuffle.partitions=1`:
   
   ```scala
   import org.apache.spark.sql.functions._
   
   def query(groups: Long) =
     spark.range(0, 3000000L, 1, 4)
       .select(
         (col("id") % groups).as("k"),
         concat(col("id").cast("string"), lit("x" * 100)).as("v"))
       .groupBy("k")
       .agg(size(collect_list("v")).as("n"))
       .agg(count(lit(1)).as("groups"), sum("n").as("values_collected"))
   
   query(16).collect()
   ```
   
   It fails with:
   
   ```
   org.apache.comet.CometNativeException: Additional allocation failed for 
FinalHashAggregateStream[0] with top memory consumers (across reservations) as:
     AggregateStream[0]#5(can spill: false) consumed 0.0 B, peak 0.0 B,
     FinalHashAggregateStream[0]#4(can spill: true) consumed 0.0 B, peak 24.0 B,
     AggregateStream[0]#6(can spill: false) consumed 0.0 B, peak 0.0 B.
   Error: Failed to acquire 164778180 bytes where this consumer already holds 0 
bytes and the fair limit is 44739242 bytes, 3 registered (0 bytes overcommitted)
   ```
   
   Controls in the same session, with the same 3,000,000 values:
   
   - `query(500000)` succeeds. The final aggregate reports 11 spills and 548 MB 
spilled.
   - `query(16)` with `spark.comet.enabled=false` succeeds.
   - `query(16)` with `spark.memory.offHeap.size=2g` succeeds, and the final 
aggregate does not spill.
   
   I ran these as a throwaway `CometTestBase` suite with the overrides above.
   
   In this repro, the final aggregate's peak is 24 bytes: the partial 
aggregates hand it batches that are already larger than its share, so it spills 
from the first one. The failure also happens after an aggregate has filled up 
to its share over many batches. There, the request is the merge's read-back 
reservation instead, twice the largest spilled batch.
   
   ### Expected behavior
   
   The query completes under the memory limit, as it does in Spark and as the 
same data does when spread over many groups. Until the upstream fix lands, the 
error should say that the request is larger than the operator's share and 
cannot be satisfied by spilling. It should also name the settings that can 
help, and the tuning guide should describe this case.
   
   ### Additional context
   
   Workarounds today:
   
   - Give the aggregate enough off-heap memory that it doesn't spill.
   - Run the aggregate function in Spark with 
`spark.comet.expression.CollectList.enabled=false` (or the matching flag for 
another function), or disable native aggregation with 
`spark.comet.exec.aggregate.enabled=false`.
   
   A smaller batch size is blocked by #6286, as noted above.
   


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