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 4c7f41373b [Fix][Connector-V2] Prevent TiDB CDC from advancing past 
slow regions (#11407)
4c7f41373b is described below

commit 4c7f41373b19a14a5f5e1ea625169784510aed2b
Author: Jast <[email protected]>
AuthorDate: Sat Aug 1 15:09:34 2026 +0800

    [Fix][Connector-V2] Prevent TiDB CDC from advancing past slow regions 
(#11407)
    
    Co-authored-by: zhangshenghang <[email protected]>
---
 .../cdc/tidb/source/reader/TiDBSourceReader.java   |  3 +-
 .../tidb/source/reader/TiDBSourceReaderTest.java   | 41 ++++++++++++++++++++++
 2 files changed, 43 insertions(+), 1 deletion(-)

diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/reader/TiDBSourceReader.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/reader/TiDBSourceReader.java
index 30b4b412ff..76d19fc5c5 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/reader/TiDBSourceReader.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/reader/TiDBSourceReader.java
@@ -264,7 +264,8 @@ public class TiDBSourceReader implements 
SourceReader<SeaTunnelRow, TiDBSourceSp
         }
         long pullCostMs = nanosToMillis(System.nanoTime() - pullStartNanos);
         long flushStartNanos = System.nanoTime();
-        resolvedTs = currentMaxResolvedTs;
+        // A split is safe to advance only after every TiKV region has reached 
the timestamp.
+        resolvedTs = cdcClient.getMinResolvedTs();
         int pendingCommitsBeforeFlush = commits.size();
         int committedEventsBeforeFlush = committedEvents.size();
         if (commits.size() > 0) {
diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/reader/TiDBSourceReaderTest.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/reader/TiDBSourceReaderTest.java
index b7915f0e77..57b4b36c5c 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/reader/TiDBSourceReaderTest.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/reader/TiDBSourceReaderTest.java
@@ -17,12 +17,22 @@
 
 package org.apache.seatunnel.connectors.seatunnel.cdc.tidb.source.reader;
 
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.source.SourceReader;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.connectors.cdc.base.option.StartupMode;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.tidb.source.config.TiDBSourceConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.cdc.tidb.source.split.TiDBSourceSplit;
+
 import org.junit.jupiter.api.Test;
+import org.tikv.cdc.CDCClient;
 import org.tikv.common.key.RowKey;
 import org.tikv.kvproto.Cdcpb;
+import org.tikv.kvproto.Coprocessor;
 
 import java.lang.reflect.Field;
 import java.lang.reflect.Method;
+import java.util.Map;
 import java.util.TreeMap;
 import java.util.concurrent.BlockingQueue;
 
@@ -30,6 +40,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
 
 class TiDBSourceReaderTest {
 
@@ -39,6 +51,27 @@ class TiDBSourceReaderTest {
     private static final long COMMIT_TS = 200L;
     private static final long RESOLVED_TS = 300L;
 
+    @Test
+    void shouldAdvanceSplitWithTheSlowestRegionResolvedTimestamp() throws 
Exception {
+        TiDBSourceConfig config =
+                
TiDBSourceConfig.builder().startupMode(StartupMode.LATEST).batchSize(1).build();
+        TiDBSourceReader reader =
+                new TiDBSourceReader(
+                        mock(SourceReader.Context.class), config, 
mock(CatalogTable.class));
+        TiDBSourceSplit split =
+                new TiDBSourceSplit(
+                        "database", "table", mock(Coprocessor.KeyRange.class), 
10L, null, true);
+        CDCClient cdcClient = mock(CDCClient.class);
+        when(cdcClient.get()).thenReturn(null);
+        when(cdcClient.getMinResolvedTs()).thenReturn(100L);
+        when(cdcClient.getMaxResolvedTs()).thenReturn(200L);
+        cdcClients(reader).put(split, cdcClient);
+
+        reader.captureStreamingEvents(split, mock(Collector.class));
+
+        assertEquals(100L, split.getResolvedTs());
+    }
+
     @Test
     void flushRowsShouldHoldCommitUntilMatchingPrewriteArrives() throws 
Exception {
         TiDBSourceReader reader = new TiDBSourceReader(null, null, null);
@@ -118,4 +151,12 @@ class TiDBSourceReaderTest {
         field.setAccessible(true);
         return (BlockingQueue<Cdcpb.Event.Row>) field.get(reader);
     }
+
+    @SuppressWarnings("unchecked")
+    private static Map<TiDBSourceSplit, CDCClient> cdcClients(TiDBSourceReader 
reader)
+            throws ReflectiveOperationException {
+        Field cacheField = 
TiDBSourceReader.class.getDeclaredField("cacheCDCClient");
+        cacheField.setAccessible(true);
+        return (Map<TiDBSourceSplit, CDCClient>) cacheField.get(reader);
+    }
 }

Reply via email to