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

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new ac65b8f13e [Fix][Connector-V2] Fix Paimon stream enumerator checkpoint 
race (#11392)
ac65b8f13e is described below

commit ac65b8f13e864cccc4d5a0e82c92b4eb705966d3
Author: QuakeWang <[email protected]>
AuthorDate: Sat Aug 1 15:14:06 2026 +0800

    [Fix][Connector-V2] Fix Paimon stream enumerator checkpoint race (#11392)
    
    Signed-off-by: QuakeWang <[email protected]>
---
 .../seatunnel/paimon/source/PaimonSourceState.java |  5 +-
 .../source/enumerator/AbstractSplitEnumerator.java | 80 +++++++++---------
 .../PaimonStreamSourceSplitEnumerator.java         | 10 ++-
 .../PaimonStreamSourceSplitEnumeratorTest.java     | 94 ++++++++++++++++++++++
 4 files changed, 148 insertions(+), 41 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/PaimonSourceState.java
 
b/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/PaimonSourceState.java
index 73991cdd9a..8b2a8959fe 100644
--- 
a/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/PaimonSourceState.java
+++ 
b/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/PaimonSourceState.java
@@ -23,6 +23,7 @@ import java.io.Serializable;
 import java.util.Collections;
 import java.util.Deque;
 import java.util.HashMap;
+import java.util.LinkedList;
 import java.util.Map;
 
 /** Paimon connector source state, saves the splits has assigned to readers. */
@@ -38,14 +39,14 @@ public class PaimonSourceState implements Serializable {
 
     public PaimonSourceState(
             Deque<PaimonSourceSplit> assignedSplits, @Nullable Long 
currentSnapshotId) {
-        this.assignedSplits = assignedSplits;
+        this.assignedSplits = new LinkedList<>(assignedSplits);
         this.currentSnapshotId = currentSnapshotId;
         this.currentSnapshotIds = Collections.emptyMap();
     }
 
     public PaimonSourceState(
             Deque<PaimonSourceSplit> assignedSplits, Map<String, Long> 
currentSnapshotIds) {
-        this.assignedSplits = assignedSplits;
+        this.assignedSplits = new LinkedList<>(assignedSplits);
         this.currentSnapshotIds = new HashMap<>(currentSnapshotIds);
         this.currentSnapshotId = getSingleSnapshotId(this.currentSnapshotIds);
     }
diff --git 
a/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/AbstractSplitEnumerator.java
 
b/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/AbstractSplitEnumerator.java
index ed1356656a..c9125cfbb1 100644
--- 
a/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/AbstractSplitEnumerator.java
+++ 
b/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/AbstractSplitEnumerator.java
@@ -153,20 +153,26 @@ public abstract class AbstractSplitEnumerator
     @Override
     public void addSplitsBack(List<PaimonSourceSplit> splits, int subtaskId) {
         log.debug("Paimon Source Enumerator adds splits back: {}", splits);
-        this.pendingSplits.addAll(splits);
-        if (context.registeredReaders().contains(subtaskId)) {
-            assignSplits();
+        synchronized (stateLock) {
+            this.pendingSplits.addAll(splits);
+            if (context.registeredReaders().contains(subtaskId)) {
+                assignSplitsLocked();
+            }
         }
     }
 
     @Override
     public int currentUnassignedSplitSize() {
-        return pendingSplits.size();
+        synchronized (stateLock) {
+            return pendingSplits.size();
+        }
     }
 
     @Override
     public void registerReader(int subtaskId) {
-        readersAwaitingSplit.add(subtaskId);
+        synchronized (stateLock) {
+            readersAwaitingSplit.add(subtaskId);
+        }
     }
 
     @Override
@@ -183,11 +189,13 @@ public abstract class AbstractSplitEnumerator
         this.pendingSplits.addAll(newSplits);
     }
 
-    /**
-     * Method should be synchronized because {@link #handleSplitRequest} and 
{@link
-     * #processDiscoveredSplits} have thread conflicts.
-     */
-    protected synchronized void assignSplits() {
+    protected void assignSplits() {
+        synchronized (stateLock) {
+            assignSplitsLocked();
+        }
+    }
+
+    private void assignSplitsLocked() {
         Iterator<Integer> pendingReaderIterator = 
readersAwaitingSplit.iterator();
         while (pendingReaderIterator.hasNext()) {
             Integer pendingReader = pendingReaderIterator.next();
@@ -234,28 +242,26 @@ public abstract class AbstractSplitEnumerator
 
     // ------------------------------------------------------------------------
 
-    // This need to be synchronized because scan object is not thread safe. 
handleSplitRequest and
-    // CompletableFuture.supplyAsync will invoke this.
-    protected synchronized List<PlanWithNextSnapshotId> scanNextSnapshot() {
+    protected List<PlanWithNextSnapshotId> scanNextSnapshot() {
 
-        List<PlanWithNextSnapshotId> snapshotIds = Lists.newArrayList();
-        if (pendingSplits.size() >= splitMaxNum) {
+        synchronized (stateLock) {
+            List<PlanWithNextSnapshotId> snapshotIds = Lists.newArrayList();
+            if (pendingSplits.size() >= splitMaxNum) {
+                return snapshotIds;
+            }
+            tableScans.forEach(
+                    (tableId, tableScan) -> {
+                        TableScan.Plan plan = tableScan.plan();
+                        Long nextSnapshotId = null;
+                        if (tableScan instanceof StreamTableScan) {
+                            nextSnapshotId = ((StreamTableScan) 
tableScan).checkpoint();
+                        }
+                        snapshotIds.add(new PlanWithNextSnapshotId(tableId, 
plan, nextSnapshotId));
+                    });
             return snapshotIds;
         }
-        tableScans.forEach(
-                (tableId, tableScan) -> {
-                    TableScan.Plan plan = tableScan.plan();
-                    Long nextSnapshotId = null;
-                    if (tableScan instanceof StreamTableScan) {
-                        nextSnapshotId = ((StreamTableScan) 
tableScan).checkpoint();
-                    }
-                    snapshotIds.add(new PlanWithNextSnapshotId(tableId, plan, 
nextSnapshotId));
-                });
-        return snapshotIds;
     }
 
-    // This method could not be synchronized, because it runs in 
coordinatorThread, which will make
-    // it serializable execution.
     protected void processDiscoveredSplits(
             List<PlanWithNextSnapshotId> planWithNextSnapshotIds, Throwable 
error) {
         if (error != null) {
@@ -269,17 +275,19 @@ public abstract class AbstractSplitEnumerator
             return;
         }
 
-        for (PlanWithNextSnapshotId planWithNextSnapshotId : 
planWithNextSnapshotIds) {
-            nextSnapshotId = planWithNextSnapshotId.nextSnapshotId;
-            nextSnapshotIds.put(
-                    planWithNextSnapshotId.tableId, 
planWithNextSnapshotId.nextSnapshotId);
-            TableScan.Plan plan = planWithNextSnapshotId.plan;
-            if (plan.splits().isEmpty()) {
-                continue;
+        synchronized (stateLock) {
+            for (PlanWithNextSnapshotId planWithNextSnapshotId : 
planWithNextSnapshotIds) {
+                nextSnapshotId = planWithNextSnapshotId.nextSnapshotId;
+                nextSnapshotIds.put(
+                        planWithNextSnapshotId.tableId, 
planWithNextSnapshotId.nextSnapshotId);
+                TableScan.Plan plan = planWithNextSnapshotId.plan;
+                if (plan.splits().isEmpty()) {
+                    continue;
+                }
+                
addSplits(splitGenerator.createSplits(planWithNextSnapshotId.tableId, plan));
             }
-            
addSplits(splitGenerator.createSplits(planWithNextSnapshotId.tableId, plan));
+            assignSplitsLocked();
         }
-        assignSplits();
     }
 
     /** The result of scan. */
diff --git 
a/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/PaimonStreamSourceSplitEnumerator.java
 
b/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/PaimonStreamSourceSplitEnumerator.java
index 67706a8e51..9d63903ac1 100644
--- 
a/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/PaimonStreamSourceSplitEnumerator.java
+++ 
b/seatunnel-connectors-v2/connector-paimon/src/main/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/PaimonStreamSourceSplitEnumerator.java
@@ -67,9 +67,13 @@ public class PaimonStreamSourceSplitEnumerator extends 
AbstractSplitEnumerator {
 
     @Override
     public void handleSplitRequest(int subtaskId) {
-        readersAwaitingSplit.add(subtaskId);
-        assignSplits();
-        if (readersAwaitingSplit.contains(subtaskId)) {
+        boolean shouldLoadNewSplits;
+        synchronized (stateLock) {
+            readersAwaitingSplit.add(subtaskId);
+            assignSplits();
+            shouldLoadNewSplits = readersAwaitingSplit.contains(subtaskId);
+        }
+        if (shouldLoadNewSplits) {
             loadNewSplits();
         }
     }
diff --git 
a/seatunnel-connectors-v2/connector-paimon/src/test/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/PaimonStreamSourceSplitEnumeratorTest.java
 
b/seatunnel-connectors-v2/connector-paimon/src/test/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/PaimonStreamSourceSplitEnumeratorTest.java
index dce0b36919..8ee1e4cf0e 100644
--- 
a/seatunnel-connectors-v2/connector-paimon/src/test/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/PaimonStreamSourceSplitEnumeratorTest.java
+++ 
b/seatunnel-connectors-v2/connector-paimon/src/test/java/org/apache/seatunnel/connectors/seatunnel/paimon/source/enumerator/PaimonStreamSourceSplitEnumeratorTest.java
@@ -22,6 +22,7 @@ import 
org.apache.seatunnel.connectors.seatunnel.paimon.source.PaimonSourceSplit
 import 
org.apache.seatunnel.connectors.seatunnel.paimon.source.PaimonSourceState;
 
 import org.apache.paimon.table.source.ReadBuilder;
+import org.apache.paimon.table.source.Split;
 import org.apache.paimon.table.source.StreamTableScan;
 import org.apache.paimon.table.source.TableScan;
 
@@ -32,9 +33,17 @@ import java.util.HashSet;
 import java.util.LinkedHashMap;
 import java.util.LinkedList;
 import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.FutureTask;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
@@ -128,6 +137,91 @@ class PaimonStreamSourceSplitEnumeratorTest {
         verify(restoredScan).restore(10L);
     }
 
+    @Test
+    void 
shouldCheckpointSnapshotIdAndPendingSplitsAtomicallyDuringAsyncDiscovery()
+            throws Exception {
+        CountDownLatch discoveryBlockedAfterSnapshotIdUpdate = new 
CountDownLatch(1);
+        CountDownLatch continueDiscovery = new CountDownLatch(1);
+        CountDownLatch snapshotStateStarted = new CountDownLatch(1);
+        AtomicInteger splitsCalls = new AtomicInteger();
+        Split split =
+                new Split() {
+                    @Override
+                    public long rowCount() {
+                        return 1L;
+                    }
+
+                    @Override
+                    public String toString() {
+                        return "split-0";
+                    }
+                };
+        TableScan.Plan plan =
+                () -> {
+                    if (splitsCalls.incrementAndGet() == 1) {
+                        discoveryBlockedAfterSnapshotIdUpdate.countDown();
+                        try {
+                            assertTrue(continueDiscovery.await(30, 
TimeUnit.SECONDS));
+                        } catch (InterruptedException e) {
+                            Thread.currentThread().interrupt();
+                            throw new AssertionError(e);
+                        }
+                    }
+                    return Collections.singletonList(split);
+                };
+
+        PaimonStreamSourceSplitEnumerator enumerator =
+                new PaimonStreamSourceSplitEnumerator(
+                        context(), new LinkedList<>(), null, 
Collections.emptyMap(), 1) {
+                    @Override
+                    public PaimonSourceState snapshotState(long checkpointId) 
throws Exception {
+                        snapshotStateStarted.countDown();
+                        return super.snapshotState(checkpointId);
+                    }
+                };
+        AtomicReference<Throwable> discoveryError = new AtomicReference<>();
+        Thread discoveryThread =
+                new Thread(
+                        () -> {
+                            try {
+                                enumerator.processDiscoveredSplits(
+                                        Collections.singletonList(
+                                                new 
AbstractSplitEnumerator.PlanWithNextSnapshotId(
+                                                        "db.table_a", plan, 
10L)),
+                                        null);
+                            } catch (Throwable throwable) {
+                                discoveryError.set(throwable);
+                            }
+                        });
+        FutureTask<PaimonSourceState> snapshotState =
+                new FutureTask<>(() -> enumerator.snapshotState(1L));
+        Thread snapshotThread = new Thread(snapshotState);
+
+        try {
+            discoveryThread.start();
+            assertTrue(discoveryBlockedAfterSnapshotIdUpdate.await(10, 
TimeUnit.SECONDS));
+
+            snapshotThread.start();
+            assertTrue(snapshotStateStarted.await(10, TimeUnit.SECONDS));
+            assertThrows(
+                    TimeoutException.class, () -> snapshotState.get(200, 
TimeUnit.MILLISECONDS));
+
+            continueDiscovery.countDown();
+            PaimonSourceState state = snapshotState.get(10, TimeUnit.SECONDS);
+            discoveryThread.join(TimeUnit.SECONDS.toMillis(10));
+
+            assertFalse(discoveryThread.isAlive());
+            if (discoveryError.get() != null) {
+                throw new AssertionError(discoveryError.get());
+            }
+            assertEquals(10L, state.getCurrentSnapshotIds().get("db.table_a"));
+            assertEquals(1, state.getAssignedSplits().size());
+        } finally {
+            continueDiscovery.countDown();
+            enumerator.close();
+        }
+    }
+
     private static Map<String, ReadBuilder> readBuilders(
             StreamTableScan firstScan, StreamTableScan secondScan) {
         Map<String, ReadBuilder> readBuilders = new LinkedHashMap<>();

Reply via email to