andygrove opened a new issue, #6466:
URL: https://github.com/apache/datafusion-comet/issues/6466
### Describe the bug
Since #5723, Comet turns on DataFusion's adaptive partial aggregation for
native shuffle plans whose partial aggregates only group, or only count one
column. DataFusion decides once per task, after the first 100,000 input rows.
If more than 80% of those rows were distinct groups, it stops aggregating and
sends every later row straight to the shuffle, and it never checks again. So a
task whose keys repeat after a mostly distinct start shuffles far more rows
than in 1.0.0, which never skipped. Several snapshot files of the same keys
packed into one split would do it.
There's no supported way to turn this off. The only route is the
testing-only `spark.comet.exec.respectDataFusionConfigs=true` together with
`spark.comet.datafusion.execution.skip_partial_aggregation_probe_ratio_threshold=1.1`,
which the tuning guide documents. The review of #5723 asked for a supported
switch.
### Steps to reproduce
```scala
withSQLConf(
SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
SQLConf.FILES_MAX_PARTITION_BYTES.key -> (1L << 30).toString,
SQLConf.SHUFFLE_PARTITIONS.key -> "4") {
withTempPath { dir =>
val path = dir.getCanonicalPath
// One file and one task. Keys 0 to 199999 cycle 10 times, so the first
100k rows are all
// distinct, but every key repeats 10 times within the task.
spark.range(0, 2000000, 1, 1).selectExpr("id % 200000 AS
k").write.parquet(path)
val df = spark.read.parquet(path).groupBy("k").count()
df.collect()
val written = collect(df.queryExecution.executedPlan) {
case e: CometShuffleExchangeExec =>
e.metrics("shuffleRecordsWritten").value
}.sum
assert(written == 200000, s"the partial aggregate shuffled $written
rows")
}
}
```
On the Spark 4.1 profile, 1.0.0 shuffles 200,000 rows and 1.1.0-rc1 shuffles
2,000,000. On rc1 the partial aggregate reports `skipped_aggregation_rows` =
1,893,504, and the shuffle writes 8.2 MB instead of 0.8 MB. The results are
correct on both.
### Expected behavior
A task like this shouldn't shuffle 10 times more rows than 1.0.0 did, and
there should be a supported config to turn the bypass off, such as
`spark.comet.exec.aggregate.skipPartial.enabled`. Having the probe look again
when later input stops being distinct would also help, but that needs a change
in DataFusion.
### Additional context
Found by the 1.1.0 regression audit (#6399) and tracked in #6402.
--
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]