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


##########
flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/over/NonTimeRangeUnboundedPrecedingFunction.java:
##########
@@ -287,7 +289,7 @@ void processRemainingElements(
             // Get previous accumulator
             RowData prevAcc = accMapState.get(curSortKey);
 
-            if (prevAcc.equals(currAcc)) {
+            if (accEqualiser.equals(prevAcc, currAcc)) {

Review Comment:
   Confirmed. src_repeated inserts 25 between 20 and 30, and on heap that gives 
ARRAY_AGG [r, p, q] for 20, [r, p, q, p, r] for 25 and [r, p, q, p, r, p] for 
30, against [r, q], [r, q, r] and [r, q, r, p] on RocksDB. LAG gives 20=p where 
RocksDB gives 20=r. Added testArrayAggInsertsInBetween and 
testLagInsertsInBetween.
   



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