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 d7dd5d5f5c [Fix][Connector-CDC] Implement MySQL CDC stop.mode =
"specific" (bounded read) (#11618)
d7dd5d5f5c is described below
commit d7dd5d5f5cb300855f49dfce9469c7a8f5a2117f
Author: Li Zhi <[email protected]>
AuthorDate: Fri Aug 14 12:15:40 2026 +0800
[Fix][Connector-CDC] Implement MySQL CDC stop.mode = "specific" (bounded
read) (#11618)
Co-authored-by: Claude Sonnet 5 <[email protected]>
Co-authored-by: David Zollo <[email protected]>
---
docs/en/connectors/source/MySQL-CDC.md | 47 +++
docs/zh/connectors/source/MySQL-CDC.md | 46 +++
.../connectors/cdc/base/source/offset/Offset.java | 10 +
.../source/reader/IncrementalSourceReader.java | 19 +-
.../external/IncrementalSourceStreamFetcher.java | 50 ++-
.../IncrementalSourceStreamFetcherTest.java | 98 +++++
.../mongodb/source/offset/ChangeStreamOffset.java | 5 +
.../source/offset/ChangeStreamOffsetTest.java} | 19 +-
.../cdc/mysql/source/offset/BinlogOffset.java | 22 +-
.../reader/fetch/binlog/MySqlBinlogFetchTask.java | 279 ++++++++++++-
.../cdc/mysql/source/offset/BinlogOffsetTest.java | 36 ++
.../cdc/oracle/source/offset/RedoLogOffset.java | 5 +
.../oracle/source/offset/RedoLogOffsetTest.java} | 18 +-
.../cdc/postgres/source/offset/LsnOffset.java | 5 +
.../cdc/postgres}/source/offset/LsnOffsetTest.java | 16 +-
.../cdc/sqlserver/source/offset/LsnOffset.java | 5 +
.../cdc/sqlserver/source/offset/LsnOffsetTest.java | 6 +
.../cdc/mysql/MysqlCDCStopModeSpecificIT.java | 434 +++++++++++++++++++++
.../resources/mysqlcdc_stop_mode_specific.conf | 54 +++
.../mysqlcdc_stop_mode_specific_timestamp.conf | 53 +++
20 files changed, 1170 insertions(+), 57 deletions(-)
diff --git a/docs/en/connectors/source/MySQL-CDC.md
b/docs/en/connectors/source/MySQL-CDC.md
index 266f7cc923..cd8c92a932 100644
--- a/docs/en/connectors/source/MySQL-CDC.md
+++ b/docs/en/connectors/source/MySQL-CDC.md
@@ -518,6 +518,53 @@ source {
}
```
+### Bounded Read: Stop at a Specific Binlog Offset
+
+Use `stop.mode = "specific"` to make the job a bounded read: it reads the
binlog between the
+startup offset (or startup timestamp) and the configured stop offset, then
terminates
+(`FINISHED`) instead of running forever.
+
+> **Note**: bounded-read termination is currently supported on the **Zeta**
engine only.
+> Flink and Spark engines do not support bounded incremental-split termination
yet.
+
+```hocon
+source {
+ MySQL-CDC {
+ server-id = 5654
+ username = "st_user_source"
+ password = "mysqlpw"
+ table-names = ["mysql_cdc.mysql_cdc_e2e_source_table"]
+ url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+ startup.mode = "specific"
+ startup.specific-offset.file = "mysql-bin.000001"
+ startup.specific-offset.pos = 154
+ stop.mode = "specific"
+ stop.specific-offset.file = "mysql-bin.000010"
+ stop.specific-offset.pos = 4096
+ }
+}
+```
+
+`stop.mode = "specific"` can also be combined with `startup.mode =
"timestamp"` to bound the
+read both by time and by binlog position:
+
+```hocon
+source {
+ MySQL-CDC {
+ server-id = 5654
+ username = "st_user_source"
+ password = "mysqlpw"
+ table-names = ["mysql_cdc.mysql_cdc_e2e_source_table"]
+ url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+ startup.mode = "timestamp"
+ startup.timestamp = 1716076800000
+ stop.mode = "specific"
+ stop.specific-offset.file = "mysql-bin.000010"
+ stop.specific-offset.pos = 4096
+ }
+}
+```
+
### Route Multiple Source Tables to JDBC
When one MySQL CDC source reads multiple tables, JDBC sink placeholders can
keep the original table name.
diff --git a/docs/zh/connectors/source/MySQL-CDC.md
b/docs/zh/connectors/source/MySQL-CDC.md
index c50393dfbe..a322cb91df 100644
--- a/docs/zh/connectors/source/MySQL-CDC.md
+++ b/docs/zh/connectors/source/MySQL-CDC.md
@@ -515,6 +515,52 @@ source {
}
```
+### 有界读取:在指定 Binlog 位置停止
+
+使用 `stop.mode = "specific"` 可以将作业变为有界读取:作业读取启动偏移量(或启动时间戳)
+与配置的停止偏移量之间的 binlog,然后自行终止(`FINISHED`),而不是一直运行下去。
+
+> **注意**:有界读取的终止行为目前仅在 **Zeta** 引擎上支持。
+> Flink 和 Spark 引擎暂不支持有界增量分片的终止。
+
+```hocon
+source {
+ MySQL-CDC {
+ server-id = 5654
+ username = "st_user_source"
+ password = "mysqlpw"
+ table-names = ["mysql_cdc.mysql_cdc_e2e_source_table"]
+ url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+ startup.mode = "specific"
+ startup.specific-offset.file = "mysql-bin.000001"
+ startup.specific-offset.pos = 154
+ stop.mode = "specific"
+ stop.specific-offset.file = "mysql-bin.000010"
+ stop.specific-offset.pos = 4096
+ }
+}
+```
+
+`stop.mode = "specific"` 也可以与 `startup.mode = "timestamp"` 组合使用,同时按时间和
+binlog 位置限定读取范围:
+
+```hocon
+source {
+ MySQL-CDC {
+ server-id = 5654
+ username = "st_user_source"
+ password = "mysqlpw"
+ table-names = ["mysql_cdc.mysql_cdc_e2e_source_table"]
+ url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+ startup.mode = "timestamp"
+ startup.timestamp = 1716076800000
+ stop.mode = "specific"
+ stop.specific-offset.file = "mysql-bin.000010"
+ stop.specific-offset.pos = 4096
+ }
+}
+```
+
### 多表读取后写入 JDBC
当一个 MySQL CDC source 读取多张表时,JDBC sink 可以使用占位符保留原始表名。
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/offset/Offset.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/offset/Offset.java
index 75ccfff2ad..7f2bfa0790 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/offset/Offset.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/offset/Offset.java
@@ -75,6 +75,16 @@ public abstract class Offset implements Comparable<Offset>,
Serializable {
return this.compareTo(that) > 0;
}
+ /**
+ * Check if this offset represents a "never stop" sentinel value.
Subclasses should override
+ * this method if they use a special sentinel value for unbounded reads.
+ *
+ * @return true if this is a never-stop sentinel value, false otherwise
+ */
+ public boolean isNeverStop() {
+ return false;
+ }
+
@Override
public String toString() {
return offset.toString();
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java
index c11cdd9518..a8899d3065 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java
@@ -167,13 +167,18 @@ public class IncrementalSourceReader<T, C extends
SourceConfig>
protected void onSplitFinished(Map<String, SourceSplitStateBase>
finishedSplitIds) {
for (SourceSplitStateBase splitState : finishedSplitIds.values()) {
SourceSplitBase sourceSplit = splitState.toSourceSplit();
- checkState(
- sourceSplit.isSnapshotSplit()
- &&
sourceSplit.asSnapshotSplit().isSnapshotReadFinished(),
- String.format(
- "Only snapshot split could finish, but the actual
split is incremental split %s",
- sourceSplit));
- finishedUnackedSplits.put(sourceSplit.splitId(),
sourceSplit.asSnapshotSplit());
+ if (sourceSplit.isSnapshotSplit()) {
+ checkState(
+ sourceSplit.asSnapshotSplit().isSnapshotReadFinished(),
+ String.format(
+ "Snapshot split should be finished, but the
actual split is %s",
+ sourceSplit));
+ finishedUnackedSplits.put(sourceSplit.splitId(),
sourceSplit.asSnapshotSplit());
+ } else {
+ log.info(
+ "Incremental split {} has finished (bounded read
completed).",
+ sourceSplit.splitId());
+ }
}
reportFinishedSnapshotSplitsIfNeed();
context.sendSplitRequest();
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/external/IncrementalSourceStreamFetcher.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/external/IncrementalSourceStreamFetcher.java
index ec960f4cb2..d724ba502b 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/external/IncrementalSourceStreamFetcher.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/external/IncrementalSourceStreamFetcher.java
@@ -64,6 +64,16 @@ public class IncrementalSourceStreamFetcher implements
Fetcher<SourceRecords, So
private final Set<TableId> pureBinlogPhaseTables;
private volatile ChangeEventQueue<DataChangeEvent> queue;
private volatile Throwable readException;
+ private volatile boolean taskStarted = false;
+ // Set synchronously in submitTask(), before the fetch task is handed to
the background
+ // thread, and cleared in a finally block once execute() returns/throws.
isFinished() must
+ // not rely on polling streamFetchTask.isRunning() alone: that flag is
owned by the fetch
+ // task and is only flipped to true from inside its own execute() method,
which the
+ // (single-threaded) poller can observe as not-yet-run. Without this flag
isFinished() could
+ // report the split as finished before the fetch task ever ran, silently
truncating the whole
+ // incremental/bounded read (see bug-002/bug-008). Setting it inside the
background thread's
+ // lambda instead of synchronously in submitTask() re-opens this exact
race.
+ private volatile boolean executing = false;
private FetchTask<SourceSplitBase> streamFetchTask;
@@ -97,6 +107,13 @@ public class IncrementalSourceStreamFetcher implements
Fetcher<SourceRecords, So
configureFilter();
taskContext.configure(currentIncrementalSplit);
this.queue = taskContext.getQueue();
+ // Set synchronously, before the task is handed to the background
thread, so there is no
+ // window in which the (single-threaded) poller can observe
taskStarted=true and
+ // executing=false before the background thread has actually started
running. Setting
+ // this flag from inside the submitted lambda left a gap between
"taskStarted = true" and
+ // the next statement (a log call) that the poller reliably hit on
every run.
+ taskStarted = true;
+ executing = true;
executorService.submit(
() -> {
try {
@@ -112,13 +129,20 @@ public class IncrementalSourceStreamFetcher implements
Fetcher<SourceRecords, So
currentIncrementalSplit),
e);
readException = e;
+ } finally {
+ executing = false;
}
});
}
@Override
public boolean isFinished() {
- return currentIncrementalSplit == null || !streamFetchTask.isRunning();
+ // Never report finished while the background thread is still inside
+ // streamFetchTask.execute() -- streamFetchTask.isRunning() is owned
by the fetch task
+ // and may not have flipped to true yet even though we are already
executing it.
+ return taskStarted
+ && !executing
+ && (currentIncrementalSplit == null ||
!streamFetchTask.isRunning());
}
@Override
@@ -126,8 +150,11 @@ public class IncrementalSourceStreamFetcher implements
Fetcher<SourceRecords, So
throws InterruptedException, SeaTunnelException {
checkReadException();
+ // Always drain the queue first, so the last batch produced before the
bounded
+ // reader stopped is not dropped. Only after the queue is drained AND
the split
+ // is finished do we signal split completion by returning null.
Iterator<SourceRecords> sourceRecordsIterator =
Collections.emptyIterator();
- if (streamFetchTask.isRunning()) {
+ if (streamFetchTask.isRunning() || isBoundedReadFinished()) {
List<DataChangeEvent> batch = queue.poll();
if (!batch.isEmpty()) {
if (schemaChangeResolver != null) {
@@ -137,12 +164,29 @@ public class IncrementalSourceStreamFetcher implements
Fetcher<SourceRecords, So
}
}
}
+
+ // If the fetch task is finished and this is a bounded read (stop.mode
= "specific"),
+ // return null to signal split completion. This is important for
bounded reads
+ // to properly terminate the task. For unbounded reads (stop.mode =
"never"),
+ // we should never return null even if the task is not running due to
errors,
+ // because that would incorrectly mark the job as FINISHED.
+ if (isBoundedReadFinished() && !sourceRecordsIterator.hasNext()) {
+ log.info("Bounded read completed, returning null to signal split
completion");
+ return null;
+ }
return sourceRecordsIterator;
}
+ private boolean isBoundedReadFinished() {
+ return isFinished()
+ && currentIncrementalSplit != null
+ && currentIncrementalSplit.getStopOffset() != null
+ && !currentIncrementalSplit.getStopOffset().isNeverStop();
+ }
+
private Iterator<SourceRecords> splitNormalStream(List<DataChangeEvent>
batchEvents) {
List<SourceRecord> sourceRecords = new ArrayList<>();
- if (streamFetchTask.isRunning()) {
+ if (streamFetchTask.isRunning() || isBoundedReadFinished()) {
for (DataChangeEvent event : batchEvents) {
if (shouldEmit(event.getRecord())) {
sourceRecords.add(event.getRecord());
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/external/IncrementalSourceStreamFetcherTest.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/external/IncrementalSourceStreamFetcherTest.java
index 23906ae6f4..265f2a433d 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/external/IncrementalSourceStreamFetcherTest.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/external/IncrementalSourceStreamFetcherTest.java
@@ -18,7 +18,10 @@
package org.apache.seatunnel.connectors.cdc.base.source.reader.external;
import org.apache.seatunnel.connectors.cdc.base.schema.SchemaChangeResolver;
+import org.apache.seatunnel.connectors.cdc.base.source.offset.Offset;
+import org.apache.seatunnel.connectors.cdc.base.source.split.IncrementalSplit;
import org.apache.seatunnel.connectors.cdc.base.source.split.SourceRecords;
+import org.apache.seatunnel.connectors.cdc.base.source.split.SourceSplitBase;
import
org.apache.seatunnel.connectors.cdc.base.source.split.wartermark.WatermarkEvent;
import org.apache.seatunnel.connectors.cdc.base.utils.SourceRecordUtils;
@@ -34,6 +37,7 @@ import org.mockito.stubbing.Answer;
import io.debezium.config.CommonConnectorConfig;
import io.debezium.config.Configuration;
import io.debezium.connector.SourceInfoStructMaker;
+import io.debezium.connector.base.ChangeEventQueue;
import io.debezium.data.Envelope;
import io.debezium.heartbeat.Heartbeat;
import io.debezium.heartbeat.HeartbeatFactory;
@@ -43,7 +47,9 @@ import io.debezium.relational.TableId;
import io.debezium.schema.TopicSelector;
import io.debezium.util.SchemaNameAdjuster;
+import java.lang.reflect.Field;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Collections;
import java.util.Iterator;
import java.util.List;
@@ -403,6 +409,98 @@ public class IncrementalSourceStreamFetcherTest {
return eventRef.get();
}
+ @Test
+ public void testPollSplitRecordsDrainsQueueBeforeBoundedCompletion()
throws Exception {
+ IncrementalSourceStreamFetcher fetcher = createFetcher();
+
+ // A bounded incremental split: the stop offset is set and is not the
+ // never-stop sentinel, so the split can finish.
+ Offset stopOffset = mock(Offset.class);
+ IncrementalSplit incrementalSplit =
+ new IncrementalSplit(
+ "incremental-0",
+ Collections.emptyList(),
+ null,
+ stopOffset,
+ Collections.emptyList());
+
+ // The fetch task has already finished (isRunning() == false),
simulating the
+ // bounded reader reaching the stop offset and stopping the binlog
context.
+ FetchTask<SourceSplitBase> fetchTask = mock(FetchTask.class);
+ when(fetchTask.isRunning()).thenReturn(false);
+ when(fetchTask.getSplit()).thenReturn(incrementalSplit);
+
+ // The queue still holds one last batch that must be drained before
completion.
+ ChangeEventQueue<DataChangeEvent> queue = mock(ChangeEventQueue.class);
+ when(queue.poll())
+ .thenReturn(
+ Arrays.asList(
+ new DataChangeEvent(createDataEvent()),
+ new DataChangeEvent(createDataEvent())),
+ Collections.emptyList());
+
+ setField(fetcher, "queue", queue);
+ setField(fetcher, "streamFetchTask", fetchTask);
+ setField(fetcher, "currentIncrementalSplit", incrementalSplit);
+ setField(fetcher, "taskStarted", true);
+ setField(fetcher, "executing", false);
+
+ // First poll: the last queued batch must be returned, not null.
+ Iterator<SourceRecords> first = fetcher.pollSplitRecords();
+ Assertions.assertNotNull(first);
+ Assertions.assertTrue(first.hasNext());
+ SourceRecords records = first.next();
+ Assertions.assertEquals(2, records.getSourceRecordList().size());
+
+ // Second poll: the queue is drained, so the bounded split completion
is
+ // signaled by returning null.
+ Assertions.assertNull(fetcher.pollSplitRecords());
+ }
+
+ @Test
+ public void testNeverStopSentinelIsNotSignaledAsBoundedCompletion() throws
Exception {
+ IncrementalSourceStreamFetcher fetcher = createFetcher();
+
+ // A real never-stop sentinel offset (e.g.
LsnOffset.NO_STOPPING_OFFSET or
+ // RedoLogOffset.NO_STOPPING_OFFSET) reports isNeverStop() == true.
The fetch task
+ // has finished without an error, but the split must NOT be marked
completed: an
+ // unbounded connector must never be turned into a bounded one by the
shared reader.
+ Offset stopOffset = mock(Offset.class);
+ when(stopOffset.isNeverStop()).thenReturn(true);
+ IncrementalSplit incrementalSplit =
+ new IncrementalSplit(
+ "incremental-0",
+ Collections.emptyList(),
+ null,
+ stopOffset,
+ Collections.emptyList());
+
+ FetchTask<SourceSplitBase> fetchTask = mock(FetchTask.class);
+ when(fetchTask.isRunning()).thenReturn(false);
+ when(fetchTask.getSplit()).thenReturn(incrementalSplit);
+
+ ChangeEventQueue<DataChangeEvent> queue = mock(ChangeEventQueue.class);
+ when(queue.poll()).thenReturn(Collections.emptyList());
+
+ setField(fetcher, "queue", queue);
+ setField(fetcher, "streamFetchTask", fetchTask);
+ setField(fetcher, "currentIncrementalSplit", incrementalSplit);
+ setField(fetcher, "taskStarted", true);
+ setField(fetcher, "executing", false);
+
+ // Even though the task is finished, an unbounded sentinel must not
cause a
+ // bounded-completion signal (null); an empty iterator is returned
instead.
+ Iterator<SourceRecords> result = fetcher.pollSplitRecords();
+ Assertions.assertNotNull(result);
+ Assertions.assertFalse(result.hasNext());
+ }
+
+ private static void setField(Object target, String fieldName, Object
value) throws Exception {
+ Field field =
IncrementalSourceStreamFetcher.class.getDeclaredField(fieldName);
+ field.setAccessible(true);
+ field.set(target, value);
+ }
+
static IncrementalSourceStreamFetcher createFetcher() {
SchemaChangeResolver schemaChangeResolver =
mock(SchemaChangeResolver.class);
when(schemaChangeResolver.support(any()))
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/offset/ChangeStreamOffset.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/offset/ChangeStreamOffset.java
index 5ec058ef51..e19633a59f 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/offset/ChangeStreamOffset.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/offset/ChangeStreamOffset.java
@@ -92,6 +92,11 @@ public class ChangeStreamOffset extends Offset {
return offset.equals(that.offset);
}
+ @Override
+ public boolean isNeverStop() {
+ return NO_STOPPING_OFFSET.equals(this);
+ }
+
@Override
public int compareTo(Offset offset) {
if (offset == null) {
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/offset/ChangeStreamOffsetTest.java
similarity index 61%
copy from
seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
copy to
seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/offset/ChangeStreamOffsetTest.java
index 7c62cb46c5..e485e9a48f 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mongodb/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mongodb/source/offset/ChangeStreamOffsetTest.java
@@ -15,24 +15,17 @@
* limitations under the License.
*/
-package org.apache.seatunnel.connectors.seatunnel.cdc.sqlserver.source.offset;
+package org.apache.seatunnel.connectors.seatunnel.cdc.mongodb.source.offset;
+import org.bson.BsonTimestamp;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
-import io.debezium.connector.sqlserver.Lsn;
-
-class LsnOffsetTest {
+class ChangeStreamOffsetTest {
@Test
- void testInitialOffsetRepresentsNoLsn() {
- LsnOffset initial = LsnOffset.INITIAL_OFFSET;
-
- // no LSN keys should be present in the offset map
- Assertions.assertTrue(initial.getOffset().isEmpty());
-
- // commit LSN resolved from the empty map should be Debezium's NULL LSN
- Lsn commitLsn = initial.getCommitLsn();
- Assertions.assertFalse(commitLsn.isAvailable());
+ void testNoStoppingOffsetIsNeverStop() {
+
Assertions.assertTrue(ChangeStreamOffset.NO_STOPPING_OFFSET.isNeverStop());
+ Assertions.assertFalse(new ChangeStreamOffset(new BsonTimestamp(0,
0)).isNeverStop());
}
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/offset/BinlogOffset.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/offset/BinlogOffset.java
index f81029daf7..c8da295228 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/offset/BinlogOffset.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/offset/BinlogOffset.java
@@ -121,6 +121,11 @@ public class BinlogOffset extends Offset {
return longOffsetValue(offset, SERVER_ID_KEY);
}
+ @Override
+ public boolean isNeverStop() {
+ return NO_STOPPING_OFFSET.equals(this);
+ }
+
/**
* This method is inspired by {@link
io.debezium.relational.history.HistoryRecordComparator}.
*/
@@ -175,10 +180,19 @@ public class BinlogOffset extends Offset {
// again we know that this offset not having GTIDs is before the
target offset ...
return -1;
} else if (StringUtils.isNotEmpty(gtidSetStr)) {
- // This offset has a GTID but the target offset does not, so per
the previous paragraph
- // we
- // assume that previous
- // is not at or before ...
+ // This offset has a GTID but the target offset does not, so we
need to compare
+ // by binlog filename and position instead
+ String thisFilename = this.getFilename();
+ String thatFilename = that.getFilename();
+ if (StringUtils.isNotEmpty(thisFilename) &&
StringUtils.isNotEmpty(thatFilename)) {
+ int filenameCompare =
thisFilename.compareToIgnoreCase(thatFilename);
+ if (filenameCompare != 0) {
+ return filenameCompare;
+ }
+ // Same binlog file, compare by position
+ return Long.compare(this.getPosition(), that.getPosition());
+ }
+ // Cannot compare, assume this is newer
return 1;
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/reader/fetch/binlog/MySqlBinlogFetchTask.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/reader/fetch/binlog/MySqlBinlogFetchTask.java
index f0ade84e4b..d02bf8479c 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/reader/fetch/binlog/MySqlBinlogFetchTask.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/reader/fetch/binlog/MySqlBinlogFetchTask.java
@@ -73,13 +73,52 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
StartupConfig startupConfig =
sourceFetchContext.getSourceConfig().getStartupConfig();
StartupMode startupMode = startupConfig.getStartupMode();
+ // Check if we need bounded read (stop at specific position or
timestamp)
+ boolean isBoundedRead =
!NO_STOPPING_OFFSET.equals(split.getStopOffset());
+
if (shouldFilterByTimestamp(startupMode, split.getStartupOffset())) {
+ if (isBoundedRead) {
+ // Both timestamp startup and a specific stop offset: apply
the timestamp filter
+ // inside the bounded reader so the job still terminates at
the stop offset
+ // (instead of running forever).
+ log.info(
+ "Starting MySQL binlog reader with timestamp filter {}
and bounded read, stop offset: {}",
+ startupConfig.getTimestamp(),
+ split.getStopOffset());
+ mySqlStreamingChangeEventSource =
+ new BoundedMySqlStreamingChangeEventSource(
+ sourceFetchContext.getDbzConnectorConfig(),
+ sourceFetchContext.getConnection(),
+ sourceFetchContext.getDispatcher(),
+ sourceFetchContext.getErrorHandler(),
+ Clock.SYSTEM,
+ sourceFetchContext.getTaskContext(),
+
sourceFetchContext.getStreamingChangeEventSourceMetrics(),
+ split,
+ startupConfig.getTimestamp());
+ } else {
+ log.info(
+ "Starting MySQL binlog reader,with timestamp filter
{}",
+ startupConfig.getTimestamp());
+
+ mySqlStreamingChangeEventSource =
+ new TimestampFilterMySqlStreamingChangeEventSource(
+ sourceFetchContext.getDbzConnectorConfig(),
+ sourceFetchContext.getConnection(),
+ sourceFetchContext.getDispatcher(),
+ sourceFetchContext.getErrorHandler(),
+ Clock.SYSTEM,
+ sourceFetchContext.getTaskContext(),
+
sourceFetchContext.getStreamingChangeEventSourceMetrics(),
+ startupConfig.getTimestamp());
+ }
+ } else if (isBoundedRead) {
+ // Bounded read: stop at specific offset
log.info(
- "Starting MySQL binlog reader,with timestamp filter {}",
- startupConfig.getTimestamp());
-
+ "Starting MySQL binlog reader with bounded read, stop
offset: {}",
+ split.getStopOffset());
mySqlStreamingChangeEventSource =
- new TimestampFilterMySqlStreamingChangeEventSource(
+ new BoundedMySqlStreamingChangeEventSource(
sourceFetchContext.getDbzConnectorConfig(),
sourceFetchContext.getConnection(),
sourceFetchContext.getDispatcher(),
@@ -87,7 +126,8 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
Clock.SYSTEM,
sourceFetchContext.getTaskContext(),
sourceFetchContext.getStreamingChangeEventSourceMetrics(),
- startupConfig.getTimestamp());
+ split,
+ null);
} else {
mySqlStreamingChangeEventSource =
new MySqlStreamingChangeEventSource(
@@ -123,6 +163,7 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
changeEventSourceContext,
sourceFetchContext.getPartition(),
sourceFetchContext.getOffsetContext());
+ taskRunning = false;
}
@Override
@@ -158,6 +199,10 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
private final JdbcSourceEventDispatcher<MySqlPartition> dispatcher;
private final ErrorHandler errorHandler;
private ChangeEventSourceContext context;
+ private long eventCount = 0;
+ private long lastLogTime = System.currentTimeMillis();
+ private static final long LOG_INTERVAL_MS = 10000;
+ private BinlogOffset lastLoggedOffset = null;
public MySqlBinlogSplitReadTask(
MySqlConnectorConfig connectorConfig,
@@ -196,6 +241,8 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
protected void handleEvent(
MySqlPartition partition, MySqlOffsetContext offsetContext,
Event event) {
super.handleEvent(partition, offsetContext, event);
+ eventCount++;
+ logBinlogProgress(offsetContext);
// check do we need to stop for fetch binlog for snapshot split.
if (isBoundedRead()) {
final BinlogOffset currentBinlogOffset =
@@ -221,6 +268,24 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
}
}
+ private void logBinlogProgress(MySqlOffsetContext offsetContext) {
+ long currentTime = System.currentTimeMillis();
+ if (currentTime - lastLogTime >= LOG_INTERVAL_MS) {
+ BinlogOffset currentOffset =
getBinlogPosition(offsetContext.getOffset());
+ if (lastLoggedOffset == null
+ ||
!currentOffset.getFilename().equals(lastLoggedOffset.getFilename())
+ || currentOffset.getPosition() !=
lastLoggedOffset.getPosition()) {
+ LOG.info(
+ "MySQL CDC binlog progress - file: {}, position:
{}, events processed: {}",
+ currentOffset.getFilename(),
+ currentOffset.getPosition(),
+ eventCount);
+ lastLoggedOffset = currentOffset;
+ }
+ lastLogTime = currentTime;
+ }
+ }
+
private boolean isBoundedRead() {
return !NO_STOPPING_OFFSET.equals(binlogSplit.getStopOffset());
}
@@ -243,6 +308,9 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
private long logTimestamp;
private boolean loggedWaitingMessage;
private final long LOG_INTERVAL_MS = 10000;
+ private long eventCount = 0;
+ private long lastLogTime = System.currentTimeMillis();
+ private BinlogOffset lastLoggedOffset = null;
public TimestampFilterMySqlStreamingChangeEventSource(
MySqlConnectorConfig connectorConfig,
@@ -272,6 +340,9 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
return;
}
+ eventCount++;
+ logBinlogProgress(offsetContext);
+
long eventTs = event.getHeader().getTimestamp();
if (eventTs == 0 || targetTimestamp == null || targetTimestamp ==
0) {
super.handleEvent(partition, offsetContext, event);
@@ -295,6 +366,122 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
super.handleEvent(partition, offsetContext, event);
}
+ private void logBinlogProgress(MySqlOffsetContext offsetContext) {
+ long currentTime = System.currentTimeMillis();
+ if (currentTime - lastLogTime >= LOG_INTERVAL_MS) {
+ BinlogOffset currentOffset =
+
MySqlBinlogSplitReadTask.getBinlogPosition(offsetContext.getOffset());
+ if (lastLoggedOffset == null
+ ||
!currentOffset.getFilename().equals(lastLoggedOffset.getFilename())
+ || currentOffset.getPosition() !=
lastLoggedOffset.getPosition()) {
+ log.info(
+ "MySQL CDC binlog progress - file: {}, position:
{}, events processed: {}",
+ currentOffset.getFilename(),
+ currentOffset.getPosition(),
+ eventCount);
+ lastLoggedOffset = currentOffset;
+ }
+ lastLogTime = currentTime;
+ }
+ }
+
+ private void updateOffsetPosition(
+ MySqlOffsetContext offsetContext, EventHeader eventHeader) {
+ try {
+ if (eventHeader instanceof EventHeaderV4) {
+ EventHeaderV4 headerV4 = (EventHeaderV4) eventHeader;
+ offsetContext.setEventPosition(
+ headerV4.getPosition(), headerV4.getEventLength());
+ }
+ offsetContext.setBinlogServerId(eventHeader.getServerId());
+ offsetContext.completeEvent();
+ } catch (Exception e) {
+ log.warn("Failed to update offset for skipped event: {}",
e.getMessage());
+ }
+ }
+ }
+
+ /**
+ * A bounded MySQL streaming change event source that stops at a specific
offset. Used when
+ * stop.mode = "specific" is configured.
+ */
+ private class BoundedMySqlStreamingChangeEventSource extends
MySqlStreamingChangeEventSource {
+
+ private final IncrementalSplit binlogSplit;
+ private final JdbcSourceEventDispatcher<MySqlPartition> dispatcher;
+ private final ErrorHandler errorHandler;
+ private final Long targetTimestamp;
+ private BoundedBinlogChangeEventSourceContext boundedContext;
+ private long eventCount = 0;
+ private long lastLogTime = System.currentTimeMillis();
+ private static final long LOG_INTERVAL_MS = 10000;
+ private BinlogOffset lastLoggedOffset = null;
+ private boolean loggedWaitingMessage;
+ private long logTimestamp;
+
+ public BoundedMySqlStreamingChangeEventSource(
+ MySqlConnectorConfig connectorConfig,
+ MySqlConnection connection,
+ JdbcSourceEventDispatcher<MySqlPartition> dispatcher,
+ ErrorHandler errorHandler,
+ Clock clock,
+ MySqlTaskContext taskContext,
+ MySqlStreamingChangeEventSourceMetrics metrics,
+ IncrementalSplit binlogSplit,
+ Long targetTimestamp) {
+ super(
+ connectorConfig,
+ connection,
+ dispatcher,
+ errorHandler,
+ clock,
+ taskContext,
+ metrics);
+ this.binlogSplit = binlogSplit;
+ this.dispatcher = dispatcher;
+ this.errorHandler = errorHandler;
+ this.targetTimestamp = targetTimestamp;
+ }
+
+ @Override
+ protected void handleEvent(
+ MySqlPartition partition, MySqlOffsetContext offsetContext,
Event event) {
+ try {
+ if (shouldSkipByTimestamp(offsetContext, event)) {
+ return;
+ }
+ super.handleEvent(partition, offsetContext, event);
+ } finally {
+ eventCount++;
+ logBinlogProgress(offsetContext);
+ checkStopOffset(partition, offsetContext);
+ }
+ }
+
+ private boolean shouldSkipByTimestamp(MySqlOffsetContext
offsetContext, Event event) {
+ if (event == null || targetTimestamp == null || targetTimestamp ==
0) {
+ return false;
+ }
+ long eventTs = event.getHeader().getTimestamp();
+ if (eventTs == 0) {
+ return false;
+ }
+ boolean shouldSkip = eventTs < targetTimestamp;
+ if (shouldSkip) {
+ if (!loggedWaitingMessage) {
+ log.info(
+ "skip binlog, currentTime:{}, filterTime:{}",
eventTs, targetTimestamp);
+ loggedWaitingMessage = true;
+ logTimestamp = eventTs;
+ }
+ if (eventTs - logTimestamp >= LOG_INTERVAL_MS) {
+ loggedWaitingMessage = false;
+ }
+ updateOffsetPosition(offsetContext, event.getHeader());
+ }
+ return shouldSkip;
+ }
+
private void updateOffsetPosition(
MySqlOffsetContext offsetContext, EventHeader eventHeader) {
try {
@@ -309,6 +496,88 @@ public class MySqlBinlogFetchTask implements
FetchTask<SourceSplitBase> {
log.warn("Failed to update offset for skipped event: {}",
e.getMessage());
}
}
+
+ private void logBinlogProgress(MySqlOffsetContext offsetContext) {
+ long currentTime = System.currentTimeMillis();
+ if (currentTime - lastLogTime >= LOG_INTERVAL_MS) {
+ BinlogOffset currentOffset =
+
MySqlBinlogSplitReadTask.getBinlogPosition(offsetContext.getOffset());
+ if (lastLoggedOffset == null
+ ||
!currentOffset.getFilename().equals(lastLoggedOffset.getFilename())
+ || currentOffset.getPosition() !=
lastLoggedOffset.getPosition()) {
+ log.info(
+ "MySQL CDC binlog progress - file: {}, position:
{}, events processed: {}",
+ currentOffset.getFilename(),
+ currentOffset.getPosition(),
+ eventCount);
+ lastLoggedOffset = currentOffset;
+ }
+ lastLogTime = currentTime;
+ }
+ }
+
+ private void checkStopOffset(MySqlPartition partition,
MySqlOffsetContext offsetContext) {
+ final BinlogOffset currentBinlogOffset =
+
MySqlBinlogSplitReadTask.getBinlogPosition(offsetContext.getOffset());
+
+ if (currentBinlogOffset.isAtOrAfter(binlogSplit.getStopOffset())) {
+ log.info(
+ "Reached stop offset {} at current position {}.
Stopping binlog reader.",
+ binlogSplit.getStopOffset(),
+ currentBinlogOffset);
+
+ // Send end watermark event
+ try {
+ dispatcher.dispatchWatermarkEvent(
+ partition.getSourcePartition(),
+ binlogSplit,
+ currentBinlogOffset,
+ WatermarkKind.END);
+ } catch (InterruptedException e) {
+ log.error("Error sending binlog end watermark event", e);
+ errorHandler.setProducerThrowable(
+ new DebeziumException("Error processing binlog end
event", e));
+ }
+
+ // Stop the task
+ if (boundedContext != null) {
+ boundedContext.finished();
+ }
+ }
+ }
+
+ @Override
+ public void execute(
+ ChangeEventSourceContext context,
+ MySqlPartition partition,
+ MySqlOffsetContext offsetContext)
+ throws InterruptedException {
+ // Wrap the context to allow stopping
+ this.boundedContext = new
BoundedBinlogChangeEventSourceContext(context);
+ super.execute(boundedContext, partition, offsetContext);
+ }
+ }
+
+ /** A context wrapper that allows stopping the binlog reader. */
+ private class BoundedBinlogChangeEventSourceContext
+ implements ChangeEventSource.ChangeEventSourceContext {
+
+ private final ChangeEventSource.ChangeEventSourceContext delegate;
+ private volatile boolean running = true;
+
+ public BoundedBinlogChangeEventSourceContext(
+ ChangeEventSource.ChangeEventSourceContext delegate) {
+ this.delegate = delegate;
+ }
+
+ @Override
+ public boolean isRunning() {
+ return running && delegate.isRunning();
+ }
+
+ public void finished() {
+ running = false;
+ }
}
private class BinlogSplitChangeEventSourceContext
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/offset/BinlogOffsetTest.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/offset/BinlogOffsetTest.java
index 4889308644..54b6165561 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/offset/BinlogOffsetTest.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-mysql/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/source/offset/BinlogOffsetTest.java
@@ -85,4 +85,40 @@ public class BinlogOffsetTest {
Assertions.assertEquals(
0,
BinlogOffset.NO_STOPPING_OFFSET.compareTo(BinlogOffset.NO_STOPPING_OFFSET));
}
+
+ @Test
+ public void
testCompareToWithThisHasGtidAndThatDoesNotFallsBackToFilename() {
+ BinlogOffset thisWithGtid =
+ new BinlogOffset("mysql-bin.000002", 4L, 1L, 5L, 0L,
GTID_SET_A, 1);
+ BinlogOffset thatWithoutGtid =
+ new BinlogOffset("mysql-bin.000001", 999L, 1L, 5L, 0L, null,
1);
+
+ Assertions.assertTrue(
+ thisWithGtid.compareTo(thatWithoutGtid) > 0,
+ "when this has a GTID but that does not, comparison must fall
back to "
+ + "binlog filename ordering");
+ Assertions.assertTrue(thatWithoutGtid.compareTo(thisWithGtid) < 0);
+ }
+
+ @Test
+ public void
testCompareToWithThisHasGtidAndThatDoesNotSameFileUsesPosition() {
+ BinlogOffset thisWithGtid =
+ new BinlogOffset("mysql-bin.000001", 100L, 1L, 5L, 0L,
GTID_SET_A, 1);
+ BinlogOffset thatWithoutGtid =
+ new BinlogOffset("mysql-bin.000001", 50L, 1L, 5L, 0L, null, 1);
+
+ Assertions.assertTrue(
+ thisWithGtid.compareTo(thatWithoutGtid) > 0,
+ "same binlog file must be ordered by position when GTIDs are
mixed");
+ Assertions.assertTrue(thatWithoutGtid.compareTo(thisWithGtid) < 0);
+ }
+
+ @Test
+ public void testIsNeverStop() {
+ Assertions.assertTrue(
+ BinlogOffset.NO_STOPPING_OFFSET.isNeverStop(),
+ "the unbounded sentinel must report itself as never-stop");
+ BinlogOffset regular = new BinlogOffset("mysql-bin.000001", 4L, 1L,
5L, 0L, GTID_SET_A, 1);
+ Assertions.assertFalse(regular.isNeverStop());
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/offset/RedoLogOffset.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/offset/RedoLogOffset.java
index 5640acae02..7bea77b867 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/offset/RedoLogOffset.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/offset/RedoLogOffset.java
@@ -118,4 +118,9 @@ public class RedoLogOffset extends Offset {
RedoLogOffset that = (RedoLogOffset) o;
return offset.equals(that.offset);
}
+
+ @Override
+ public boolean isNeverStop() {
+ return NO_STOPPING_OFFSET.equals(this);
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/offset/RedoLogOffsetTest.java
similarity index 61%
copy from
seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
copy to
seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/offset/RedoLogOffsetTest.java
index 7c62cb46c5..e3f6a30716 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-oracle/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/oracle/source/offset/RedoLogOffsetTest.java
@@ -15,24 +15,16 @@
* limitations under the License.
*/
-package org.apache.seatunnel.connectors.seatunnel.cdc.sqlserver.source.offset;
+package org.apache.seatunnel.connectors.seatunnel.cdc.oracle.source.offset;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
-import io.debezium.connector.sqlserver.Lsn;
-
-class LsnOffsetTest {
+class RedoLogOffsetTest {
@Test
- void testInitialOffsetRepresentsNoLsn() {
- LsnOffset initial = LsnOffset.INITIAL_OFFSET;
-
- // no LSN keys should be present in the offset map
- Assertions.assertTrue(initial.getOffset().isEmpty());
-
- // commit LSN resolved from the empty map should be Debezium's NULL LSN
- Lsn commitLsn = initial.getCommitLsn();
- Assertions.assertFalse(commitLsn.isAvailable());
+ void testNoStoppingOffsetIsNeverStop() {
+ Assertions.assertTrue(RedoLogOffset.NO_STOPPING_OFFSET.isNeverStop());
+ Assertions.assertFalse(RedoLogOffset.INITIAL_OFFSET.isNeverStop());
}
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/offset/LsnOffset.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/offset/LsnOffset.java
index f65e916609..bcad272af1 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/offset/LsnOffset.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/offset/LsnOffset.java
@@ -130,4 +130,9 @@ public class LsnOffset extends Offset {
LsnOffset that = (LsnOffset) o;
return offset.equals(that.offset);
}
+
+ @Override
+ public boolean isNeverStop() {
+ return NO_STOPPING_OFFSET.equals(this);
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/offset/LsnOffsetTest.java
similarity index 63%
copy from
seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
copy to
seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/offset/LsnOffsetTest.java
index 7c62cb46c5..cd6753547e 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-postgres/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/source/offset/LsnOffsetTest.java
@@ -15,24 +15,16 @@
* limitations under the License.
*/
-package org.apache.seatunnel.connectors.seatunnel.cdc.sqlserver.source.offset;
+package org.apache.seatunnel.connectors.seatunnel.cdc.postgres.source.offset;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
-import io.debezium.connector.sqlserver.Lsn;
-
class LsnOffsetTest {
@Test
- void testInitialOffsetRepresentsNoLsn() {
- LsnOffset initial = LsnOffset.INITIAL_OFFSET;
-
- // no LSN keys should be present in the offset map
- Assertions.assertTrue(initial.getOffset().isEmpty());
-
- // commit LSN resolved from the empty map should be Debezium's NULL LSN
- Lsn commitLsn = initial.getCommitLsn();
- Assertions.assertFalse(commitLsn.isAvailable());
+ void testNoStoppingOffsetIsNeverStop() {
+ Assertions.assertTrue(LsnOffset.NO_STOPPING_OFFSET.isNeverStop());
+ Assertions.assertFalse(LsnOffset.INITIAL_OFFSET.isNeverStop());
}
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffset.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffset.java
index 7fa7a59f88..df046c303a 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffset.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffset.java
@@ -95,4 +95,9 @@ public class LsnOffset extends Offset {
prime * result + ((getEventSerialNo() == null) ? 0 :
getEventSerialNo().hashCode());
return result;
}
+
+ @Override
+ public boolean isNeverStop() {
+ return NO_STOPPING_OFFSET.equals(this);
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
index 7c62cb46c5..00ceb29389 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetTest.java
@@ -35,4 +35,10 @@ class LsnOffsetTest {
Lsn commitLsn = initial.getCommitLsn();
Assertions.assertFalse(commitLsn.isAvailable());
}
+
+ @Test
+ void testNoStoppingOffsetIsNeverStop() {
+ Assertions.assertTrue(LsnOffset.NO_STOPPING_OFFSET.isNeverStop());
+ Assertions.assertFalse(LsnOffset.INITIAL_OFFSET.isNeverStop());
+ }
}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCStopModeSpecificIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCStopModeSpecificIT.java
new file mode 100644
index 0000000000..cd3a6b4117
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/mysql/MysqlCDCStopModeSpecificIT.java
@@ -0,0 +1,434 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.cdc.mysql;
+
+import org.apache.seatunnel.connectors.cdc.base.config.JdbcSourceConfig;
+import org.apache.seatunnel.connectors.cdc.base.config.JdbcSourceConfigFactory;
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.config.MySqlSourceConfigFactory;
+import org.apache.seatunnel.connectors.seatunnel.cdc.mysql.source.MySqlDialect;
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.source.offset.BinlogOffset;
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.testutils.MySqlContainer;
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.testutils.MySqlVersion;
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.testutils.UniqueDatabase;
+import
org.apache.seatunnel.connectors.seatunnel.cdc.mysql.utils.MySqlConnectionUtils;
+import org.apache.seatunnel.e2e.common.TestResource;
+import org.apache.seatunnel.e2e.common.TestSuiteBase;
+import org.apache.seatunnel.e2e.common.container.ContainerExtendedFactory;
+import org.apache.seatunnel.e2e.common.container.EngineType;
+import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.junit.DisabledOnContainer;
+import org.apache.seatunnel.e2e.common.junit.TestContainerExtension;
+import org.apache.seatunnel.e2e.common.util.JobIdGenerator;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.TestTemplate;
+import org.testcontainers.containers.Container;
+import org.testcontainers.containers.output.Slf4jLogConsumer;
+import org.testcontainers.lifecycle.Startables;
+import org.testcontainers.utility.DockerLoggerFactory;
+
+import com.github.shyiko.mysql.binlog.BinaryLogClient;
+import com.github.shyiko.mysql.binlog.event.EventData;
+import com.github.shyiko.mysql.binlog.event.EventHeaderV4;
+import com.github.shyiko.mysql.binlog.event.FormatDescriptionEventData;
+import com.github.shyiko.mysql.binlog.event.RotateEventData;
+import io.debezium.jdbc.JdbcConnection;
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.IOException;
+import java.sql.Connection;
+import java.sql.DriverManager;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.ArrayBlockingQueue;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import java.util.stream.Stream;
+
+import static org.awaitility.Awaitility.await;
+
+/**
+ * Integration test for MySQL CDC {@code stop.mode = "specific"}.
+ *
+ * <p>Starts a MySQL CDC job that reads the binlog between a configured start
offset and a
+ * configured stop offset, and asserts that the job terminates (FINISHED) at
the stop offset instead
+ * of running forever, and that only the rows written before the stop offset
are synced.
+ */
+@Slf4j
+@DisabledOnContainer(
+ value = {},
+ type = {EngineType.FLINK, EngineType.SPARK},
+ disabledReason =
+ "Currently FLINK and SPARK do not support the bounded
incremental split "
+ + "termination for MySQL CDC stop.mode=specific")
+public class MysqlCDCStopModeSpecificIT extends TestSuiteBase implements
TestResource {
+
+ private static final String MYSQL_HOST = "mysql_cdc_e2e";
+ private static final String MYSQL_USER_NAME = "mysqluser";
+ private static final String MYSQL_USER_PASSWORD = "mysqlpw";
+ private static final String MYSQL_DATABASE = "mysql_cdc";
+ private static final MySqlContainer MYSQL_CONTAINER =
createMySqlContainer(MySqlVersion.V8_0);
+
+ private final UniqueDatabase inventoryDatabase =
+ new UniqueDatabase(
+ MYSQL_CONTAINER, MYSQL_DATABASE, "mysqluser", "mysqlpw",
MYSQL_DATABASE);
+
+ private static final String SOURCE_TABLE = "mysql_cdc_e2e_source_table";
+ private static final String SINK_TABLE = "mysql_cdc_e2e_sink_table";
+
+ private static MySqlContainer createMySqlContainer(MySqlVersion version) {
+ return new MySqlContainer(version)
+ .withConfigurationOverride("docker/server-gtids/my.cnf")
+ .withSetupSQL("docker/setup.sql")
+ .withNetwork(NETWORK)
+ .withNetworkAliases(MYSQL_HOST)
+ .withDatabaseName(MYSQL_DATABASE)
+ .withUsername(MYSQL_USER_NAME)
+ .withPassword(MYSQL_USER_PASSWORD)
+ .withLogConsumer(
+ new
Slf4jLogConsumer(DockerLoggerFactory.getLogger("mysql-docker-image")));
+ }
+
+ @TestContainerExtension
+ protected final ContainerExtendedFactory extendedFactory =
+ MysqlCDCDriverResolver::copyMySQLDriverToContainer;
+
+ @BeforeAll
+ @Override
+ public void startUp() {
+ log.info("Starting Mysql container for stop.mode=specific e2e
test...");
+ Startables.deepStart(Stream.of(MYSQL_CONTAINER)).join();
+ inventoryDatabase.createAndInitialize();
+ log.info("Mysql ddl execution is complete");
+ }
+
+ @AfterAll
+ @Override
+ public void tearDown() {
+ if (MYSQL_CONTAINER != null) {
+ MYSQL_CONTAINER.close();
+ }
+ }
+
+ @TestTemplate
+ public void
testMysqlCdcStopModeSpecificTerminatesAtStopOffset(TestContainer container)
+ throws Exception {
+ String jobId = String.valueOf(JobIdGenerator.newJobId());
+ String jobConfigFile = "/mysqlcdc_stop_mode_specific.conf";
+
+ clearTable(MYSQL_DATABASE, SOURCE_TABLE);
+ clearTable(MYSQL_DATABASE, SINK_TABLE);
+
+ // Current binlog position, used as the start offset.
+ BinlogOffset startOffset = getCurrentBinlogOffset();
+
+ // Rows written before the stop offset: must be synced.
+ executeSql(
+ String.format(
+ "INSERT INTO %s.%s (id) VALUES (21), (22)",
MYSQL_DATABASE, SOURCE_TABLE));
+
+ // Current binlog position after the rows above: the stop offset.
+ BinlogOffset stopOffset = getCurrentBinlogOffset();
+
+ String[] variables = {
+ "start_offset_file=" + startOffset.getFilename(),
+ "start_offset_pos=" + startOffset.getPosition(),
+ "stop_offset_file=" + stopOffset.getFilename(),
+ "stop_offset_pos=" + stopOffset.getPosition()
+ };
+
+ CompletableFuture<Container.ExecResult> jobFuture =
+ CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return container.executeJob(jobConfigFile,
jobId, variables);
+ } catch (Exception e) {
+ log.error("Commit task exception :" +
e.getMessage());
+ throw new RuntimeException(e);
+ }
+ });
+
+ // Wait until the rows before the stop offset are synced, so the
snapshot phase
+ // (if any) has completed and only the binlog phase is still running.
+ await().atMost(60000, TimeUnit.MILLISECONDS)
+ .untilAsserted(
+ () -> {
+ List<List<Object>> sinkIds = queryIds();
+ Assertions.assertTrue(
+
sinkIds.contains(Collections.singletonList(21))
+ &&
sinkIds.contains(Collections.singletonList(22)),
+ "rows before the stop offset must be
synced, got: " + sinkIds);
+ });
+
+ // Rows written after the stop offset: must NOT be synced. They are
inserted only
+ // after the snapshot phase has completed, so they can only be picked
up by the
+ // binlog phase, which must stop at the configured stop offset.
+ executeSql(
+ String.format(
+ "INSERT INTO %s.%s (id) VALUES (23), (24)",
MYSQL_DATABASE, SOURCE_TABLE));
+
+ // The bounded job must terminate on its own at the stop offset.
+ Container.ExecResult result = jobFuture.get(120, TimeUnit.SECONDS);
+ Assertions.assertEquals(0, result.getExitCode(), result.getStderr());
+
+ // Only the rows written before the stop offset must be present in the
sink.
+ await().atMost(30000, TimeUnit.MILLISECONDS)
+ .untilAsserted(
+ () -> {
+ List<List<Object>> sinkIds = queryIds();
+ Assertions.assertTrue(
+
sinkIds.contains(Collections.singletonList(21))
+ &&
sinkIds.contains(Collections.singletonList(22)),
+ "rows before the stop offset must be
synced, got: " + sinkIds);
+ Assertions.assertFalse(
+
sinkIds.contains(Collections.singletonList(23))
+ ||
sinkIds.contains(Collections.singletonList(24)),
+ "rows after the stop offset must not be
synced, got: "
+ + sinkIds);
+ });
+ }
+
+ @TestTemplate
+ public void testMysqlCdcStopModeSpecificWithTimestampStartup(TestContainer
container)
+ throws Exception {
+ String jobId = String.valueOf(JobIdGenerator.newJobId());
+ String jobConfigFile = "/mysqlcdc_stop_mode_specific_timestamp.conf";
+
+ clearTable(MYSQL_DATABASE, SOURCE_TABLE);
+ clearTable(MYSQL_DATABASE, SINK_TABLE);
+
+ // Rows written before the startup timestamp: must NOT be synced.
+ executeSql(
+ String.format(
+ "INSERT INTO %s.%s (id) VALUES (31), (32)",
MYSQL_DATABASE, SOURCE_TABLE));
+
+ // MySQL binlog timestamps have second granularity, wait so the
startup timestamp
+ // is clearly after the rows above.
+ Thread.sleep(3000L);
+
+ // Take the startup timestamp before inserting the rows that must be
synced,
+ // so their binlog event timestamps are greater than the startup
timestamp.
+ long startTimestamp = getCurrentBinlogTimestamp() + 2000L;
+
+ // Rows written after the startup timestamp: must be synced.
+ executeSql(
+ String.format(
+ "INSERT INTO %s.%s (id) VALUES (33), (34)",
MYSQL_DATABASE, SOURCE_TABLE));
+
+ // Current binlog position after the rows above: the stop offset.
+ BinlogOffset stopOffset = getCurrentBinlogOffset();
+
+ String[] variables = {
+ "start_timestamp=" + startTimestamp,
+ "stop_offset_file=" + stopOffset.getFilename(),
+ "stop_offset_pos=" + stopOffset.getPosition()
+ };
+ log.info("Startup timestamp :{}", variables[0]);
+
+ CompletableFuture<Container.ExecResult> jobFuture =
+ CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return container.executeJob(jobConfigFile,
jobId, variables);
+ } catch (Exception e) {
+ log.error("Commit task exception :" +
e.getMessage());
+ throw new RuntimeException(e);
+ }
+ });
+
+ // Wait until the rows after the startup timestamp are synced, so the
snapshot phase
+ // (if any) has completed and only the binlog phase is still running.
+ await().atMost(60000, TimeUnit.MILLISECONDS)
+ .untilAsserted(
+ () -> {
+ List<List<Object>> sinkIds = queryIds();
+ Assertions.assertTrue(
+
sinkIds.contains(Collections.singletonList(33))
+ &&
sinkIds.contains(Collections.singletonList(34)),
+ "rows after the startup timestamp must be
synced, got: "
+ + sinkIds);
+ });
+
+ // Rows written after the stop offset: must NOT be synced. They are
inserted only
+ // after the snapshot phase has completed, so they can only be picked
up by the
+ // binlog phase, which must stop at the configured stop offset.
+ executeSql(
+ String.format(
+ "INSERT INTO %s.%s (id) VALUES (35), (36)",
MYSQL_DATABASE, SOURCE_TABLE));
+
+ // The bounded job must terminate on its own at the stop offset.
+ Container.ExecResult result = jobFuture.get(120, TimeUnit.SECONDS);
+ Assertions.assertEquals(0, result.getExitCode(), result.getStderr());
+
+ // Only the rows written after the startup timestamp and before the
stop offset
+ // must be present in the sink.
+ await().atMost(30000, TimeUnit.MILLISECONDS)
+ .untilAsserted(
+ () -> {
+ List<List<Object>> sinkIds = queryIds();
+ Assertions.assertTrue(
+
sinkIds.contains(Collections.singletonList(33))
+ &&
sinkIds.contains(Collections.singletonList(34)),
+ "rows after the startup timestamp must be
synced, got: "
+ + sinkIds);
+ Assertions.assertFalse(
+
sinkIds.contains(Collections.singletonList(31))
+ ||
sinkIds.contains(Collections.singletonList(32)),
+ "rows before the startup timestamp must
not be synced, got: "
+ + sinkIds);
+ Assertions.assertFalse(
+
sinkIds.contains(Collections.singletonList(35))
+ ||
sinkIds.contains(Collections.singletonList(36)),
+ "rows after the stop offset must not be
synced, got: "
+ + sinkIds);
+ });
+ }
+
+ private long getCurrentBinlogTimestamp() {
+ BinlogOffset binlogOffset = getCurrentBinlogOffset();
+
+ JdbcSourceConfigFactory configFactory =
+ new MySqlSourceConfigFactory()
+ .hostname(MYSQL_CONTAINER.getHost())
+ .port(MYSQL_CONTAINER.getDatabasePort())
+ .username(MYSQL_CONTAINER.getUsername())
+ .password(MYSQL_CONTAINER.getPassword())
+ .databaseList(MYSQL_CONTAINER.getDatabaseName());
+ JdbcSourceConfig jdbcSourceConfig = configFactory.create(0);
+ MySqlDialect mySqlDialect =
+ new MySqlDialect((MySqlSourceConfigFactory) configFactory,
Collections.emptyList());
+ BinaryLogClient client =
+
MySqlConnectionUtils.createBinaryClient(jdbcSourceConfig.getDbzConfiguration());
+
+ final String showBinaryLogStmt =
+ "SHOW BINLOG EVENTS IN '" + binlogOffset.getFilename() + "'";
+ List<Long> logPosList = new ArrayList<>();
+ JdbcConnection.ResultSetConsumer rsc =
+ rs -> {
+ while (rs.next()) {
+ logPosList.add(rs.getLong(5));
+ }
+ };
+ try (JdbcConnection jdbc =
mySqlDialect.openJdbcConnection(jdbcSourceConfig)) {
+ jdbc.query(showBinaryLogStmt, rsc);
+ if (logPosList.isEmpty()) {
+ return System.currentTimeMillis();
+ }
+ Long pos =
+ logPosList.stream()
+ .distinct()
+ .sorted(Collections.reverseOrder())
+ .collect(Collectors.toList())
+ .get(1);
+
+ ArrayBlockingQueue<Long> binlogTimestamps = new
ArrayBlockingQueue<>(1);
+ BinaryLogClient.EventListener eventListener =
+ event -> {
+ EventData data = event.getData();
+ if (data instanceof RotateEventData
+ || data instanceof FormatDescriptionEventData)
{
+ return;
+ }
+ EventHeaderV4 header = event.getHeader();
+ long timestamp = header.getTimestamp();
+ if (timestamp > 0) {
+ binlogTimestamps.offer(timestamp);
+ try {
+ client.disconnect();
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+ }
+ };
+ try {
+ client.registerEventListener(eventListener);
+ client.setBinlogFilename(binlogOffset.getFilename());
+ client.setBinlogPosition(pos);
+ client.connect();
+ } finally {
+ client.unregisterEventListener(eventListener);
+ }
+ return binlogTimestamps.take();
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private Connection getJdbcConnection() throws SQLException {
+ return DriverManager.getConnection(
+ MYSQL_CONTAINER.getJdbcUrl(),
+ MYSQL_CONTAINER.getUsername(),
+ MYSQL_CONTAINER.getPassword());
+ }
+
+ private List<List<Object>> queryIds() {
+ try (Connection connection = getJdbcConnection();
+ Statement statement = connection.createStatement();
+ ResultSet resultSet =
+ statement.executeQuery(
+ String.format(
+ "select id from %s.%s",
MYSQL_DATABASE, SINK_TABLE))) {
+ List<List<Object>> result = new ArrayList<>();
+ while (resultSet.next()) {
+ result.add(Collections.singletonList(resultSet.getObject(1)));
+ }
+ return result;
+ } catch (SQLException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private void executeSql(String sql) {
+ try (Connection connection = getJdbcConnection()) {
+ connection.createStatement().execute(sql);
+ } catch (SQLException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private void clearTable(String database, String tableName) {
+ try (Connection connection = getJdbcConnection()) {
+ connection
+ .createStatement()
+ .execute(String.format("TRUNCATE TABLE %s.%s", database,
tableName));
+ } catch (SQLException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ private BinlogOffset getCurrentBinlogOffset() {
+ JdbcSourceConfigFactory configFactory =
+ new MySqlSourceConfigFactory()
+ .hostname(MYSQL_CONTAINER.getHost())
+ .port(MYSQL_CONTAINER.getDatabasePort())
+ .username(MYSQL_CONTAINER.getUsername())
+ .password(MYSQL_CONTAINER.getPassword())
+ .databaseList(MYSQL_CONTAINER.getDatabaseName());
+ MySqlDialect mySqlDialect =
+ new MySqlDialect((MySqlSourceConfigFactory) configFactory,
Collections.emptyList());
+ JdbcConnection jdbcConnection =
mySqlDialect.openJdbcConnection(configFactory.create(0));
+ return MySqlConnectionUtils.currentBinlogOffset(jdbcConnection);
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/resources/mysqlcdc_stop_mode_specific.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/resources/mysqlcdc_stop_mode_specific.conf
new file mode 100644
index 0000000000..fc1a143b0d
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/resources/mysqlcdc_stop_mode_specific.conf
@@ -0,0 +1,54 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 5000
+}
+
+source {
+ MySQL-CDC {
+ plugin_output = "customers_mysql_cdc"
+ server-id = 5656
+ username = "st_user_source"
+ password = "mysqlpw"
+ table-names = ["mysql_cdc.mysql_cdc_e2e_source_table"]
+ url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+ startup.mode = "specific"
+ startup.specific-offset.file = ${start_offset_file}
+ startup.specific-offset.pos = ${start_offset_pos}
+ stop.mode = "specific"
+ stop.specific-offset.file = ${stop_offset_file}
+ stop.specific-offset.pos = ${stop_offset_pos}
+ }
+}
+
+sink {
+ jdbc {
+ plugin_input = "customers_mysql_cdc"
+ url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+ driver = "com.mysql.cj.jdbc.Driver"
+ user = "st_user_sink"
+ password = "mysqlpw"
+
+ generate_sink_sql = true
+ database = mysql_cdc
+ table = mysql_cdc_e2e_sink_table
+ primary_keys = ["id"]
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/resources/mysqlcdc_stop_mode_specific_timestamp.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/resources/mysqlcdc_stop_mode_specific_timestamp.conf
new file mode 100644
index 0000000000..5815a0bd06
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-mysql-e2e/src/test/resources/mysqlcdc_stop_mode_specific_timestamp.conf
@@ -0,0 +1,53 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+ checkpoint.interval = 5000
+}
+
+source {
+ MySQL-CDC {
+ plugin_output = "customers_mysql_cdc"
+ server-id = 5657
+ username = "st_user_source"
+ password = "mysqlpw"
+ table-names = ["mysql_cdc.mysql_cdc_e2e_source_table"]
+ url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+ startup.mode = "timestamp"
+ startup.timestamp = ${start_timestamp}
+ stop.mode = "specific"
+ stop.specific-offset.file = ${stop_offset_file}
+ stop.specific-offset.pos = ${stop_offset_pos}
+ }
+}
+
+sink {
+ jdbc {
+ plugin_input = "customers_mysql_cdc"
+ url = "jdbc:mysql://mysql_cdc_e2e:3306/mysql_cdc"
+ driver = "com.mysql.cj.jdbc.Driver"
+ user = "st_user_sink"
+ password = "mysqlpw"
+
+ generate_sink_sql = true
+ database = mysql_cdc
+ table = mysql_cdc_e2e_sink_table
+ primary_keys = ["id"]
+ }
+}