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);
}
}
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/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/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-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);
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/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