SEPURI-SAI-KRISHNA commented on code in PR #29412:
URL: https://github.com/apache/flink/pull/29412#discussion_r4217586318
##########
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, logDataViewOnNonTimeOver. It fires only for a non-time window on a
version 1 plan, and names the aggregates that keep a data view, including the
distinct ones, which sit in distinctInfos rather than in viewSpecs.
--
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]