andygrove opened a new pull request, #6474: URL: https://github.com/apache/datafusion-comet/pull/6474
## 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=true` together 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.1` once it merges. ## What changes are included in this PR? - A new config, `spark.comet.exec.aggregate.skipPartial.enabled`, in the `exec` category and `false` by default. It joins the configs that `CometExecIterator.serializeCometSQLConfs` sends to native code resolved, so its default crosses JNI. - `configure_skip_partial_aggregation` in `jni_api.rs` takes the flag and keeps skipping off unless it is set, on top of the existing eligibility checks. It still runs after the `spark.comet.datafusion.*` pass-through, so the DataFusion thresholds can tune an eligible plan but can't turn skipping on by themselves. - The Adaptive Partial Aggregation section of the operator tuning guide now describes the feature as opt-in, explains how the one-way decision can shuffle more rows, and drops the ratio-1.1 workaround. The plan, the protobuf and the metrics don't change. With the config on, the behavior is the same as on `main` today. ## How are these changes tested? - A new `CometAggregateSuite` test, `skip partial aggregation is disabled by default`, is the reproducer from #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. - The two existing skip-partial tests now turn the config on. `skip partial aggregation admits only supported native shuffle plans` also checks that with the config off, an eligible plan doesn't skip even when a DataFusion ratio threshold of 0.8 is passed through. - The native `skip_partial_eligibility_is_fail_closed` test checks that a disabled flag forces the ratio threshold to 1.1 for eligible plans. - `CometExecSuite`'s `SQLConf serde resolves the configs that native code parses` covers the new key. It is the only test that catches the key missing from the resolved list, because an explicitly set `spark.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 full `CometAggregateSuite` passes locally on the default Spark 4.1 profile. -- 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]
