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]