bvarghese1 commented on code in PR #29373:
URL: https://github.com/apache/flink/pull/29373#discussion_r4209060208
##########
flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/over/NonTimeRowsUnboundedPrecedingFunctionTest.java:
##########
@@ -96,6 +96,58 @@ void testInsertOnlyRecordsWithCustomSortKey() throws
Exception {
validateRows(actualRows, expectedRows);
}
+ @Test
+ void testInsertWithDuplicateSortKeyAndLastValueAgg() throws Exception {
+ KeyedProcessOperator<RowData, RowData, RowData> operator =
+ new KeyedProcessOperator<>(
+ new NonTimeRowsUnboundedPrecedingFunction<RowData>(
+ 0L,
+ lastValueAggsHandleFunction,
+ GENERATED_ROW_VALUE_EQUALISER,
+ GENERATED_SORT_KEY_EQUALISER,
+ GENERATED_SORT_KEY_COMPARATOR_ASC,
+ lastValueAccTypes,
+ inputFieldTypes,
+ SORT_KEY_TYPES,
+ SORT_KEY_SELECTOR) {});
+
+ OneInputStreamOperatorTestHarness<RowData, RowData> testHarness =
+ createTestHarness(operator);
+ testHarness.open();
+
+ testHarness.processElement(insertRecord("key1", 1L, 100L));
+ testHarness.processElement(insertRecord("key1", 2L, 200L));
+ testHarness.processElement(insertRecord("key1", 5L, 500L));
+ testHarness.processElement(insertRecord("key1", 6L, 600L));
+ testHarness.processElement(insertRecord("key1", 4L, 400L));
+ testHarness.processElement(insertRecord("key1", 5L, 503L));
+ testHarness.processElement(updateBeforeRecord("key1", 5L, 500L));
Review Comment:
Good catch, I was able to reproduce it. Fixed.
Parameterized the tests so that the tests use both the `heap` and
`EmbeddedRocksDBStateBackend` state backends.
Note: early out tests are skipped on RocksDB for now. It can be updated once
[FLINK-40735](https://issues.apache.org/jira/browse/FLINK-40735) is merged.
--
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]