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]

Reply via email to