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]

Reply via email to