Skip to content

fix: [branch-1.1] make adaptive aggregation skipping opt-in (#6474) - #6488

Merged
andygrove merged 1 commit into
apache:branch-1.1from
andygrove:backport-6474-branch-1.1
Oct 2, 2026
Merged

andygrove merged 1 commit into
apache:branch-1.1from
andygrove:backport-6474-branch-1.1

Conversation

@andygrove

Copy link
Copy Markdown
Member

Which issue does this PR close?

Refs #6466 and #6402. Backport of merged #6474.

Rationale for this change

Disable adaptive aggregation skipping by default for RC2; retain explicit opt-in for eligible aggregates.

What changes are included in this PR?

Clean cherry-pick of c19ebbd onto branch-1.1, including configuration, native eligibility checks, documentation and regression tests.

How are these changes tested?

Fresh native build passed; native skip_partial_eligibility_is_fail_closed test passed. JDK 17/default Spark 4.1: three CometAggregateSuite "skip partial" tests and two CometExecSuite "SQLConf serde" tests passed. Full release CI pending. git diff --check passed.

…ggregate.skipPartial.enabled

DataFusion decides once per task, after the first 100,000 input rows, whether
to stop partial aggregation, and never checks again, so a task whose keys
repeat after a mostly distinct start can shuffle many times more rows than in
1.0.0. Put the bypass behind a new config that defaults to false, which
restores the 1.0.0 behavior.

Closes apache#6466.

(cherry picked from commit c19ebbd)
@github-actions github-actions Bot added the bug Something isn't working label Oct 1, 2026
@andygrove
andygrove marked this pull request as ready for review October 1, 2026 23:28

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summary

  • Prior state and problem: Adaptive partial aggregation could stop aggregating after a mostly distinct prefix and substantially increase shuffle volume when later keys repeated.
  • Design approach: Add spark.comet.exec.aggregate.skipPartial.enabled, defaulting to false, while retaining explicit opt-in for eligible native-shuffle plans.
  • Correctness / compatibility analysis: The setting crosses JNI with its resolved default and normalized boolean value. Native enforcement runs after DataFusion overrides. Existing restrictions on accumulators, PartialMerge, mixed modes and non-native-shuffle plans remain intact. Compared Spark’s Count and aggregation planning sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No introduced compatibility problem found.
  • Key design decisions: Reusing the existing session configuration gate keeps the change small. Eligibility is checked during session setup, and DataFusion 55.1 disables its aggregation probe when the ratio threshold is 1.1.
  • Implementation sketch: Add the Scala configuration, resolved serialization entry and native key, extend the existing gate, and update regression tests and tuning documentation.
  • Behavioral changes worth calling out: Ordinary aggregation is now the default. Opting in preserves the previous bypass behavior and its documented aggregation-work versus shuffle-volume tradeoff.
  • Suggested improvements: None meeting the reporting bar. No introduced P1/P2 issues found within this review.

Reviewed the full seven-file diff at bbb6806d6de4421a44ffe0c42f511d15da1407d7 against 66800830ea30cc75fddeed1b24e50d6d7ec61c57. The PR is not a draft. Snapshot and live discussion checks contained no reviews, issue comments, inline comments or review threads.

Routed skills: review-comet-pr, review-comet-expression-pr, review-comet-ffi-pr and review-comet-shuffle-pr.

Exact-head CI: run 36801073318 completed with 104 successful checks and seven skipped. Its merge commit has the same source tree as the reviewed head. Spark SQL passed on 3.5, 4.0 and 4.1, and Iceberg checks passed. Inspected artifacts confirm the native eligibility test, all three “skip partial” tests and both SQLConf serialization tests passed.

Validation limits: No fresh local native/JVM build, test rerun or benchmark was performed. Validation used source inspection, exact-head CI logs/reports and a passing local git diff --check. The upstream Spark 3.4 SQL suite was skipped.

@andygrove
andygrove merged commit 49bdf63 into apache:branch-1.1 Oct 2, 2026
111 checks passed
@andygrove

Copy link
Copy Markdown
Member Author

Thanks @sunchao

andygrove added a commit to andygrove/datafusion-comet that referenced this pull request Oct 2, 2026
…ease notes

Their backports merged into branch-1.1, so 1.1.0 won't ship them: apache#6423
(apache#6489), apache#6426 (apache#6486), apache#5507 (apache#6487), apache#6505 (apache#6511), apache#6334 (apache#6491) and
apache#6466 (apache#6488). The intro now says that apache#6402 also lists the regressions
fixed before the release.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants