damccorm commented on code in PR #40149:
URL: https://github.com/apache/beam/pull/40149#discussion_r4209071008
##########
runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferingDoFnRunner.java:
##########
@@ -299,7 +299,12 @@ public void checkpointCompleted(long checkpointId) throws
Exception {
bufferingElementsHandler.clear();
}
}
- minBufferedElementTimestamp = Long.MAX_VALUE;
+ minBufferedElementTimestamp =
+ currentBufferingElementsHandler
+ .getElements()
+ .mapToLong(e -> e.getTimestamp().getMillis())
+ .min()
+ .orElse(Long.MAX_VALUE);
Review Comment:
A few thoughts on how `minBufferedElementTimestamp` is recomputed here:
1. **Concurrent checkpoints (`maxConcurrentCheckpoints > 1`):** If
checkpoint `N` completes while checkpoint `N+1` has already been snapshotted
(and is still pending in `notYetAcknowledgedSnapshots`),
`currentBufferingElementsHandler` has already rotated to the buffer for `N+2`.
Only checking `currentBufferingElementsHandler` will miss elements still
buffered in unacknowledged checkpoint `N+1` and can drop their watermark hold.
2. **Locking (`Locker`):** When using `KeyedBufferingElementsHandler` (e.g.,
in `ExecutableStageDoFnOperator`), `getElements()` iterates keys and calls
`backend.setCurrentKey(key)` / `state.get()` on the `KeyedStateBackend`. All
other accesses to `BufferingElementsHandler` in this class are guarded by `try
(Locker lock = locker != null ? locker.get() : null)`.
3. **Avoiding extra state reads/deserialization:** Calling
`currentBufferingElementsHandler.getElements()` on every `checkpointCompleted`
forces Flink to read and deserialize all buffered elements in the active buffer
(and scan all keys in RocksDB when keyed), only to read and deserialize them
again when the next checkpoint completes.
Could we instead track the minimum timestamp per buffer index in memory
(e.g., a `long[] minTimestampPerBuffer` of size `numCheckpointBuffers`
initialized to `Long.MAX_VALUE`)?
- In `processElement` / `onTimer`, update
`minTimestampPerBuffer[currentStateIndex]` (and `minBufferedElementTimestamp`).
- In `checkpointCompleted`, reset
`minTimestampPerBuffer[toBeAcked.internalId] = Long.MAX_VALUE` as each buffer
is cleared, and then set `minBufferedElementTimestamp` to the min across
`minTimestampPerBuffer`.
That would avoid scanning/deserializing Flink state in
`checkpointCompleted`, avoid the `Locker` issue, and naturally preserve
watermark holds across multiple concurrent checkpoints. (If we also want holds
preserved across job restarts, we could populate `minTimestampPerBuffer` from
`pendingSnapshots` once in `initializeState`.)
##########
runners/flink/src/test/java/org/apache/beam/runners/flink/translation/wrappers/streaming/DoFnOperatorTest.java:
##########
@@ -2233,6 +2234,87 @@ public void finishBundle(FinishBundleContext context) {
WindowedValues.valueInGlobalWindow(KV.of("key3",
"finishBundle"))));
}
+ /** Ensures that output watermark does not pass buffered elements */
+ @Test
+ public void testExactlyOnceBufferingOutputWatermarkCorrectness() throws
Exception {
+ FlinkPipelineOptions options = FlinkPipelineOptions.defaults();
+ options.setStreaming(true);
+ options.setMaxBundleSize(1L);
+ options.setCheckpointingInterval(1L);
+ options.setCheckpointingMode("EXACTLY_ONCE");
+ options.setNumConcurrentCheckpoints(1);
+
+ TupleTag<String> outputTag = new TupleTag<>("main-output");
+
+ WindowedValues.FullWindowedValueCoder<String> windowedValueCoder =
+ WindowedValues.getFullCoder(StringUtf8Coder.of(),
GlobalWindow.Coder.INSTANCE);
+
+ DoFn<String, String> doFn =
+ new DoFn<String, String>() {
+ @ProcessElement
+ // Use RequiresStableInput to force buffering elements
+ @RequiresStableInput
+ public void processElement(ProcessContext context) {
+ context.output(context.element());
+ }
+ };
+
+ DoFnOperator.MultiOutputOutputManagerFactory<String> outputManagerFactory =
+ new DoFnOperator.MultiOutputOutputManagerFactory<>(
+ outputTag,
+ WindowedValues.getFullCoder(StringUtf8Coder.of(),
GlobalWindow.Coder.INSTANCE),
+ new SerializablePipelineOptions(options));
+
+ Supplier<DoFnOperator<String, String, String>> doFnOperatorSupplier =
+ () ->
+ new DoFnOperator<>(
+ doFn,
+ "stepName",
+ windowedValueCoder,
+ Collections.emptyMap(),
+ outputTag,
+ Collections.emptyList(),
+ outputManagerFactory,
+ WindowingStrategy.globalDefault(),
+ new HashMap<>(), /* side-input mapping */
+ Collections.emptyList(), /* side inputs */
+ options,
+ null,
+ null,
+ DoFnSchemaInformation.create(),
+ Collections.emptyMap());
+
+ DoFnOperator<String, String, String> doFnOperator =
doFnOperatorSupplier.get();
+ OneInputStreamOperatorTestHarness<WindowedValue<String>,
WindowedValue<String>> testHarness =
+ new OneInputStreamOperatorTestHarness<>(doFnOperator);
+
+ testHarness.open();
+
+ 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(
+ "A@10 has not yet completed the checkpoint; output watermark must be
held at 10 until checkpoint completion",
+ doFnOperator.getCurrentOutputWatermark(),
+ is(10L));
+
+ doFnOperator.notifyCheckpointComplete(1L);
+ assertThat(
+ "A@10 has been checked but B@20 is still buffered; output watermark
must not pass it",
+ doFnOperator.getCurrentOutputWatermark(),
+ is(20L));
+
+ testHarness.snapshot(2L, 0L);
+ doFnOperator.notifyCheckpointComplete(2L);
+ assertThat(
+ "B@20 has been checked and the buffer is empty; output watermark must
advance to input watermark",
+ doFnOperator.getCurrentOutputWatermark(),
+ is(100L));
+ }
Review Comment:
Nit: Since this test doesn't restore from a snapshot, `doFnOperatorSupplier`
can be inlined to construct `doFnOperator` directly, and we should call
`testHarness.close()` at the end (or use try-with-resources). Also might be
nice to assert on `testHarness.getOutput()` at the end to verify that `A@10`
and `B@20` were emitted in order before watermark `100`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]