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 9e9524957f50..b6a579072512 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 @@ -282,7 +282,7 @@ public FinishBundleContext finishBundleContext(DoFn doFn) { processContext.tracker.checkDone(); } if (residual == null) { - return new Result(null, cont, null, null); + return new Result(null, cont, null, null, 0.0); } final KV> 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 a750b01963f6..a93dd8ecd509 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 @@ -606,6 +606,9 @@ public String getErrorContext() { 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 52ac6b1a8199..6ff40fae91dd 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 void testInvokeProcessElementVoluntaryReturnStop() throws Exception { 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 void testBacklogBytes() throws Exception { 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.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 381e41c98705..e01a61ef4a65 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 @@ -763,6 +763,10 @@ public void testReportsBacklog() throws Exception { // 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); } } @@ -788,6 +792,34 @@ public void testReportsBacklogWithoutGetSize() throws Exception { // 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 fn = new GetSizeFn(); + Instant base = Instant.now(); + final List backlogs = new ArrayList<>(); + + try (ProcessFnTester 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); } } }