fix: make adaptive partial aggregation opt-in with spark.comet.exec.aggregate.skipPartial.enabled - #6474
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.
|
@peterxcli @sunchao could you review? |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Adaptive partial aggregation could permanently bypass aggregation after a mostly distinct prefix, increasing shuffle volume when keys subsequently repeat.
- Design approach: Introduce
spark.comet.exec.aggregate.skipPartial.enabled, defaulting tofalse, while retaining the existing eligibility checks. - Correctness / compatibility analysis: The resolved boolean crosses JNI correctly. Applying the gate after DataFusion overrides prevents those overrides from enabling skipping when disabled. Spark count, partial/final aggregation and deduplication semantics remain intact.
- Key design decisions: Reusing the existing session configuration gate keeps the change small. Both DataFusion aggregation implementations omit the adaptive probe when the ratio threshold is
1.1. - Implementation sketch: Declare and serialize the flag, read it during native session creation, extend eligibility and serialization tests, and update tuning documentation. The added configuration work occurs during initialization, without introducing per-row JNI calls or copies.
- Behavioral changes worth calling out: Ordinary partial aggregation becomes the default. Opting in preserves the previous adaptive behavior, including its documented irreversible bypass decision.
- Suggested improvements: No introduced P1/P2 issues found within this review.
Reviewed the entire seven-file diff from b38a60b37770df6efb4c66faac751160e92990bf to c19ebbdc209de59b00c158917bfa0e27214575c8. The PR is not a draft. Read the supplied discussion snapshot and refreshed existing reviews and inline comments. No unresolved substantive concerns were present. Routed skills: review-comet-pr, review-comet-ffi-pr, and review-comet-shuffle-pr.
Exact-head CI: 21 checks succeeded, 14 were skipped, and three remained running: Spark 4.1 execution tests, expression tests, and TPC-DS verification. No failures were reported. Native CI passed, including skip_partial_eligibility_is_fail_closed. Spark 4.1 scan and shuffle suites also passed. Spark’s own SQL suites, macOS, and Iceberg suites were skipped.
Validation: Compared relevant Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0, and checked DataFusion 55.1.0 probe behavior. git diff --check passed. Local integration tests were not run: the checkout lacks build artifacts and Spark dependencies, and the filesystem exhausted available space during review. Consequently, the new JVM regression tests and broader SQL compatibility checks lack a completed verdict in this review.
|
Thanks @sunchao |
Which issue does this PR close?
Closes #6466.
Found by the 1.1.0 regression audit (#6399) and tracked in #6402.
Rationale for this change
#5723 turned on DataFusion's adaptive partial aggregation for native shuffle plans whose partial aggregates only group, or only compute single-argument
COUNT. DataFusion starts checking after the first 100,000 input rows of a task. As soon as more than 80% of the rows it has seen are distinct groups, it stops aggregating and sends every later row to the shuffle as it is, and it never checks again. So a task whose keys repeat after a mostly distinct start, such as several snapshot files of the same keys packed into one split, can shuffle many times more rows than in 1.0.0, which never skipped. The reproducer in #6466 shuffles 2,000,000 rows on 1.1.0-rc1 and 200,000 on 1.0.0.The only way to turn it off was the testing-only
spark.comet.exec.respectDataFusionConfigs=truetogether with a DataFusion ratio threshold of 1.1, and the review of #5723 asked for a supported switch. Rather than tune the heuristic now, this puts the optimization behind a config that is off by default, which restores the 1.0.0 behavior. Probing again when later input stops being distinct needs a change in DataFusion, and the default can be revisited after that.This is meant for 1.1.0-rc2, through a backport to
branch-1.1once it merges.What changes are included in this PR?
spark.comet.exec.aggregate.skipPartial.enabled, in theexeccategory andfalseby default. It joins the configs thatCometExecIterator.serializeCometSQLConfssends to native code resolved, so its default crosses JNI.configure_skip_partial_aggregationinjni_api.rstakes the flag and keeps skipping off unless it is set, on top of the existing eligibility checks. It still runs after thespark.comet.datafusion.*pass-through, so the DataFusion thresholds can tune an eligible plan but can't turn skipping on by themselves.The plan, the protobuf and the metrics don't change. With the config on, the behavior is the same as on
maintoday.How are these changes tested?
CometAggregateSuitetest,skip partial aggregation is disabled by default, is the reproducer from Adaptive partial aggregation can shuffle 10x more rows than 1.0.0, with no supported way to turn it off #6466. One task reads 2,000,000 rows whose 200,000 keys cycle 10 times, and the test asserts that the shuffle writes 200,000 records. With skipping on, it writes exactly 2,000,000, the rc1 number. It runs in under 2 seconds.skip partial aggregation admits only supported native shuffle plansalso checks that with the config off, an eligible plan doesn't skip even when a DataFusion ratio threshold of 0.8 is passed through.skip_partial_eligibility_is_fail_closedtest checks that a disabled flag forces the ratio threshold to 1.1 for eligible plans.CometExecSuite'sSQLConf serde resolves the configs that native code parsescovers the new key. It is the only test that catches the key missing from the resolved list, because an explicitly setspark.comet.*value crosses JNI anyway.I checked the tests against three broken versions of the fix: the default flipped to
true, the key dropped from the resolved list, and the native gate removed. Each one fails at least one of the tests. The fullCometAggregateSuitepasses locally on the default Spark 4.1 profile.