Skip to content

Adaptive partial aggregation can shuffle 10x more rows than 1.0.0, with no supported way to turn it off #6466

Description

@andygrove

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

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.

Activity

  1. self-assigned this
    on Sep 30, 2026
  2. added a commit that references this issue on Sep 30, 2026
    961fbfe
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

area:aggregationHash aggregates, aggregate expressionsbugSomething isn't workingperformanceregressionA bug that did not affect the most recent Comet releaserequires-triage

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions