sunchao opened a new pull request, #6432:
URL: https://github.com/apache/datafusion-comet/pull/6432

   ## Which issue does this PR close?
   
   No linked issue.
   
   ## Rationale for this change
   
   Partial aggregation saves shuffle work by combining rows with the same key 
before sending them to the next stage. That saving depends on how often keys 
repeat within each input partition. For example:
   
   ```sql
   SELECT customer_id, MAX(bytes) AS largest_event
   FROM events
   GROUP BY customer_id;
   ```
   
   If one partition contains ten million events for 100 customers, local 
aggregation can reduce it to about 100 intermediate results. If those events 
belong to nine million customers, it still builds a large hash table but 
removes only about a tenth of the rows. These numbers illustrate the tradeoff; 
they are not benchmark results.
   
   DataFusion already detects this situation: it starts aggregating, measures 
the reduction, and can stop grouping locally when too many input rows become 
separate groups. Comet currently permits that path mainly for grouping-only 
queries and single-argument `COUNT`. An unsupported partial aggregate also 
disables it for the other aggregates in the same native execution block. 
Consequently, an eligible `MAX` such as the example above cannot benefit.
   
   This PR extends that existing adaptive behavior to more aggregates and makes 
eligibility independent for each aggregate operator. Final aggregation still 
combines the intermediate results into the answer. The potential saving is less 
unproductive local grouping; the tradeoff is more intermediate rows for shuffle 
and final aggregation.
   
   ## What changes are included in this PR?
   
   An eligible operator starts with ordinary partial aggregation. If 
DataFusion's reduction probe decides that grouping is no longer useful, it 
emits the groups already accumulated, then converts subsequent rows into the 
intermediate states its downstream merge expects.
   
   ```mermaid
   flowchart TD
       A[Input rows] --> B[Partial aggregation measures reduction]
       B --> C{Does local grouping reduce enough rows?}
       C -->|Yes| D[Continue local grouping]
       C -->|No| E[Emit accumulated states, then one state per later row]
       D --> F[Shuffle by grouping key]
       E --> F
       F --> G[Final aggregation merges states]
   ```
   
   For `MAX`, merging states `10`, `30`, and `20` gives the same answer as 
merging a locally combined state `30`. More complex aggregates need several 
fields: an `AVG` state contains a sum and count. The converter preserves those 
state formats by using an aggregate's existing conversion method or 
constructing its ordinary one-row state. Decimal128 `SUM` and `AVG` have direct 
converters, and `PartialMerge` can forward the states it already receives.
   
   Each aggregate operator now gets its own qualification decision. An eligible 
child can bypass even when its parent cannot. Spark also identifies which 
physical operators may emit repeated states: post-shuffle DISTINCT 
deduplication stages must still produce unique groups and remain ineligible. 
Global aggregates, unsupported boundaries, and order-sensitive or otherwise 
unqualified functions retain ordinary aggregation.
   
   The default policy supports grouping-only operators, `COUNT` with one or 
more arguments, `MIN`, `MAX`, bitwise aggregates, exact `percentile`, 
`collect_set`, and legacy integer `SUM`.
   
   The numerical policy is explicit. 
`spark.comet.exec.aggregate.partialBypass.enabled` defaults to `true`; 
`spark.comet.exec.aggregate.partialBypass.allowNumericalDifferences` defaults 
to `false`. The second setting is required for floating-point `SUM`, every 
`AVG`, statistical aggregates, decimal `SUM`, and ANSI/TRY integer `SUM`. 
Moving arithmetic between partial and final aggregation can change rounding or 
intermediate overflow, including whether a result is null or an error is 
raised. The opt-in accepts those numerical differences while retaining the 
structural and DISTINCT restrictions.
   
   Short intermediate-state batches are combined before shuffle when memory 
permits. The buffer retains at most 8 MiB of input array-size estimates and 
reserves space for both those inputs and concatenation output. It flushes when 
it cannot admit another batch, and full, oversized or refused individual 
batches pass through without copying. One already-read input may remain pending 
during a flush, so this bounds the accumulated fragments rather than all 
pipeline memory. Metrics distinguish eligibility from actual bypass and report 
grouped input, bypassed rows and reduction.
   
   ## How are these changes tested?
   
   Validated head `198d33d98b7e87684611d2facb80f05b277b1eca` against the public 
dependencies from base `7c9129b540331ec6da303ea95f83705bdafa432e` (DataFusion 
55.1.0, Arrow 59.3.0). No dependency versions change.
   
   - Native execution tests: **308 passed**, with 4 existing HDFS-dependent 
tests ignored.
   - Native aggregate-function tests: **172 passed**.
   - Full `CometAggregateSuite` on Spark 4.1.3 / Scala 2.13.17 / JDK 21: **131 
successful test executions**, including a repeat of the three new bypass 
regressions; 2 existing tests ignored.
   - Workspace Clippy with warnings denied, native build, Rust formatting, 
Maven Spotless/Scalastyle, and edited Markdown formatting passed.
   
   The Spark tests used the native library built from this PR's native source 
tree; its SHA-256 matched the library staged in the JVM resources. Native 
validation preceded the final one-line Scala test fixture adjustment, which 
leaves that native source tree unchanged.
   
   Native regressions cover operator-local policy and context restoration, 
conversion of scalar-only and multi-field aggregate states, nullable and 
filtered decimal inputs, state/schema preservation, numerical-policy gates, and 
the transition from accumulated groups to bypassed rows. Buffering tests 
exercise asynchronous polls, memory refusal, early stream drop, large list 
states and output order.
   
   Spark/JNI regressions compare query results with Spark and require positive 
bypass metrics for eligible cases. They cover AQE on and off, mixed 
Partial/PartialMerge producers, DISTINCT deduplication, filters, numerical 
opt-in, the disable switch, and JVM shuffle boundaries.
   
   No whole-query performance result is claimed by this contribution.
   


-- 
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