apache / apache/beam

[Bug]: Flink drops watermark holds for records awaiting a checkpoint

Open
#39,991 1 comment 0 reactions 1 assignee Claimed by @yashgoyal0110 View on GitHub
flink P1
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

### What happened?

The Flink runner can emit an output watermark past records still buffered for `@RequiresStableInput`. Revalidated against Apache Beam commit `ffcbd7e464f8244bd623d23e7fa9fcc750e81bbb`.

[`BufferingDoFnRunner.checkpointCompleted()`](https://github.com/apache/beam/blob/ffcbd7e464f8244bd623d23e7fa9fcc750e81bbb/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferingDoFnRunner.java#L290-L316) clears the minimum buffered timestamp unconditionally at line 312, after releasing only the acknowledged buffers. The current buffer can still contain records received after that checkpoint's snapshot. [`DoFnOperator.notifyCheckpointComplete()`](https://github.com/apache/beam/blob/ffcbd7e464f8244bd623d23e7fa9fcc750e81bbb/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/DoFnOperator.java#L1140-L1144) immediately recomputes the output watermark using this now-missing hold.

For a streaming, non-keyed identity `DoFn` with `@RequiresStableInput`, exactly-once checkpointing, one concurrent checkpoint, and `maxBundleSize=1`:

1. Receive `A@10`, then snapshot checkpoint 1.
2. Receive `B@20` into the new active buffer.
3. Receive input watermark 100. The buffered timestamps currently hold output watermark at 10.
4. Complete checkpoint 1. A is emitted, but B remains buffered. The hold is reset to `Long.MAX_VALUE`, allowing output watermark 100.
5. Snapshot and complete checkpoint 2. B is emitted at timestamp 20, after output watermark 100.

Actual output from the source-method reproduction:

```text
[watermark:10, record:A@10, watermark:100, record:B@20]
```

Expected: retain B's hold at 20 (or an earlier conservative hold) until B is emitted. Advancing the watermark to 100 can cause downstream event-time windows to finalize early or discard otherwise on-time records as late. The sequence occurs during normal processing with one checkpoint in flight.

The [`bundleStarted` guard](https://github.com/apache/beam/blob/ffcbd7e464f8244bd623d23e7fa9fcc750e81bbb/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/DoFnOperator.java#L913-L929) does not prevent this: ordinary bundle completion sets it to false, while [`BufferingDoFnRunner.finishBundle()`](https://github.com/apache/beam/blob/ffcbd7e464f8244bd623d23e7fa9fcc750e81bbb/runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferingDoFnRunner.java#L263-L266) leaves these records buffered. A size of 1 makes the sequence deterministic; a count or timeout boundary can also finish the bundle with larger sizes.

To reproduce in the existing test suite, copy the operator/harness setup from [`DoFnOperatorTest.testExactlyOnceBuffering()`](https://github.com/apache/beam/blob/ffcbd7e464f8244bd623d23e7fa9fcc750e81bbb/runners/flink/src/test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/DoFnOperatorTest.java#L1900-L1957), with these changes before constructing the operator:

```java
options.setStreaming(true);
options.setMaxBundleSize(1L);
options.setCheckpointingInterval(1L);
options.setCheckpointingMode("EXACTLY_ONCE");
options.setNumConcurrentCheckpoints(1);
```

Use an identity `@ProcessElement @RequiresStableInput` method and omit the example's `@StartBundle` and `@FinishBundle` hooks. Replace its value-only input coder with a full coder so checkpoint buffering preserves the test timestamps:

```java
WindowedValues.FullWindowedValueCoder windowedValueCoder =
WindowedValues.getFullCoder(StringUtf8Coder.of(), GlobalWindow.Coder.INSTANCE);
```

After `testHarness.open()`, the regression sequence is:

```java
testHarness.processElement(new StreamRecord<>(
WindowedValues.timestampedValueInGlobalWindow("A", new Instant(10))));
testHarness.snapshot(1L, 0L);
testHarness.processElement(new StreamRecord<>(
WindowedValues.timestampedValueInGlobalWindow("B", new Instant(20))));
testHarness.processWatermark(new Watermark(100));
assertThat(doFnOperator.getCurrentOutputWatermark(), is(10L));

doFnOperator.notifyCheckpointComplete(1L);
org.junit.Assert.assertTrue(
"B@20 is still buffered; output watermark must not pass it",
doFnOperator.getCurrentOutputWatermark() <= 20L);

// With the preceding assertion omitted, inspect getOutput() after this
// checkpoint to observe B@20 arriving after watermark 100.
testHarness.snapshot(2L, 0L);
doFnOperator.notifyCheckpointComplete(2L);
```

Validation: compiled and ran a standalone Java harness extracting unchanged checkpoint, buffering, bundle-finishing, watermark, and checkpoint-completion methods from the commit above, with in-memory dependency stubs. It confirmed `bundleStarted=false`, B still buffered, and output watermark 100 after checkpoint 1. The proposed Flink JUnit regression above has not been run; this report does not claim a full Flink integration-test result.

### Issue Priority

Priority: 1 (data loss / total loss of function)

### Issue Components

- [x] Component: Flink Runner
- [x] Component: Java SDK

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.