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

Reply via email to