From 68bc0b0a0d51b79789263e8949da1115624d7e1c Mon Sep 17 00:00:00 2001 From: developer-rpai Date: Fri, 9 Oct 2026 16:05:45 -0700 Subject: [PATCH] Flink: Forward watermarks in EqualityConvertCommitter when no cycle is active Fixes #18435. --- .../operator/EqualityConvertCommitter.java | 10 +++++++++- .../operator/TestEqualityConvertCommitter.java | 18 ++++++++++++++++++ 2 files changed, 27 insertions(+), 1 deletion(-) diff --git a/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/EqualityConvertCommitter.java b/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/EqualityConvertCommitter.java index ada16ec09384..5c080ac64c9a 100644 --- a/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/EqualityConvertCommitter.java +++ b/flink/v2.3/flink/src/main/java/org/apache/iceberg/flink/maintenance/operator/EqualityConvertCommitter.java @@ -148,7 +148,15 @@ public void processElement2(StreamRecord record) { @Override public void processWatermark(Watermark mark) throws Exception { - if (planResult == null || mark.getTimestamp() < planResult.doneTimestamp()) { + if (planResult == null) { + // No conversion cycle is active: forward the watermark. It may originate from another + // task's trigger sharing the maintenance lock, and holding it back would stall lock + // release for that task. + super.processWatermark(mark); + return; + } + + if (mark.getTimestamp() < planResult.doneTimestamp()) { // Hold back watermarks until the cycle commits so the LockRemover keeps the maintenance lock // for the whole cycle. Forwarding the planner's mid-cycle phase watermarks will release the // lock early and could let the next trigger run a concurrent cycle on the same staging diff --git a/flink/v2.3/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestEqualityConvertCommitter.java b/flink/v2.3/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestEqualityConvertCommitter.java index 53e98a06b478..1bca49a39fc8 100644 --- a/flink/v2.3/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestEqualityConvertCommitter.java +++ b/flink/v2.3/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestEqualityConvertCommitter.java @@ -133,6 +133,24 @@ void holdsBackWatermarkUntilCommit() throws Exception { } } + @Test + void forwardsWatermarkWithoutActiveCycle() throws Exception { + Table table = createTable(3, FileFormat.PARQUET); + insert(table, 1, "a"); + + try (TwoInputStreamOperatorTestHarness harness = + createHarness()) { + harness.open(); + + // No plan was received, so no conversion cycle is active. A watermark here originates + // from another task's trigger sharing the maintenance lock; dropping it would stall + // the unioned watermark and the lock would never be released. + long markTs = System.currentTimeMillis(); + harness.processBothWatermarks(new Watermark(markTs)); + assertThat(watermarks(harness)).containsExactly(new Watermark(markTs)); + } + } + @Test void skipsCommitForEmptyCycle() throws Exception { Table table = createTable(3, FileFormat.PARQUET);