This is an automated email from the ASF dual-hosted git repository.

claudevdm 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 3e5cfce0b5a Fix race condition in test. (#40201)
3e5cfce0b5a is described below

commit 3e5cfce0b5a12d5c1281713997bc96ee23bb2d24
Author: claudevdm <[email protected]>
AuthorDate: Mon Sep 21 17:30:44 2026 -0400

    Fix race condition in test. (#40201)
    
    * init
    
    * add retry limit
---
 .../org/apache/beam/sdk/io/iceberg/BoundedAsyncTasks.java |  5 ++++-
 .../apache/beam/sdk/io/iceberg/BoundedAsyncTasksTest.java | 15 +++++++++------
 2 files changed, 13 insertions(+), 7 deletions(-)

diff --git 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasks.java
 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasks.java
index 3349c9e07fe..2fb76bf6ac8 100644
--- 
a/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasks.java
+++ 
b/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasks.java
@@ -25,6 +25,7 @@ import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
 import java.util.function.Consumer;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.util.concurrent.ThreadFactoryBuilder;
 
@@ -100,7 +101,9 @@ class BoundedAsyncTasks<T> {
     executor.shutdownNow();
   }
 
-  private void drainFinished(Consumer<T> onDone) throws Exception {
+  /** Delivers every task that has finished, in queue order, without blocking. 
*/
+  @VisibleForTesting
+  void drainFinished(Consumer<T> onDone) throws Exception {
     Iterator<Future<T>> iterator = active.iterator();
     while (iterator.hasNext()) {
       Future<T> future = iterator.next();
diff --git 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasksTest.java
 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasksTest.java
index 09471c1f414..478a616d58a 100644
--- 
a/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasksTest.java
+++ 
b/sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/BoundedAsyncTasksTest.java
@@ -88,14 +88,17 @@ public class BoundedAsyncTasksTest {
         },
         delivered::add);
 
-    // Wait until fast task has finished running in the background
+    // The latch fires inside the task body; the future is marked done a 
moment later, so a
+    // single drain can miss it. Drain until it shows up while "slow" is still 
blocked.
     fastTaskDone.await();
-
-    // Trigger a drain by submitting another task
     tasks.submit(() -> "noop", delivered::add);
-
-    // "fast" should have been drained while "slow" is still blocked
-    assertTrue(delivered.contains("fast"));
+    int retries = 0;
+    while (!delivered.contains("fast") && retries < 1000) {
+      Thread.sleep(1);
+      tasks.drainFinished(delivered::add);
+      retries++;
+    }
+    assertTrue("fast task was never drained", delivered.contains("fast"));
     assertFalse(delivered.contains("slow"));
 
     slowTaskHold.countDown();

Reply via email to