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