This is an automated email from the ASF dual-hosted git repository.
Abacn 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 69bfcacf0be Make AsyncWrapperTest deterministic (#40448) (#40449)
69bfcacf0be is described below
commit 69bfcacf0be9369940dfd6aa076335c7ec76db25
Author: Tobias Kaymak <[email protected]>
AuthorDate: Thu Oct 8 22:36:01 2026 +0200
Make AsyncWrapperTest deterministic (#40448) (#40449)
testBufferStopsAcceptingItems asserted the buffer count after a fixed
100ms sleep, which failed on slow Windows runners (expected 5 but was 0).
Add an optional CountDownLatch gate to BasicDofn so tests hold elements
in-flight until they release it, instead of relying on wall-clock sleeps.
testBufferStopsAcceptingItems now joins all producers before asserting
that exactly 5 items are buffered, then verifies the 5 throttled items
are rescheduled and emitted. testBufferCount, testBufferWithCancellation
and testLongItem use the same gate; testSlowDuplicates waits for the
buffer to drain instead of sleeping.
Fixes #40448. Supersedes the timing-based approach in #40397.
---
.../beam/sdk/transforms/AsyncWrapperTest.java | 72 ++++++++++++++--------
1 file changed, 48 insertions(+), 24 deletions(-)
diff --git
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/AsyncWrapperTest.java
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/AsyncWrapperTest.java
index d098104a8cf..cfa38f5e5dd 100644
---
a/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/AsyncWrapperTest.java
+++
b/sdks/java/core/src/test/java/org/apache/beam/sdk/transforms/AsyncWrapperTest.java
@@ -27,6 +27,7 @@ import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Random;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
@@ -50,13 +51,23 @@ public class AsyncWrapperTest implements Serializable {
private final boolean useThreadPool = true;
// Used for testing basic DoFn processing logic with optional latency.
+ // An optional gate holds elements in-flight until the test releases it. If
a test fails before
+ // releasing it, AsyncWrapper.resetState() in setUp() interrupts the blocked
worker.
private static class BasicDofn extends DoFn<String, String> {
private final long sleepTimeMs;
+ // Not Serializable; the DoFn is only invoked in-process here.
+ private final transient CountDownLatch gate;
private int processed = 0;
private final ReentrantLock lock = new ReentrantLock();
BasicDofn(long sleepTimeMs) {
this.sleepTimeMs = sleepTimeMs;
+ this.gate = new CountDownLatch(0);
+ }
+
+ BasicDofn(CountDownLatch gate) {
+ this.sleepTimeMs = 0;
+ this.gate = gate;
}
BasicDofn() {
@@ -65,6 +76,11 @@ public class AsyncWrapperTest implements Serializable {
@ProcessElement
public void processElement(@Element String element, OutputReceiver<String>
receiver) {
+ try {
+ gate.await();
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
if (sleepTimeMs > 0) {
try {
Thread.sleep(sleepTimeMs);
@@ -419,7 +435,8 @@ public class AsyncWrapperTest implements Serializable {
// execution task has not finished processing yet.
@Test
public void testLongItem() {
- BasicDofn dofn = new BasicDofn(500);
+ CountDownLatch gate = new CountDownLatch(1);
+ BasicDofn dofn = new BasicDofn(gate);
AsyncWrapper<String, String, String> asyncWrapper =
new AsyncWrapper<>(
dofn, 1, Duration.standardSeconds(5), null, null, null, null,
useThreadPool);
@@ -438,7 +455,8 @@ public class AsyncWrapperTest implements Serializable {
assertEquals(0, dofn.getProcessed());
assertEquals(1, fakeBagState.items.size());
- waitForEmpty(asyncWrapper, 2);
+ gate.countDown();
+ waitForEmpty(asyncWrapper);
result =
asyncWrapper.commitFinishedItemsDirect(
@@ -580,11 +598,7 @@ public class AsyncWrapperTest implements Serializable {
asyncWrapper.processDirect(msg, GlobalWindow.INSTANCE, Instant.now(),
fakeBagState, fakeTimer);
- try {
- Thread.sleep(100);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- }
+ waitForEmpty(asyncWrapper);
fakeBagState.clear();
List<String> result =
@@ -610,7 +624,8 @@ public class AsyncWrapperTest implements Serializable {
// and decrement immediately upon execution completion.
@Test
public void testBufferCount() {
- BasicDofn dofn = new BasicDofn(10);
+ CountDownLatch gate = new CountDownLatch(1);
+ BasicDofn dofn = new BasicDofn(gate);
AsyncWrapper<String, String, String> asyncWrapper =
new AsyncWrapper<>(
dofn, 1, Duration.standardSeconds(5), null, null, null, null,
useThreadPool);
@@ -623,6 +638,7 @@ public class AsyncWrapperTest implements Serializable {
asyncWrapper.processDirect(msg, GlobalWindow.INSTANCE, Instant.now(),
fakeBagState, fakeTimer);
checkItemsInBuffer(asyncWrapper, 1);
+ gate.countDown();
waitForEmpty(asyncWrapper);
checkItemsInBuffer(asyncWrapper, 0);
@@ -634,18 +650,21 @@ public class AsyncWrapperTest implements Serializable {
// Test 11: testBufferStopsAcceptingItems
// Verifies queue boundaries and backpressure throttling.
// When concurrent threads push elements exceeding the capacity limit,
- // the scheduler must block and delay submissions appropriately.
+ // the scheduler must block, then time out and leave the element in state
for the timer to
+ // reschedule.
@Test
public void testBufferStopsAcceptingItems() {
- BasicDofn dofn = new BasicDofn(500);
+ // Gated items cannot complete, so exactly 5 occupy the buffer and the
other 5 time out.
+ CountDownLatch gate = new CountDownLatch(1);
+ BasicDofn dofn = new BasicDofn(gate);
AsyncWrapper<String, String, String> asyncWrapper =
new AsyncWrapper<>(
dofn,
1,
Duration.standardSeconds(5),
5, // max buffer capacity
- null,
- null,
+ Duration.millis(100), // short timeout so throttled producers give
up fast
+ Duration.millis(20), // max backoff
null,
useThreadPool);
asyncWrapper.setup(null);
@@ -669,16 +688,6 @@ public class AsyncWrapperTest implements Serializable {
}));
}
- try {
- Thread.sleep(100);
- } catch (InterruptedException e) {
- Thread.currentThread().interrupt();
- }
-
- assertEquals(5, asyncWrapper.getItemsInBufferCount());
-
- waitForEmpty(asyncWrapper, 100);
-
// Verify that all background tasks completed successfully without
throwing exceptions
for (Future<?> future : futures) {
try {
@@ -687,10 +696,22 @@ public class AsyncWrapperTest implements Serializable {
throw new AssertionError("Background task failed", e);
}
}
+ poolExecutor.shutdown();
+ checkItemsInBuffer(asyncWrapper, 5);
+ assertEquals(10, fakeBagState.items.size());
+ assertEquals(0, dofn.getProcessed());
+
+ gate.countDown();
+ waitForEmpty(asyncWrapper, 100);
+ assertEquals(5, dofn.getProcessed());
+
+ // Emits the 5 finished items and reschedules the 5 throttled ones.
List<String> result =
asyncWrapper.commitFinishedItemsDirect(
fakeTimer.getCurrentRelativeTime(), fakeBagState, fakeTimer);
+ assertEquals(5, result.size());
+ assertEquals(5, fakeBagState.items.size());
waitForEmpty(asyncWrapper, 100);
@@ -700,14 +721,15 @@ public class AsyncWrapperTest implements Serializable {
checkOutput(result, expectedOutput);
checkItemsInBuffer(asyncWrapper, 0);
- poolExecutor.shutdown();
+ assertEquals(0, fakeBagState.items.size());
}
// Test 12: testBufferWithCancellation
// Verifies actively cancelled elements are cleanly dropped from the buffer
during throttling.
@Test
public void testBufferWithCancellation() {
- BasicDofn dofn = new BasicDofn(10);
+ CountDownLatch gate = new CountDownLatch(1);
+ BasicDofn dofn = new BasicDofn(gate);
AsyncWrapper<String, String, String> asyncWrapper =
new AsyncWrapper<>(
dofn, 1, Duration.standardSeconds(5), null, null, null, null,
useThreadPool);
@@ -731,7 +753,9 @@ public class AsyncWrapperTest implements Serializable {
fakeTimer.getCurrentRelativeTime(), fakeBagState, fakeTimer);
checkOutput(result, Collections.emptyList());
assertEquals(1, fakeBagState.items.size());
+ checkItemsInBuffer(asyncWrapper, 1);
+ gate.countDown();
waitForEmpty(asyncWrapper);
result =