Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down Expand Up @@ -108,8 +122,24 @@ public boolean updateCombinedWatermark() {

/** Per-output watermark state. */
static class PartialWatermark {

/**
* The idleness the output has reported, if any.
*
* <p>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(
Expand Down Expand Up @@ -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);
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -153,14 +153,27 @@ public void onPeriodicEmit() {
*
* <p>It also handles scenarios where both emitting a watermark and entering the idle state
* occur within the same invocation.
*
* <p>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.
*
* <p>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();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<T> implements WatermarkGenerator<T> {
Expand All @@ -43,6 +43,13 @@ public class WatermarksWithIdleness<T> implements WatermarkGenerator<T> {

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.
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Watermark> 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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Long> 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);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down
Loading