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);
+ }
}