From cd775eff8778a38e2e0b758a779fd163d1162eb2 Mon Sep 17 00:00:00 2001 From: "yongfu.gao" <2642474295@qq.com> Date: Sun, 27 Sep 2026 01:36:33 +0800 Subject: [PATCH 1/4] [FLINK-40504][core] Make PartialWatermark idleness tri-state An output's idleness is now UNKNOWN, ACTIVE or IDLE instead of a boolean that defaults to "not idle", and CombinedWatermarkStatus exposes hasActiveOutput() for callers that need to know whether an output actually reported being active. This is behaviour-neutral on its own: UNKNOWN combines exactly like the previous idle = false did, so it still holds the combined watermark back, and no caller uses hasActiveOutput() yet. Generated-by: DeepSeek Harness (deepseek-v4.1-flash) --- .../eventtime/CombinedWatermarkStatus.java | 41 +++++++++++++++++-- 1 file changed, 38 insertions(+), 3 deletions(-) diff --git a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/CombinedWatermarkStatus.java b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/CombinedWatermarkStatus.java index fb290ef957b1d..c94211305a11f 100644 --- a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/CombinedWatermarkStatus.java +++ b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/CombinedWatermarkStatus.java @@ -46,6 +46,20 @@ public boolean isIdle() { return idle; } + /** + * Returns true if any output has reported that it is active, which is what makes it safe to + * announce activity downstream. An output that has not said anything yet, or that only reported + * idleness, is not evidence that anything is producing. + */ + public boolean hasActiveOutput() { + for (PartialWatermark partialWatermark : partialWatermarks) { + if (partialWatermark.isActive()) { + return true; + } + } + return false; + } + public boolean remove(PartialWatermark o) { return partialWatermarks.remove(o); } @@ -108,8 +122,24 @@ public boolean updateCombinedWatermark() { /** Per-output watermark state. */ static class PartialWatermark { + + /** + * The idleness the output has reported, if any. + * + *

A newly registered output is {@link #UNKNOWN}. It holds the combined watermark back + * the same way an active output does, but it is not evidence that anything is producing + * yet: outputs are registered for every assigned split, even when the reader never reports + * anything through them (for example when it emits everything through the main output, or + * when the watermark strategy does not use idleness at all). + */ + enum IdlenessState { + UNKNOWN, + ACTIVE, + IDLE + } + private long watermark = Long.MIN_VALUE; - private boolean idle = false; + private IdlenessState idlenessState = IdlenessState.UNKNOWN; private final WatermarkOutputMultiplexer.WatermarkUpdateListener onWatermarkUpdate; public PartialWatermark( @@ -139,11 +169,16 @@ public boolean setWatermark(long watermark) { } private boolean isIdle() { - return idle; + return idlenessState == IdlenessState.IDLE; + } + + /** Returns true if this output has reported that it is active. */ + private boolean isActive() { + return idlenessState == IdlenessState.ACTIVE; } public void setIdle(boolean idle) { - this.idle = idle; + this.idlenessState = idle ? IdlenessState.IDLE : IdlenessState.ACTIVE; this.onWatermarkUpdate.onIdleUpdate(idle); } } From b4e3b2631530814ade392f168c64d072593b97de Mon Sep 17 00:00:00 2001 From: "yongfu.gao" <2642474295@qq.com> Date: Sun, 27 Sep 2026 01:36:33 +0800 Subject: [PATCH 2/4] [FLINK-40504][runtime] Mark WatermarkToDataOutput active on non-advancing watermarks emitWatermark() returned at the monotonicity guard before marking the output active, although the WatermarkOutput contract states that emitting a watermark implicitly marks the stream active. A source that resumes with watermarks that are not larger than the max watermark so far therefore stayed announced as idle downstream, and its records could be dropped as late. The activation is now reported before the guard. Source implementers can observe this change, so it needs a release note. Generated-by: DeepSeek Harness (deepseek-v4.1-flash) --- .../source/WatermarkToDataOutput.java | 18 +++++++++++++----- .../source/WatermarkToDataOutputTest.java | 15 +++++++++++++++ 2 files changed, 28 insertions(+), 5 deletions(-) diff --git a/flink-runtime/src/main/java/org/apache/flink/streaming/api/operators/source/WatermarkToDataOutput.java b/flink-runtime/src/main/java/org/apache/flink/streaming/api/operators/source/WatermarkToDataOutput.java index a283443634a35..2aced695b0b24 100644 --- a/flink-runtime/src/main/java/org/apache/flink/streaming/api/operators/source/WatermarkToDataOutput.java +++ b/flink-runtime/src/main/java/org/apache/flink/streaming/api/operators/source/WatermarkToDataOutput.java @@ -74,16 +74,24 @@ public WatermarkToDataOutput( @Override public void emitWatermark(Watermark watermark) { final long newWatermark = watermark.getTimestamp(); - if (newWatermark <= maxWatermarkSoFar) { - return; + final boolean watermarkAdvanced = newWatermark > maxWatermarkSoFar; + if (watermarkAdvanced) { + maxWatermarkSoFar = newWatermark; + watermarkEmitted.updateCurrentEffectiveWatermark(maxWatermarkSoFar); } - maxWatermarkSoFar = newWatermark; - watermarkEmitted.updateCurrentEffectiveWatermark(maxWatermarkSoFar); - try { + // Emitting a watermark implicitly marks this output as active, independently of whether + // the watermark advances the effective watermark (see WatermarkOutput). Otherwise a + // source that resumes with watermarks that are not larger than the max watermark so far + // would stay idle downstream. The activation must therefore not be gated on the + // monotonicity check below. markActiveInternally(); + if (!watermarkAdvanced) { + return; + } + output.emitWatermark( new org.apache.flink.streaming.api.watermark.Watermark(newWatermark)); } catch (ExceptionInChainedOperatorException e) { diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/WatermarkToDataOutputTest.java b/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/WatermarkToDataOutputTest.java index 2e261dfa08d7b..4b534f99c98c1 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/WatermarkToDataOutputTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/WatermarkToDataOutputTest.java @@ -65,4 +65,19 @@ void becomingActiveEmitsStatus() { assertThat(testingOutput.events) .contains(WatermarkStatus.IDLE, WatermarkStatus.ACTIVE, new Watermark(100L)); } + + @Test + void becomingActiveEmitsStatusOnNonAdvancingWatermark() { + final CollectingDataOutput testingOutput = new CollectingDataOutput<>(); + final WatermarkToDataOutput wmOutput = new WatermarkToDataOutput(testingOutput); + + wmOutput.emitWatermark(new org.apache.flink.api.common.eventtime.Watermark(100L)); + wmOutput.markIdle(); + + // the watermark does not advance, but the source is active again + wmOutput.emitWatermark(new org.apache.flink.api.common.eventtime.Watermark(100L)); + + assertThat(testingOutput.events) + .containsExactly(new Watermark(100L), WatermarkStatus.IDLE, WatermarkStatus.ACTIVE); + } } From d9e7c4bb5a91641f70cea4394129925914aabd58 Mon Sep 17 00:00:00 2001 From: "yongfu.gao" <2642474295@qq.com> Date: Sun, 27 Sep 2026 01:36:45 +0800 Subject: [PATCH 3/4] [FLINK-40504][core] Report activity from WatermarksWithIdleness on events onEvent() reset the idleness timer without undoing the markIdle() it had emitted earlier, so the output stayed announced as idle until the wrapped generator produced an advancing watermark. It now reports activity on the first record and after each idle period, so an output whose generator never emits a watermark is still known to be active. WatermarksWithIdleness is @Public, so this is a separate commit to keep its release note easy to trace. Generated-by: DeepSeek Harness (deepseek-v4.1-flash) --- .../eventtime/WatermarksWithIdleness.java | 20 +++++++++++++++++-- .../eventtime/WatermarksWithIdlenessTest.java | 20 +++++++++++++++++++ 2 files changed, 38 insertions(+), 2 deletions(-) diff --git a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java index 4cc651cf6e913..586190f14d7c2 100644 --- a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java +++ b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java @@ -32,7 +32,7 @@ /** * A WatermarkGenerator that adds idleness detection to another WatermarkGenerator. If no events * come within a certain time (timeout duration) then this generator marks the stream as idle, until - * the next watermark is generated. + * an event arrives or the wrapped generator emits a watermark. */ @Public public class WatermarksWithIdleness implements WatermarkGenerator { @@ -43,6 +43,13 @@ public class WatermarksWithIdleness implements WatermarkGenerator { private boolean isIdleNow = false; + /** + * Whether the wrapped generator has ever seen an event. The first event has to report activity + * even though nothing was marked idle before, otherwise an output whose generator never emits a + * watermark would never be known to be active. + */ + private boolean isActivityReported = false; + /** * Creates a new WatermarksWithIdleness generator to the given generator idleness detection with * the given timeout. @@ -67,7 +74,16 @@ public WatermarksWithIdleness( public void onEvent(T event, long eventTimestamp, WatermarkOutput output) { watermarks.onEvent(event, eventTimestamp, output); idlenessTimer.activity(); - isIdleNow = false; + + if (isIdleNow || !isActivityReported) { + // A record is evidence of activity in its own right. Waiting for the wrapped generator + // to produce an advancing watermark instead would leave the output announced as idle + // for as long as the resumed input stays behind the watermark it reached before, and an + // output whose generator never emits a watermark would never be known to be active. + output.markActive(); + isIdleNow = false; + isActivityReported = true; + } } @Override diff --git a/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarksWithIdlenessTest.java b/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarksWithIdlenessTest.java index f2af73de2b0fe..1a89a841df013 100644 --- a/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarksWithIdlenessTest.java +++ b/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarksWithIdlenessTest.java @@ -107,6 +107,26 @@ void testIdleActiveIdle() { assertThat(timer.checkIfIdle()).isTrue(); } + @Test + void testMarksActiveOnFirstEventAfterIdleness() { + final Duration idleTimeout = Duration.ofMillis(10); + final ManualClock clock = new ManualClock(System.nanoTime()); + final WatermarksWithIdleness watermarks = + new WatermarksWithIdleness<>(new NoWatermarksGenerator<>(), idleTimeout, clock); + final TestingWatermarkOutput output = new TestingWatermarkOutput(); + + watermarks.onPeriodicEmit(output); // start the idleness timer + clock.advanceTime(idleTimeout.plusMillis(1)); + watermarks.onPeriodicEmit(output); + assertThat(output.isIdle()).isTrue(); + + watermarks.onEvent(1L, 1L, output); + + assertThat(output.isIdle()) + .as("a record is activity in its own right, no watermark has to advance") + .isFalse(); + } + private static IdlenessTimer createTimerAndMakeIdle(ManualClock clock, Duration idleTimeout) { final IdlenessTimer timer = new IdlenessTimer(clock, idleTimeout); From a93ea2dbd73da6f9d540d1d4254b836d109cb4a2 Mon Sep 17 00:00:00 2001 From: "yongfu.gao" <2642474295@qq.com> Date: Sun, 27 Sep 2026 01:36:45 +0800 Subject: [PATCH 4/4] [FLINK-40504][core] Propagate active combined status from WatermarkOutputMultiplexer updateCombinedWatermark() reported the combined idle state to the underlying output, but nothing when it was active. An output that is active without advancing its watermark therefore produced neither an emitWatermark nor a markIdle call, and the underlying output never learned that it was active again. Because the split branch of IdlenessManager starts out idle and is only cleared by one of those two calls, an active split was then ANDed into IDLE by the main output's idleness timer. The combined status is now reported whenever it is idle or active, but activity is only announced when some output actually reported being active (hasActiveOutput()), so an output that is registered for an assigned split but never reports anything holds the combined watermark back without keeping the downstream output active. Generated-by: DeepSeek Harness (deepseek-v4.1-flash) --- .../eventtime/WatermarkOutputMultiplexer.java | 13 + .../WatermarkOutputMultiplexerTest.java | 107 +++++++ ...tampsAndWatermarksSubtaskIdlenessTest.java | 271 ++++++++++++++++++ .../source/SourceOperatorEventTimeTest.java | 33 +++ 4 files changed, 424 insertions(+) create mode 100644 flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest.java diff --git a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarkOutputMultiplexer.java b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarkOutputMultiplexer.java index 91a372722d742..e4df5241266f8 100644 --- a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarkOutputMultiplexer.java +++ b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarkOutputMultiplexer.java @@ -153,14 +153,27 @@ public void onPeriodicEmit() { * *

It also handles scenarios where both emitting a watermark and entering the idle state * occur within the same invocation. + * + *

Idleness is propagated in both directions: the combined status is reported to the {@link + * #underlyingOutput} whenever it is idle or active. An output that is active but whose + * watermark does not advance produces neither a watermark nor an idle update, so without + * reporting the active state the underlying output would never learn that it is active again. + * + *

Activity is only announced for outputs that reported being active. An output is registered + * for every assigned split, even when the reader never reports anything through it, and + * reporting idleness on one output says nothing about the others: such outputs hold the + * combined watermark back, but they must not keep the downstream output active. */ private void updateCombinedWatermark() { if (combinedWatermarkStatus.updateCombinedWatermark()) { underlyingOutput.emitWatermark( new Watermark(combinedWatermarkStatus.getCombinedWatermark())); } + if (combinedWatermarkStatus.isIdle()) { underlyingOutput.markIdle(); + } else if (combinedWatermarkStatus.hasActiveOutput()) { + underlyingOutput.markActive(); } } diff --git a/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarkOutputMultiplexerTest.java b/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarkOutputMultiplexerTest.java index 4ff83f9c24a5e..22bf029fb3147 100644 --- a/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarkOutputMultiplexerTest.java +++ b/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarkOutputMultiplexerTest.java @@ -20,6 +20,8 @@ import org.junit.jupiter.api.Test; +import java.util.ArrayList; +import java.util.List; import java.util.UUID; import static org.assertj.core.api.Assertions.assertThat; @@ -444,6 +446,111 @@ void testNotEmittingIdleAfterAllSplitsRemoved() { assertThat(underlyingWatermarkOutput.isIdle()).isFalse(); } + @Test + void whenIdleImmediateOutputsBecomeActiveUnderlyingOutputIsMarkedActive() { + TestingWatermarkOutput underlyingWatermarkOutput = createTestingWatermarkOutput(); + WatermarkOutputMultiplexer multiplexer = + new WatermarkOutputMultiplexer(underlyingWatermarkOutput); + + WatermarkOutput watermarkOutput1 = createImmediateOutput(multiplexer); + WatermarkOutput watermarkOutput2 = createImmediateOutput(multiplexer); + + watermarkOutput1.emitWatermark(new Watermark(5)); + watermarkOutput2.emitWatermark(new Watermark(2)); + watermarkOutput1.markIdle(); + watermarkOutput2.markIdle(); + + assertThat(underlyingWatermarkOutput.lastWatermark()).isEqualTo(new Watermark(5)); + assertThat(underlyingWatermarkOutput.isIdle()).isTrue(); + + watermarkOutput1.markActive(); + + assertThat(underlyingWatermarkOutput.isIdle()).isFalse(); + } + + @Test + void whenIdleDeferredOutputResumesUnderlyingOutputIsMarkedActive() { + TestingWatermarkOutput underlyingWatermarkOutput = createTestingWatermarkOutput(); + WatermarkOutputMultiplexer multiplexer = + new WatermarkOutputMultiplexer(underlyingWatermarkOutput); + + WatermarkOutput watermarkOutput1 = createDeferredOutput(multiplexer); + WatermarkOutput watermarkOutput2 = createDeferredOutput(multiplexer); + + watermarkOutput1.emitWatermark(new Watermark(5)); + watermarkOutput2.emitWatermark(new Watermark(2)); + watermarkOutput1.markIdle(); + watermarkOutput2.markIdle(); + + multiplexer.onPeriodicEmit(); + + assertThat(underlyingWatermarkOutput.lastWatermark()).isEqualTo(new Watermark(5)); + assertThat(underlyingWatermarkOutput.isIdle()).isTrue(); + + // the resumed output only has backlog that does not advance the combined watermark + watermarkOutput1.emitWatermark(new Watermark(3)); + multiplexer.onPeriodicEmit(); + + assertThat(underlyingWatermarkOutput.lastWatermark()).isEqualTo(new Watermark(5)); + assertThat(underlyingWatermarkOutput.isIdle()).isFalse(); + } + + @Test + void whenThereAreNoOutputsNothingIsReported() { + final RecordingWatermarkOutput underlyingWatermarkOutput = new RecordingWatermarkOutput(); + final WatermarkOutputMultiplexer multiplexer = + new WatermarkOutputMultiplexer(underlyingWatermarkOutput); + + multiplexer.onPeriodicEmit(); + + assertThat(underlyingWatermarkOutput.watermarks) + .as("without any output there is no combined status to report") + .isEmpty(); + assertThat(underlyingWatermarkOutput.idleUpdates).isZero(); + assertThat(underlyingWatermarkOutput.activeUpdates).isZero(); + } + + @Test + void whenRegisteredOutputReportsNothingNothingIsReported() { + final RecordingWatermarkOutput underlyingWatermarkOutput = new RecordingWatermarkOutput(); + final WatermarkOutputMultiplexer multiplexer = + new WatermarkOutputMultiplexer(underlyingWatermarkOutput); + + // an output is registered for every assigned split, even when the reader never reports + // anything through it, for example when it emits everything through the main output + multiplexer.registerNewOutput("silent-output"); + multiplexer.onPeriodicEmit(); + + assertThat(underlyingWatermarkOutput.watermarks) + .as("an output that never reported anything must not keep the stream active") + .isEmpty(); + assertThat(underlyingWatermarkOutput.idleUpdates).isZero(); + assertThat(underlyingWatermarkOutput.activeUpdates).isZero(); + } + + /** A {@link WatermarkOutput} that records everything it is asked to report. */ + private static final class RecordingWatermarkOutput implements WatermarkOutput { + + private final List watermarks = new ArrayList<>(); + private int idleUpdates; + private int activeUpdates; + + @Override + public void emitWatermark(Watermark watermark) { + watermarks.add(watermark); + } + + @Override + public void markIdle() { + idleUpdates++; + } + + @Override + public void markActive() { + activeUpdates++; + } + } + /** * Convenience method so we don't have to go through the output ID dance when we only want an * immediate output for a given output ID. diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest.java b/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest.java new file mode 100644 index 0000000000000..88d6597371e4c --- /dev/null +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest.java @@ -0,0 +1,271 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.streaming.api.operators.source; + +import org.apache.flink.api.common.eventtime.WatermarkGenerator; +import org.apache.flink.api.common.eventtime.WatermarkOutput; +import org.apache.flink.api.common.eventtime.WatermarkStrategy; +import org.apache.flink.api.connector.source.ReaderOutput; +import org.apache.flink.api.connector.source.SourceOutput; +import org.apache.flink.metrics.groups.UnregisteredMetricsGroup; +import org.apache.flink.runtime.metrics.groups.UnregisteredMetricGroups; +import org.apache.flink.streaming.runtime.tasks.TestProcessingTimeService; +import org.apache.flink.util.clock.Clock; +import org.apache.flink.util.clock.ManualClock; +import org.apache.flink.util.clock.RelativeClock; +import org.apache.flink.util.clock.SystemClock; + +import org.junit.jupiter.api.Test; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests for the subtask-level idleness signal that {@link ProgressiveTimestampsAndWatermarks} + * announces downstream. + * + *

That signal is the conjunction of two independent branches - the merged per-split outputs and + * the main output - so these tests pin down when the conjunction may and may not conclude that the + * whole subtask is idle. + */ +class ProgressiveTimestampsAndWatermarksSubtaskIdlenessTest { + + private static final Duration IDLE_TIMEOUT = Duration.ofMillis(10); + + /** Timestamp of the record that a split emits before it falls idle. */ + private static final long FIRST_EVENT_TIMESTAMP = 100L; + + /** Timestamp of the record the split resumes with, behind {@link #FIRST_EVENT_TIMESTAMP}. */ + private static final long RESUMED_EVENT_TIMESTAMP = 50L; + + /** + * A subtask must not announce idleness while one of its splits is still producing records, even + * if that split's watermark never advances. The main output of a split-based source never sees + * a record, so its activity timer always expires; the split branch is what has to hold the + * subtask active. + */ + @Test + void subtaskMustNotGoIdleWhileRecordsFlowThroughItsSplits() { + final ManualClock mainInputActivityClock = new ManualClock(System.nanoTime()); + final RecordingListener listener = new RecordingListener(); + // a generator that never emits: the split stays active but its watermark stays at MIN_VALUE + final TimestampsAndWatermarks eventTimeLogic = + createEventTimeLogic( + WatermarkStrategy.forGenerator(context -> new NeverEmits()) + .withTimestampAssigner((event, timestamp) -> event) + .withIdleness(IDLE_TIMEOUT), + mainInputActivityClock, + SystemClock.getInstance()); + + final ReaderOutput mainOutput = + eventTimeLogic.createMainOutput(new CollectingDataOutput<>(), listener); + final SourceOutput splitOutput = mainOutput.createOutputForSplit("split-A"); + + splitOutput.collect(1L, 1L); + eventTimeLogic.emitImmediateWatermark(0L); + assertThat(listener.idleUpdates).isEmpty(); + + mainInputActivityClock.advanceTime(IDLE_TIMEOUT.plusMillis(1)); + splitOutput.collect(2L, 2L); + eventTimeLogic.emitImmediateWatermark(0L); + + assertThat(listener.idleUpdates) + .as("records kept arriving on split-A, so the subtask is not idle") + .isEmpty(); + } + + /** Once every split really is idle, the subtask must still stand down. */ + @Test + void subtaskGoesIdleOnceAllOfItsSplitsAreIdle() { + final ManualClock clock = new ManualClock(System.nanoTime()); + final RecordingListener listener = new RecordingListener(); + final TimestampsAndWatermarks eventTimeLogic = createEventTimeLogic(clock, clock); + + final ReaderOutput mainOutput = + eventTimeLogic.createMainOutput(new CollectingDataOutput<>(), listener); + final SourceOutput splitOutput = mainOutput.createOutputForSplit("split-A"); + + splitOutput.collect(1L, 1L); + eventTimeLogic.emitImmediateWatermark(0L); + assertThat(listener.idleUpdates).isEmpty(); + + // the idleness timers need one emit to start counting and one more to expire + for (int i = 0; i < 2; i++) { + clock.advanceTime(IDLE_TIMEOUT.plusMillis(1)); + eventTimeLogic.emitImmediateWatermark(0L); + } + + assertThat(listener.idleUpdates) + .as("split-A stopped producing, so the subtask has nothing left to hold it active") + .containsExactly(true); + } + + /** + * A subtask that was never assigned a split has no split branch to speak for it, so the main + * output's activity timer alone has to make it stand down - otherwise it would pin the + * downstream watermark forever. + */ + @Test + void subtaskWithoutSplitsGoesIdleViaTheMainOutputTimer() { + final ManualClock mainInputActivityClock = new ManualClock(System.nanoTime()); + final RecordingListener listener = new RecordingListener(); + final TimestampsAndWatermarks eventTimeLogic = + createEventTimeLogic(mainInputActivityClock, SystemClock.getInstance()); + + eventTimeLogic.createMainOutput(new CollectingDataOutput<>(), listener); + + eventTimeLogic.emitImmediateWatermark(0L); + assertThat(listener.idleUpdates).isEmpty(); + + mainInputActivityClock.advanceTime(IDLE_TIMEOUT.plusMillis(1)); + eventTimeLogic.emitImmediateWatermark(0L); + + assertThat(listener.idleUpdates) + .as("a subtask without splits must still signal idleness downstream") + .containsExactly(true); + } + + /** + * When all splits fall idle their watermarks are flushed to the maximum seen so far. A split + * that resumes behind that maximum therefore produces no advancing watermark, and the subtask + * would stay announced as idle even though it is emitting records again. + */ + @Test + void subtaskGoesActiveAgainWhenAnIdleSplitResumesBehindTheFlushedWatermark() { + final ManualClock clock = new ManualClock(System.nanoTime()); + final RecordingListener listener = new RecordingListener(); + final TimestampsAndWatermarks eventTimeLogic = createEventTimeLogic(clock, clock); + + final ReaderOutput mainOutput = + eventTimeLogic.createMainOutput(new CollectingDataOutput<>(), listener); + final SourceOutput splitOutput = mainOutput.createOutputForSplit("split-A"); + + splitOutput.collect(FIRST_EVENT_TIMESTAMP, FIRST_EVENT_TIMESTAMP); + eventTimeLogic.emitImmediateWatermark(0L); + + // the idleness timers need one emit to start counting and one more to expire + for (int i = 0; i < 2; i++) { + clock.advanceTime(IDLE_TIMEOUT.plusMillis(1)); + eventTimeLogic.emitImmediateWatermark(0L); + } + assertThat(listener.idleUpdates).containsExactly(true); + + splitOutput.collect(RESUMED_EVENT_TIMESTAMP, RESUMED_EVENT_TIMESTAMP); + eventTimeLogic.emitImmediateWatermark(0L); + + assertThat(listener.idleUpdates) + .as("split-A produces records again, so the subtask has to be announced as active") + .containsExactly(true, false); + } + + /** + * A split that is registered but never reports anything holds the combined watermark back, but + * it is not evidence that anything is producing. Reporting idleness on another split must + * therefore not re-activate the subtask. + */ + @Test + void subtaskStaysIdleWhenOnlyAnUnknownAndAnIdleSplitAreRegistered() { + final ManualClock clock = new ManualClock(System.nanoTime()); + final RecordingListener listener = new RecordingListener(); + // generators that never emit, so that the registered splits really stay silent + final TimestampsAndWatermarks eventTimeLogic = + createEventTimeLogic( + WatermarkStrategy.forGenerator(context -> new NeverEmits()) + .withTimestampAssigner((event, timestamp) -> event) + .withIdleness(IDLE_TIMEOUT), + clock, + clock); + + final ReaderOutput mainOutput = + eventTimeLogic.createMainOutput(new CollectingDataOutput<>(), listener); + + // no split is assigned yet, so the subtask stands down through the main output's timer + for (int i = 0; i < 2; i++) { + clock.advanceTime(IDLE_TIMEOUT.plusMillis(1)); + eventTimeLogic.emitImmediateWatermark(0L); + } + assertThat(listener.idleUpdates).containsExactly(true); + + // one split never reports anything, the other only reports that it is idle + mainOutput.createOutputForSplit("silent"); + mainOutput.createOutputForSplit("idle").markIdle(); + eventTimeLogic.emitImmediateWatermark(0L); + + assertThat(listener.idleUpdates) + .as("no split has reported activity, so the subtask stays idle") + .containsExactly(true); + } + + // ------------------------------------------------------------------------ + + private static TimestampsAndWatermarks createEventTimeLogic( + RelativeClock mainInputActivityClock, Clock clock) { + return createEventTimeLogic( + WatermarkStrategy.forMonotonousTimestamps() + .withTimestampAssigner((event, timestamp) -> event) + .withIdleness(IDLE_TIMEOUT), + mainInputActivityClock, + clock); + } + + private static TimestampsAndWatermarks createEventTimeLogic( + WatermarkStrategy strategy, RelativeClock mainInputActivityClock, Clock clock) { + return TimestampsAndWatermarks.createProgressiveEventTimeLogic( + strategy, + new UnregisteredMetricsGroup(), + new TestProcessingTimeService(), + 0L, + mainInputActivityClock, + clock, + UnregisteredMetricGroups.createUnregisteredTaskMetricGroup().getIOMetricGroup()); + } + + private static final class NeverEmits implements WatermarkGenerator { + @Override + public void onEvent(Long event, long eventTimestamp, WatermarkOutput output) {} + + @Override + public void onPeriodicEmit(WatermarkOutput output) {} + } + + private static final class RecordingListener + implements TimestampsAndWatermarks.WatermarkUpdateListener { + private final List idleUpdates = new ArrayList<>(); + + @Override + public void updateIdle(boolean isIdle) { + idleUpdates.add(isIdle); + } + + @Override + public void updateCurrentEffectiveWatermark(long watermark) {} + + @Override + public void updateCurrentSplitWatermark(String splitId, long watermark) {} + + @Override + public void updateCurrentSplitIdle(String splitId, boolean idle) {} + + @Override + public void splitFinished(String splitId) {} + } +} diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/SourceOperatorEventTimeTest.java b/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/SourceOperatorEventTimeTest.java index 9e295abc50f1c..9819bf755df24 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/SourceOperatorEventTimeTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/api/operators/source/SourceOperatorEventTimeTest.java @@ -244,6 +244,39 @@ void testMainAndPerSplitWatermarkIdleness() throws Exception { new Watermark(300)); } + @Test + void testResumingSplitWithoutAdvancingWatermarkEmitsActive() throws Exception { + final WatermarkStrategy watermarkStrategy = + WatermarkStrategy.forGenerator((ctx) -> new OnEventTestWatermarkGenerator<>()); + + InterpretingSourceReader reader = + new InterpretingSourceReader( + // Both splits emit a watermark and then go idle. The combined watermark is + // flushed to the maximum watermark (200) and the edge goes IDLE. + output -> output.createOutputForSplit("1").collect(0, 200L), + output -> output.createOutputForSplit("2").collect(0, 100L), + output -> output.createOutputForSplit("1").markIdle(), + output -> output.createOutputForSplit("2").markIdle(), + // Split 1 resumes with backlog that does not advance the combined + // watermark. The edge has to be re-activated, otherwise its records are + // dropped as late downstream. + output -> output.createOutputForSplit("1").collect(0, 150L)); + + SourceOperator sourceOperator = + createTestOperator(reader, watermarkStrategy, true); + + List events = testSequenceOfEvents(sourceOperator); + + assertThat(events) + .containsExactly( + new StreamRecord<>(0, 200L), + new Watermark(200L), + new StreamRecord<>(0, 100L), + new WatermarkStatus(WatermarkStatus.IDLE_STATUS), + new StreamRecord<>(0, 150L), + new WatermarkStatus(WatermarkStatus.ACTIVE_STATUS)); + } + // ------------------------------------------------------------------------ // test execution helpers // ------------------------------------------------------------------------