This is an automated email from the ASF dual-hosted git repository. pnowojski pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit d42a1ed5daf4b75caf1f574049584a054dc5dd6d Author: David Christle <[email protected]> AuthorDate: Sun Jan 28 00:04:09 2024 -0800 [FLINK-34252][table] Fix lastRecordTime tracking in WatermarkAssignerOperator to prevent erroneous idle WatermarkStatus --- .../wmassigners/WatermarkAssignerOperator.java | 9 +++++--- .../wmassigners/WatermarkAssignerOperatorTest.java | 25 ++++++++++++++++++++++ .../WatermarkAssignerOperatorTestBase.java | 8 +++++++ 3 files changed, 39 insertions(+), 3 deletions(-) diff --git a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperator.java b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperator.java index 55b7ca4c7fe..9a7a5067973 100644 --- a/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperator.java +++ b/flink-table/flink-table-runtime/src/main/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperator.java @@ -100,11 +100,14 @@ public class WatermarkAssignerOperator extends AbstractStreamOperator<RowData> @Override public void processElement(StreamRecord<RowData> element) throws Exception { - if (idleTimeout > 0 && currentStatus.equals(WatermarkStatus.IDLE)) { - // mark the channel active - emitWatermarkStatus(WatermarkStatus.ACTIVE); + if (idleTimeout > 0) { + if (currentStatus.equals(WatermarkStatus.IDLE)) { + // mark the channel active + emitWatermarkStatus(WatermarkStatus.ACTIVE); + } lastRecordTime = getProcessingTimeService().getCurrentProcessingTime(); } + RowData row = element.getValue(); if (row.isNullAt(rowtimeFieldIndex)) { throw new RuntimeException( diff --git a/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperatorTest.java b/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperatorTest.java index 450840201b8..e5deacb3999 100644 --- a/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperatorTest.java +++ b/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperatorTest.java @@ -84,6 +84,31 @@ public class WatermarkAssignerOperatorTest extends WatermarkAssignerOperatorTest assertThat(filterOutRecords(output)).isEqualTo(expectedOutput); } + @Test + public void testIdleStateAvoidanceWithConsistentDataFlow() throws Exception { + OneInputStreamOperatorTestHarness<RowData, RowData> testHarness = + createTestHarness(0, WATERMARK_GENERATOR, 1000); + testHarness.getExecutionConfig().setAutoWatermarkInterval(50); + testHarness.open(); + + ConcurrentLinkedQueue<Object> output = testHarness.getOutput(); + List<Object> expectedOutput = new ArrayList<>(); + + // Process elements at intervals less than idleTimeout (1000ms) + for (long i = 1; i <= 10; i++) { + long timestamp = i * 900; + testHarness.processElement(new StreamRecord<>(GenericRowData.of(timestamp), timestamp)); + testHarness.setProcessingTime(timestamp); + + // Expect a watermark for each element, lagging by 1 + expectedOutput.add(new Watermark(timestamp - 1)); + } + + // Check if the status ever becomes IDLE (it shouldn't) + assertThat(extractWatermarkStatuses(output)).doesNotContain(WatermarkStatus.IDLE); + assertThat(filterOutRecords(output)).isEqualTo(expectedOutput); + } + @Test public void testWatermarkAssignerOperator() throws Exception { OneInputStreamOperatorTestHarness<RowData, RowData> testHarness = diff --git a/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperatorTestBase.java b/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperatorTestBase.java index 1aeba2e6c86..6f64d2c6b24 100644 --- a/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperatorTestBase.java +++ b/flink-table/flink-table-runtime/src/test/java/org/apache/flink/table/runtime/operators/wmassigners/WatermarkAssignerOperatorTestBase.java @@ -22,6 +22,7 @@ import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.watermark.Watermark; import org.apache.flink.streaming.runtime.streamrecord.StreamElement; import org.apache.flink.streaming.runtime.streamrecord.StreamRecord; +import org.apache.flink.streaming.runtime.watermarkstatus.WatermarkStatus; import org.apache.flink.table.data.RowData; import java.util.ArrayList; @@ -60,6 +61,13 @@ public abstract class WatermarkAssignerOperatorTestBase { return watermarks; } + protected List<WatermarkStatus> extractWatermarkStatuses(Collection<Object> collection) { + return collection.stream() + .filter(obj -> obj instanceof WatermarkStatus) + .map(obj -> (WatermarkStatus) obj) + .collect(Collectors.toList()); + } + protected List<Object> filterOutRecords(Collection<Object> collection) { return collection.stream() .filter(obj -> !(obj instanceof StreamElement && ((StreamElement) obj).isRecord()))
