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

Reply via email to