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