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