SEPURI-SAI-KRISHNA opened a new pull request, #29410:
URL: https://github.com/apache/flink/pull/29410
## What is the purpose of the change
A non-time `OVER` window compares the previous and current accumulator with
`RowData#equals` to
decide whether it can skip emitting updates. That comparison is wrong in two
ways.
On the heap backend it throws for any aggregate whose accumulator holds a
RAW field:
```
java.lang.UnsupportedOperationException: Unmaterialized BinaryRawValueData
cannot be compared.
at
org.apache.flink.table.data.binary.BinaryRawValueData.equals(BinaryRawValueData.java:85)
at java.base/java.util.Arrays.deepEquals(Arrays.java:4625)
at
org.apache.flink.table.data.GenericRowData.equals(GenericRowData.java:229)
at
...NonTimeRangeUnboundedPrecedingFunction.processRemainingElements(...:290)
```
On RocksDB it does not throw but never matches, because the stored
accumulator is deserialized
as `BinaryRowData` while the current one is a `GenericRowData`. The
early-out never fires, so
every following row is recomputed and emits an
`UPDATE_BEFORE`/`UPDATE_AFTER` pair carrying
identical values.
A generated `RecordEqualiser` over the accumulator types fixes both: it
materializes RAW fields
before comparing them, and it reads fields positionally so the two row
classes may differ.
## Brief change log
- `StreamExecOverAggregate` generates a `RecordEqualiser` for the flattened
accumulator types and
passes it to both non-time functions.
- `AbstractNonTimeUnboundedPrecedingOver` and
`NonTimeRangeUnboundedPrecedingFunction` use it for
the early-out comparison.
- `NonTimeOverAggregateITCase` covers `LAG`, `ARRAY_AGG`, `COLLECT` and
`PERCENTILE` on a window
ordered by a non-time attribute, under both state backends.
The two remaining comparisons in `NonTimeRowsUnboundedPrecedingFunction`
compare an aggregate
value against an accumulator, which FLINK-40883 is correcting. They are left
alone here; the
equaliser is threaded through its constructor so it is available once that
lands. The two changes
touch different lines and can land in either order.
`COLLECT` and `PERCENTILE` values are not asserted, only that the queries
run. They return the
current row rather than the running window, which is FLINK-40737.
## Verifying this change
Red/green on the new ITCase, which runs 8 cases (4 queries x 2 state
backends). With the two
comparison changes reverted, 4 of the 8 fail: every heap case, each with the
exception above. The
4 RocksDB cases pass, emitting the redundant pairs described. All 8 pass
with the change.
Existing coverage: `NonTimeRangeUnboundedPrecedingFunctionTest` and
`NonTimeRowsUnboundedPrecedingFunctionTest` 24/24; the
`OverAggregate`/`OverWindow` planner, ITCase,
harness and restore suites 386 tests, 0 failures, including
`OverWindowRestoreTest` and
`OverAggregateRestoreTest`.
## 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): (yes, the
early-out comparison in a
non-time OVER window is now a generated equaliser instead of
`RowData#equals`)
- Anything that affects deployment or recovery: (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 Code (Opus 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]