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


##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/over/AbstractNonTimeUnboundedPrecedingOver.java:
##########
@@ -418,10 +445,24 @@ Tuple2<Integer, Boolean> findIndexOfSortKey(
      * @throws Exception
      */
     RowData setAccumulatorAndGetValue(RowData accumulator) throws Exception {
-        aggFuncs.setAccumulators(accumulator);
+        setAccumulators(accumulator);
         return aggFuncs.getValue();
     }
 
+    /**
+     * Hands an accumulator to the aggregate functions, copying it first if it 
can hold a data view.
+     *
+     * <p>Accumulators are kept per sort key in {@link #accMapState}. 
Accumulating into one that was
+     * read from state would also change the entry that state still holds for 
an earlier sort key,
+     * because the heap state backend returns the stored object rather than a 
copy.
+     *
+     * @param accumulator the accumulator to continue from
+     */
+    void setAccumulators(RowData accumulator) throws Exception {

Review Comment:
   Confirmed, heap gave 2 for ord 20. accMapState held the live distinct view, 
the reset at the end of processElement cleared it, and the count field kept its 
1, so the second row counted a again. Added src_append and 
testCountDistinctAppendsAfterMaximum. Fixed by the get and put copy on 29410.
   



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