fix: [branch-1.1] make adaptive aggregation skipping opt-in (#6474) - #6488
Conversation
…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)
sunchao
left a comment
There was a problem hiding this comment.
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 tofalse, 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’sCountand 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.
|
Thanks @sunchao |
…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.
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 --checkpassed.