SEPURI-SAI-KRISHNA commented on code in PR #29412:
URL: https://github.com/apache/flink/pull/29412#discussion_r4230574624


##########
flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/stream/StreamExecOverAggregate.java:
##########
@@ -572,4 +603,36 @@ private KeyedProcessFunction<RowData, RowData, RowData> 
createBoundedOverProcess
                 throw new TableException("Unsupported bounded operation for 
OVER window.");
         }
     }
+
+    /**
+     * Warns that a compiled plan from before the data views of a non-time 
OVER window moved into
+     * the accumulator keeps returning wrong results.
+     *
+     * <p>Such a plan stays on version 1 so that its state still restores. The 
results of an
+     * aggregate with a data view do not become correct by restoring it, so 
the plan has to be
+     * recompiled.
+     */
+    private void logDataViewOnNonTimeOver(AggregateInfoList aggInfoList) {
+        final Stream<String> withView =
+                Arrays.stream(aggInfoList.aggInfos())
+                        .filter(aggInfo -> aggInfo.viewSpecs().length > 0)
+                        .map(aggInfo -> 
aggInfo.agg().getAggregation().getName());
+        // A distinct aggregate keeps its view next to the aggregates, not in 
its own view specs
+        final Stream<String> distinct =
+                Arrays.stream(aggInfoList.distinctInfos())
+                                .anyMatch(distinctInfo -> 
distinctInfo.dataViewSpec().isDefined())
+                        ? Stream.of("DISTINCT")
+                        : Stream.empty();
+        final String aggregates =
+                Stream.concat(withView, 
distinct).distinct().collect(Collectors.joining(", "));
+        if (!aggregates.isEmpty()) {
+            LOG.warn(
+                    "The OVER window ordered by a non-time attribute 
aggregates with {}, which "
+                            + "keeps data in a data view. This compiled plan 
is on version 1 of "
+                            + "'stream-exec-over-aggregate', where those 
aggregates return wrong "
+                            + "results. Recompile the plan to get correct 
results. The job then "
+                            + "starts without the state of this operator.",

Review Comment:
   Right, generateUid leaves the version out on purpose. Fixed in the warning, 
the version 2 comment, the javadoc, the description and the release note.
   



-- 
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