SEPURI-SAI-KRISHNA commented on code in PR #29412:
URL: https://github.com/apache/flink/pull/29412#discussion_r4217636989
##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecOverAggregate.java:
##########
@@ -321,7 +321,11 @@ private KeyedProcessFunction<RowData, RowData, RowData>
createUnboundedOverProce
JavaScalaConversionUtil.toScala(aggCalls),
new boolean[aggCalls.size()],
false, // needInputCount
- true, // isStateBackendDataViews
+ // The non-time functions keep one accumulator per
sort key in state.
+ // State backed data views are bound to the key, not
to the sort key, so
+ // all of those accumulators would share a single
view. Keep the views in
+ // the accumulator instead, so each sort key gets its
own copy.
+ timeAttribute != TimeAttribute.NON_TIME, //
isStateBackendDataViews
Review Comment:
Added version 2 with minPlanVersion and minStateVersion 2.4, plus
over-aggregate-non-time-range-unbounded-collect and -count-distinct, both
ignored for version 1 where they would be wrong.
11 existing programs have no version 2 files and are ignored there. Four
stopped compiling when FLINK-38927 made a sink primary key that differs from
the upsert key an error. The other seven are the rowtime ones in
OverAggregateRestoreTest, where a fresh run emits rows their
consumedBeforeRestore does not list, so generation never reaches the savepoint.
Same seven twice in a row, and five of them use the bounded operator this PR
does not touch, so I think it predates this change.
Reuse their version 1 savepoints for version 2, or file the rowtime one
separately?
--
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]