SEPURI-SAI-KRISHNA opened a new pull request, #29412:
URL: https://github.com/apache/flink/pull/29412
## What is the purpose of the change
`COLLECT` over an `OVER` window ordered by a non-time attribute returns a
multiset holding only
the current row instead of the running window, with no error:
```sql
SELECT ord, COLLECT(v) OVER (PARTITION BY k ORDER BY ord),
ARRAY_AGG(v) OVER (PARTITION BY k ORDER BY ord)
FROM (VALUES ('a',10,'p'),('a',20,'q'),('a',30,'r')) AS t(k,ord,v);
```
```
ord COLLECT ARRAY_AGG
10 {p=1} [p]
20 {q=1} [p, q]
30 {r=1} [p, q, r]
```
`ARRAY_AGG` accumulates, `COLLECT` does not, on the same rows in the same
query.
The non-time functions keep one accumulator per sort key in `accMapState`,
but the planner asks
for state backed data views. A state backed view is bound to the key, not to
the sort key, so all
of those accumulators share a single view, and
`DataViewUtils.adjustDataViews` gives the view
field a `NullSerializer` so nothing of the accumulator is actually written
to state. On top of
that, `processElement` ends every record with `aggFuncs.cleanup()`, which is
generated as
`view.clear()`, so the one place the accumulator lived is wiped after each
row. The next row
starts from an empty view, which is why the result is just the current row.
This affects every aggregate whose accumulator holds a `MapView` or
`ListView`, so also
`PERCENTILE`, `FIRST_VALUE`/`LAST_VALUE` with retraction, `LISTAGG` with
retraction,
`JSON_OBJECTAGG`, and `DISTINCT` aggregates. `COUNT(DISTINCT v)` over the
same kind of window is
wrong on both state backends today. Aggregates whose accumulator holds plain
fields, such as
`SUM`, `COUNT` and `ARRAY_AGG`, are unaffected because they round-trip
through `accMapState`
normally.
The fix keeps the data views inside the accumulator for the non-time path,
so each sort key gets
its own copy and the per-record `cleanup()` has nothing to wipe.
## Brief change log
- `StreamExecOverAggregate` asks for non state backed data views for the
`NON_TIME` branch only.
The time based unbounded functions keep one accumulator per key, so they
are left as they were.
- `AbstractNonTimeUnboundedPrecedingOver` copies the accumulator before
handing it to the
aggregate functions, because the aggregate functions change a data view in
place and the heap
state backend returns the object it stores rather than a copy. The copy is
skipped when no
accumulator type contains a RAW field, which is the case for `SUM`,
`COUNT` and friends.
- `DistinctAggCodeGen` reads a heap distinct view through the external
converter instead of
`getJavaObject()`. `getJavaObject()` is null once the accumulator has been
read back from
RocksDB, where the raw value only holds bytes.
- `NonTimeOverAggregateITCase` asserts the `COLLECT` and `PERCENTILE` values
that were left
unasserted before, and adds a repeated value case for multiset counts and
`COUNT(DISTINCT)`.
## Verifying this change
`NonTimeOverAggregateITCase` runs on both the heap and the RocksDB backend.
Reverting the three
production changes and keeping the tests gives 12 cases run, 8 failures,
which is every case of
`testCollect`, `testCollectCountsRepeatedValues`, `testPercentile` and
`testCountDistinct` on
both backends. `testLag` and `testArrayAgg` keep passing, so the existing
behaviour is untouched.
With the change all 12 pass, and the heap and RocksDB changelogs are
identical.
Regression runs on Java 17: 1109 tests in the over window, aggregate and
distinct aggregate
suites, and 283 tests in the batch aggregate suites, 0 failures. That
includes
`OverWindowRestoreTest` and `OverAggregateRestoreTest`, whose 20 non-time
savepoints restore
unchanged: their programs use `SUM` and `AVG`, whose accumulators hold no
data view, so the
accumulator state schema is the same either way. Only aggregates with a data
view change schema,
and those cannot have a correct savepoint today.
## 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
accumulator is copied
once per sort key visit, and only for aggregates whose accumulator holds
a RAW field
- Anything that affects deployment or recovery: yes, see above, the
accumulator state schema
changes for aggregates that hold a data view
- 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]