This is an automated email from the ASF dual-hosted git repository.
tvalentyn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new 790468a6591 Ensure 0 backlog is sent when finishing processing
restrictions. (#40102)
790468a6591 is described below
commit 790468a6591cd1dea07365c2294f7de07918e723
Author: Andrew Crites <[email protected]>
AuthorDate: Sat Sep 12 19:04:13 2026 -0700
Ensure 0 backlog is sent when finishing processing restrictions. (#40102)
---
...TimeBoundedSplittableProcessElementInvoker.java | 2 +-
.../core/SplittableParDoViaKeyedWorkItems.java | 3 ++
...BoundedSplittableProcessElementInvokerTest.java | 13 +++++++++
.../runners/core/SplittableParDoProcessFnTest.java | 32 ++++++++++++++++++++++
4 files changed, 49 insertions(+), 1 deletion(-)
diff --git
a/runners/core-java/src/main/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvoker.java
b/runners/core-java/src/main/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvoker.java
index 49b873c015c..59d33c596c8 100644
---
a/runners/core-java/src/main/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvoker.java
+++
b/runners/core-java/src/main/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvoker.java
@@ -281,7 +281,7 @@ public class
OutputAndTimeBoundedSplittableProcessElementInvoker<
processContext.tracker.checkDone();
}
if (residual == null) {
- return new Result(null, cont, null, null);
+ return new Result(null, cont, null, null, 0.0);
}
final KV<RestrictionT, KV<Instant, WatermarkEstimatorStateT>>
residualForGetSize = residual;
// For a list of all DoFnInvoker arguments, see DoFn.java.
diff --git
a/runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java
b/runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java
index bdeed8aab77..99fbc6a14f6 100644
---
a/runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java
+++
b/runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java
@@ -673,6 +673,9 @@ public class SplittableParDoViaKeyedWorkItems {
restrictionState.clear();
watermarkEstimatorState.clear();
holdState.clear();
+ if (backlogBytesCallback != null) {
+ backlogBytesCallback.accept(0.0);
+ }
return;
}
diff --git
a/runners/core-java/src/test/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvokerTest.java
b/runners/core-java/src/test/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvokerTest.java
index 52ac6b1a819..6ff40fae91d 100644
---
a/runners/core-java/src/test/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvokerTest.java
+++
b/runners/core-java/src/test/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvokerTest.java
@@ -214,6 +214,7 @@ public class
OutputAndTimeBoundedSplittableProcessElementInvokerTest {
runTest(5, Duration.ZERO, Integer.MAX_VALUE, Duration.millis(100));
assertFalse(res.getContinuation().shouldResume());
assertNull(res.getResidualRestriction());
+ assertEquals(0.0, res.getBacklogBytes(), 0.001);
}
@Test
@@ -274,4 +275,16 @@ public class
OutputAndTimeBoundedSplittableProcessElementInvokerTest {
assertEquals(7.0, res.getBacklogBytes(), 0.001);
assertEquals(new OffsetRange(3, 10), res.getResidualRestriction());
}
+
+ @Test
+ public void testBacklogBytesWhenDone() throws Exception {
+ GetSizeFn fn = new GetSizeFn();
+ OffsetRange initialRestriction = new OffsetRange(0, 2);
+ // Set a high checkpoint duration to prevent flakiness caused by early
checkpointing.
+ SplittableProcessElementInvoker<Void, String, OffsetRange, Long,
Void>.Result res =
+ runTest(fn, initialRestriction, Duration.standardMinutes(3));
+ // GetSizeFn claims 2 elements and finishes.
+ assertEquals(0.0, res.getBacklogBytes(), 0.001);
+ assertNull(res.getResidualRestriction());
+ }
}
diff --git
a/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java
b/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java
index 247412a85dc..73b3f9ddd57 100644
---
a/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java
+++
b/runners/core-java/src/test/java/org/apache/beam/runners/core/SplittableParDoProcessFnTest.java
@@ -775,6 +775,10 @@ public class SplittableParDoProcessFnTest {
// The residual range should be [3, 10), so size is 7.
assertEquals(1, backlogs.size());
assertEquals(7.0, backlogs.get(0), 0.001);
+
+ assertTrue(tester.advanceProcessingTimeBy(Duration.standardSeconds(1)));
+ assertEquals(2, backlogs.size());
+ assertEquals(0.0, backlogs.get(1), 0.001);
}
}
@@ -800,6 +804,34 @@ public class SplittableParDoProcessFnTest {
// The residual range should be [3, 10), so size is 7.
assertEquals(1, backlogs.size());
assertEquals(7.0, backlogs.get(0), 0.001);
+
+ assertTrue(tester.advanceProcessingTimeBy(Duration.standardSeconds(1)));
+ assertEquals(2, backlogs.size());
+ assertEquals(0.0, backlogs.get(1), 0.001);
+ }
+ }
+
+ @Test
+ public void testReportsZeroBacklogWhenDone() throws Exception {
+ DoFn<Integer, String> fn = new GetSizeFn();
+ Instant base = Instant.now();
+ final List<Double> backlogs = new ArrayList<>();
+
+ try (ProcessFnTester<Integer, String, OffsetRange, Long, Void> tester =
+ new ProcessFnTester<>(
+ base,
+ fn,
+ BigEndianIntegerCoder.of(),
+ SerializableCoder.of(OffsetRange.class),
+ VoidCoder.of(),
+ MAX_OUTPUTS_PER_BUNDLE,
+ MAX_BUNDLE_DURATION)) {
+ tester.processFn.setBacklogBytesCallback(backlogs::add);
+
+ // OffsetRange(0, 2) completes immediately without resume.
+ tester.startElement(42, new OffsetRange(0, 2));
+ assertEquals(1, backlogs.size());
+ assertEquals(0.0, backlogs.get(0), 0.001);
}
}