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]