kotwal-itpro opened a new pull request, #29395:
URL: https://github.com/apache/flink/pull/29395
## What is the purpose of the change
`AVG(x)` on an INT column returns a wrong value when the planner rewrites it
into `SUM(x) / COUNT(x)`, e.g. when SUM and COUNT of the same column are
selected or when AVG is pushed below a join. The SUM keeps type INT and
overflows, while the AVG function accumulates in BIGINT.
Opened as a draft because the fix has a planning trade-off that I'd like a
maintainer to weigh in on (see below).
## Brief change log
- Add `FlinkAggregateReduceFunctionsRule`: Calcite's
`AggregateReduceFunctionsRule` configured with an extra condition that skips
AVG on TINYINT, SMALLINT and INT. Use it in the batch and stream rule sets and
in `WindowAggregateReduceFunctionsRule`.
## Trade-off and alternative
Not reducing these AVG calls means `AGGREGATE_VALUES` can no longer fold a
global aggregate over a statically empty input into a literal row when it
contains AVG on an INT column (`testGlobalAggOverEmptyInputReplacedByValues`
plans change), and AVG no longer shares a SUM(x) of the same query
(`SplitAggregateRuleTest#testSingleDistinctAggWithAllNonDistinctAgg`).
The alternative keeps both optimizations: in the reduction itself, sum
`CAST(x AS BIGINT)` and divide before casting back to the AVG type. That has to
change `reduceAvg` in the copied
`org.apache.calcite.rel.rules.AggregateReduceFunctionsRule`, which AGENTS.md
asks contributors not to modify, so I did not do that here. I have that version
ready if it is preferred.
Separately, `SplitAggregateRule` has the same overflow when
`table.optimizer.distinct-agg.split.enabled` is set (`SELECT COUNT(DISTINCT k),
AVG(x)` returns 0 for `x = MAX_INT, MAX_INT, 1`). I will file a separate issue
for it.
## Verifying this change
- Added
`AggregateITCaseBase#testAvgOnIntDoesNotOverflowWhenReducedToSumAndCount`
covering both triggers from the JIRA, for hash and sort aggregation. It fails
without the fix.
- Aggregate, window, join and expression-reduction suites pass (6,638
tests); the four plan changes listed above are the only differences.
## Does this pull request potentially affect one of the following parts:
- Dependencies (does it add or upgrade a dependency): no
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: no
- The serializers: no
- The runtime per-record code paths (performance sensitive): no
- Anything that affects deployment or recovery: JobManager (and its
components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no
- The S3 file system connector: no
## Documentation
- Does this pull request introduce a new feature? no
- If yes, how is the feature documented? not applicable
---
##### Was generative AI tooling used to co-author this PR?
- [X] Yes (please specify the tool below)
Generated-by: Claude Opus 5.5
--
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]