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 =

Reply via email to