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]