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

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new ec9a78174c [core] Fix flaky concurrent duplicate file discard test 
(#9285)
ec9a78174c is described below

commit ec9a78174cb6ffeea889fe4cb7113e9fea6a3364
Author: Arnav Balyan <[email protected]>
AuthorDate: Sun Aug 23 20:57:37 2026 +0530

    [core] Fix flaky concurrent duplicate file discard test (#9285)
---
 .../paimon/table/AppendOnlySimpleTableTest.java    | 60 ++++++++++------------
 1 file changed, 26 insertions(+), 34 deletions(-)

diff --git 
a/paimon-core/src/test/java/org/apache/paimon/table/AppendOnlySimpleTableTest.java
 
b/paimon-core/src/test/java/org/apache/paimon/table/AppendOnlySimpleTableTest.java
index c0c3f5e738..9d74f31ee0 100644
--- 
a/paimon-core/src/test/java/org/apache/paimon/table/AppendOnlySimpleTableTest.java
+++ 
b/paimon-core/src/test/java/org/apache/paimon/table/AppendOnlySimpleTableTest.java
@@ -100,9 +100,11 @@ import java.util.Optional;
 import java.util.PriorityQueue;
 import java.util.Random;
 import java.util.UUID;
+import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Consumer;
@@ -325,57 +327,47 @@ public class AppendOnlySimpleTableTest extends 
SimpleTableTestBase {
                             
options.set(CoreOptions.COMMIT_DISCARD_DUPLICATE_FILES, true);
                             options.set(CoreOptions.COMMIT_MAX_RETRIES, 50);
                             options.set(CoreOptions.COMMIT_MAX_RETRY_WAIT, 
Duration.ofMillis(100));
-                            // Keep all snapshots so concurrent expiry does 
not race readers.
-                            options.set(CoreOptions.SNAPSHOT_NUM_RETAINED_MIN, 
1000);
-                            options.set(CoreOptions.SNAPSHOT_NUM_RETAINED_MAX, 
1000);
                         });
         BatchWriteBuilder writeBuilder = table.newBatchWriteBuilder();
-        List<List<CommitMessage>> messages = new ArrayList<>();
-        for (int i = 0; i < 10; i++) {
-            try (BatchTableWrite write = writeBuilder.newWrite()) {
-                write.write(rowData(1, 10, 100L));
-                messages.add(write.prepareCommit());
-            }
+        List<CommitMessage> messages;
+        try (BatchTableWrite write = writeBuilder.newWrite()) {
+            write.write(rowData(1, 10, 100L));
+            messages = write.prepareCommit();
         }
-        int commitThreadNum = 10;
-        int commitsPerThread = 10;
-        Runnable asserter =
-                () -> {
-                    List<Split> splits = 
table.newReadBuilder().newScan().plan().splits();
-                    assertThat(splits.size()).isEqualTo(1);
-                    assertThat(splits.get(0).convertToRawFiles().get().size())
-                            .isLessThanOrEqualTo(messages.size());
-                };
 
+        int commitThreadNum = 5;
+        CountDownLatch ready = new CountDownLatch(commitThreadNum);
+        CountDownLatch start = new CountDownLatch(1);
         ExecutorService pool = Executors.newFixedThreadPool(commitThreadNum);
+        List<Future<?>> futures = new ArrayList<>();
         try {
-            List<Future<?>> futures = new ArrayList<>();
-            for (int thread = 0; thread < commitThreadNum; thread++) {
-                int threadId = thread;
+            for (int i = 0; i < commitThreadNum; i++) {
                 futures.add(
                         pool.submit(
                                 () -> {
-                                    for (int round = 0; round < 
commitsPerThread; round++) {
-                                        int messageIndex = (threadId + round) 
% messages.size();
-                                        try (BatchTableCommit commit = 
writeBuilder.newCommit()) {
-                                            
commit.commit(messages.get(messageIndex));
-                                        } catch (Exception e) {
-                                            throw new RuntimeException(
-                                                    String.format(
-                                                            "Failed to commit 
message %s in thread %s round %s.",
-                                                            messageIndex, 
threadId, round),
-                                                    e);
-                                        }
+                                    try (BatchTableCommit commit = 
writeBuilder.newCommit()) {
+                                        ready.countDown();
+                                        assertThat(start.await(10, 
TimeUnit.SECONDS)).isTrue();
+                                        commit.commit(messages);
+                                    } catch (Exception e) {
+                                        throw new RuntimeException(e);
                                     }
                                 }));
             }
+            assertThat(ready.await(10, TimeUnit.SECONDS)).isTrue();
+            start.countDown();
             for (Future<?> future : futures) {
                 future.get();
             }
         } finally {
-            pool.shutdownNow();
+            start.countDown();
+            pool.shutdown();
+            assertThat(pool.awaitTermination(1, TimeUnit.MINUTES)).isTrue();
         }
-        asserter.run();
+
+        List<Split> splits = table.newReadBuilder().newScan().plan().splits();
+        assertThat(splits).hasSize(1);
+        assertThat(splits.get(0).convertToRawFiles().get()).hasSize(1);
     }
 
     @Test

Reply via email to