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]

Reply via email to