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 fbd82a4eb0 [Fix][Connector-V2] Preserve SQL Server CDC resume offsets
(#11410)
fbd82a4eb0 is described below
commit fbd82a4eb07cc9ab8d7e6ba00032d8474a31846b
Author: Jast <[email protected]>
AuthorDate: Sun Aug 23 15:29:01 2026 +0800
[Fix][Connector-V2] Preserve SQL Server CDC resume offsets (#11410)
Co-authored-by: zhangshenghang <[email protected]>
Co-authored-by: davidzollo <[email protected]>
Co-authored-by: davidzollo <[email protected]>
---
.../cdc/sqlserver/source/offset/LsnOffset.java | 78 ++++++++++++++-
.../sqlserver/source/offset/LsnOffsetFactory.java | 3 +-
.../fetch/SqlServerSourceFetchTaskContext.java | 14 ++-
.../cdc/sqlserver/utils/SqlServerUtils.java | 5 +-
.../cdc/sqlserver/source/offset/LsnOffsetTest.java | 109 +++++++++++++++++++++
.../cdc/sqlserver/utils/SqlServerUtilsTest.java | 21 ++++
6 files changed, 223 insertions(+), 7 deletions(-)
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 df046c303a..6e2ef9b08d 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
@@ -37,6 +37,42 @@ public class LsnOffset extends Offset {
return new LsnOffset(Lsn.valueOf(commitLsn), null, null);
}
+ /**
+ * Creates an offset from the full SQL Server position reported by
Debezium.
+ *
+ * <p>SQL Server can emit multiple change events for the same commit LSN.
The change LSN and
+ * event serial number are therefore required to resume without skipping
records.
+ */
+ public static LsnOffset valueOf(Map<String, ?> offset) {
+ Object eventSerialNo = offset.get(SourceInfo.EVENT_SERIAL_NO_KEY);
+ Object commitLsn = offset.get(SourceInfo.COMMIT_LSN_KEY);
+ Object changeLsn = offset.get(SourceInfo.CHANGE_LSN_KEY);
+ return new LsnOffset(
+ Lsn.valueOf(commitLsn == null ? null : commitLsn.toString()),
+ Lsn.valueOf(changeLsn == null ? null : changeLsn.toString()),
+ eventSerialNo == null ? null :
Long.valueOf(eventSerialNo.toString()));
+ }
+
+ /**
+ * Creates a boundary offset suitable for {@code startup.mode=timestamp}.
+ *
+ * <p>{@code sys.fn_cdc_map_time_to_lsn('smallest greater than or equal',
ts)} returns the
+ * COMMIT lsn of the first transaction whose commit time is at or after
the requested timestamp.
+ * Rows belonging to that transaction must be emitted, so the boundary
must order itself BEFORE
+ * any real change event at the same commit. This is achieved by leaving
the commit LSN intact
+ * while using the smallest available change LSN and event serial number
as the in-commit
+ * position. With {@link #compareTo(Offset)}, every real in-commit event
then compares as {@code
+ * isAfter(this)}.
+ *
+ * <p>Compare with {@link #valueOf(String)} (commit-only), which is used
for {@code
+ * startup.mode=latest} and intentionally orders AFTER same-commit events
to avoid replaying
+ * rows that already existed before startup.
+ */
+ public static LsnOffset timestampBoundary(String commitLsn) {
+ return new LsnOffset(
+ Lsn.valueOf(commitLsn), Lsn.valueOf(new byte[] {0, 0, 0, 0, 0,
0, 0, 0, 0, 1}), 0L);
+ }
+
private LsnOffset(Lsn commitLsn, Lsn changeLsn, Long eventSerialNo) {
Map<String, String> offsetMap = new HashMap<>();
@@ -68,7 +104,47 @@ public class LsnOffset extends Offset {
public int compareTo(Offset o) {
LsnOffset that = (LsnOffset) o;
final int comparison = getCommitLsn().compareTo(that.getCommitLsn());
- return comparison == 0 ? getChangeLsn().compareTo(that.getChangeLsn())
: comparison;
+ if (comparison != 0) {
+ return comparison;
+ }
+ // A commit-only offset carries no in-commit position. It can
represent the latest
+ // startup boundary, a timestamp-derived boundary, a legacy coarse
checkpoint or a
+ // snapshot watermark, all of which address a whole commit (everything
in that commit
+ // is treated as already processed before the boundary takes effect).
+ //
+ // Comparing a commit-only boundary against a same-commit complete
event position must
+ // therefore order the complete event BEFORE the boundary, otherwise:
+ // * the non-exactly-once path would skip same-commit rows
(acceptable, but stricter
+ // than needed),
+ // * the exactly-once pure-binlog transition (`isAtOrAfter`) would
flip a table into
+ // "emit everything" mode on the very first record of the boundary
commit and
+ // replay rows already committed before startup. That is a
data-correctness
+ // regression on the normal streaming path for
`startup.mode=latest`.
+ final boolean thisComplete = hasCompletePosition();
+ final boolean thatComplete = that.hasCompletePosition();
+ if (!thisComplete && !thatComplete) {
+ return 0;
+ }
+ if (!thisComplete) {
+ return 1;
+ }
+ if (!thatComplete) {
+ return -1;
+ }
+ final int changeLsnComparison =
getChangeLsn().compareTo(that.getChangeLsn());
+ if (changeLsnComparison != 0) {
+ return changeLsnComparison;
+ }
+ return Long.compare(eventSerialNo(), that.eventSerialNo());
+ }
+
+ private boolean hasCompletePosition() {
+ return getChangeLsn().isAvailable() && getEventSerialNo() != null;
+ }
+
+ private long eventSerialNo() {
+ Object eventSerialNo = getEventSerialNo();
+ return eventSerialNo == null ? 0L :
Long.parseLong(eventSerialNo.toString());
}
public boolean equals(Object obj) {
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetFactory.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetFactory.java
index f24075ebb8..674b055162 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetFactory.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/offset/LsnOffsetFactory.java
@@ -24,7 +24,6 @@ import
org.apache.seatunnel.connectors.seatunnel.cdc.sqlserver.config.SqlServerS
import
org.apache.seatunnel.connectors.seatunnel.cdc.sqlserver.source.SqlServerDialect;
import
org.apache.seatunnel.connectors.seatunnel.cdc.sqlserver.utils.SqlServerUtils;
-import io.debezium.connector.sqlserver.SourceInfo;
import io.debezium.connector.sqlserver.SqlServerConnection;
import io.debezium.jdbc.JdbcConnection;
@@ -62,7 +61,7 @@ public class LsnOffsetFactory extends OffsetFactory {
@Override
public Offset specific(Map<String, String> offset) {
- return LsnOffset.valueOf(offset.get(SourceInfo.COMMIT_LSN_KEY));
+ return LsnOffset.valueOf(offset);
}
@Override
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/reader/fetch/SqlServerSourceFetchTaskContext.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/reader/fetch/SqlServerSourceFetchTaskContext.java
index 07e5f14f9a..9b5623bb52 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/reader/fetch/SqlServerSourceFetchTaskContext.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/source/reader/fetch/SqlServerSourceFetchTaskContext.java
@@ -60,6 +60,7 @@ import lombok.extern.slf4j.Slf4j;
import java.sql.SQLException;
import java.time.Instant;
+import java.util.HashMap;
import java.util.Map;
/** The context for fetch task that fetching data of snapshot split from MySQL
data source. */
@@ -272,7 +273,18 @@ public class SqlServerSourceFetchTaskContext extends
JdbcSourceFetchTaskContext
? LsnOffset.INITIAL_OFFSET
: split.asIncrementalSplit().getStartupOffset();
- SqlServerOffsetContext sqlServerOffsetContext =
loader.load(offset.getOffset());
+ Map<String, Object> restoredOffset = new HashMap<>(offset.getOffset());
+ // Debezium's SqlServerOffsetContext.Loader expects `event_serial_no`
to be a Number;
+ // SeaTunnel's checkpoint/offset map only stores String values, so the
previously
+ // serialized String must be converted back to Long here. Without this
coercion the
+ // loader silently rejects the offset and the original `#10571`
resume-precision bug
+ // reappears on the restore path.
+ String eventSerialNo =
offset.getOffset().get(SourceInfo.EVENT_SERIAL_NO_KEY);
+ if (eventSerialNo != null) {
+ restoredOffset.put(SourceInfo.EVENT_SERIAL_NO_KEY,
Long.valueOf(eventSerialNo));
+ }
+
+ SqlServerOffsetContext sqlServerOffsetContext =
loader.load(restoredOffset);
return sqlServerOffsetContext;
}
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/utils/SqlServerUtils.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/utils/SqlServerUtils.java
index 5867a84c87..95908aebee 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/utils/SqlServerUtils.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/utils/SqlServerUtils.java
@@ -27,7 +27,6 @@ import
org.apache.seatunnel.connectors.seatunnel.cdc.sqlserver.source.offset.Lsn
import org.apache.kafka.connect.source.SourceRecord;
import io.debezium.connector.sqlserver.Lsn;
-import io.debezium.connector.sqlserver.SourceInfo;
import io.debezium.connector.sqlserver.SqlServerConnection;
import io.debezium.connector.sqlserver.SqlServerConnectorConfig;
import io.debezium.connector.sqlserver.SqlServerDatabaseSchema;
@@ -272,7 +271,7 @@ public class SqlServerUtils {
offsetStrMap.put(
entry.getKey(), entry.getValue() == null ? null :
entry.getValue().toString());
}
- return LsnOffset.valueOf(offsetStrMap.get(SourceInfo.COMMIT_LSN_KEY));
+ return LsnOffset.valueOf(offsetStrMap);
}
/** Fetch current largest log sequence number (LSN) of the database. */
@@ -331,7 +330,7 @@ public class SqlServerUtils {
timestampMs,
new Timestamp(timestampMs),
lsn);
- return LsnOffset.valueOf(lsn.toString());
+ return LsnOffset.timestampBoundary(lsn.toString());
});
} catch (SQLException e) {
throw new SeaTunnelException(
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 00ceb29389..c4949fa62b 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
@@ -21,9 +21,29 @@ import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import io.debezium.connector.sqlserver.Lsn;
+import io.debezium.connector.sqlserver.SourceInfo;
+
+import java.util.HashMap;
+import java.util.Map;
class LsnOffsetTest {
+ private static final String COMMIT_LSN = "00000027:00000a80:0003";
+
+ private static final String NEXT_COMMIT_LSN = "00000027:00000a80:0004";
+
+ private static final String CHANGE_LSN = "00000027:00000a80:0005";
+
+ private static final String NEXT_CHANGE_LSN = "00000027:00000a80:0006";
+
+ private static LsnOffset completeOffset(String commitLsn, String
changeLsn, long serialNo) {
+ Map<String, Object> offset = new HashMap<>();
+ offset.put(SourceInfo.COMMIT_LSN_KEY, commitLsn);
+ offset.put(SourceInfo.CHANGE_LSN_KEY, changeLsn);
+ offset.put(SourceInfo.EVENT_SERIAL_NO_KEY, serialNo);
+ return LsnOffset.valueOf(offset);
+ }
+
@Test
void testInitialOffsetRepresentsNoLsn() {
LsnOffset initial = LsnOffset.INITIAL_OFFSET;
@@ -36,6 +56,95 @@ class LsnOffsetTest {
Assertions.assertFalse(commitLsn.isAvailable());
}
+ @Test
+ void testCompleteOffsetsCompareChangeLsnAndEventSerialNo() {
+ LsnOffset first = completeOffset(COMMIT_LSN, CHANGE_LSN, 2L);
+ LsnOffset second = completeOffset(COMMIT_LSN, CHANGE_LSN, 3L);
+
+ Assertions.assertTrue(second.isAfter(first));
+ Assertions.assertFalse(first.isAfter(second));
+
+ LsnOffset laterChange = completeOffset(COMMIT_LSN, NEXT_CHANGE_LSN,
1L);
+ Assertions.assertTrue(laterChange.isAfter(second));
+ Assertions.assertFalse(second.isAfter(laterChange));
+ }
+
+ @Test
+ void testCommitOnlyBoundaryDoesNotReplayInCommitEvents() {
+ // startup.mode=latest records the current max commit as a commit-only
boundary; the
+ // streaming query then re-reads that commit with an inclusive lower
bound. Events of
+ // the boundary commit must be ordered BEFORE the boundary, otherwise
records that
+ // already existed before startup would be replayed. This applies to
both the
+ // non-exactly-once (`isAfter`) and the exactly-once (`isAtOrAfter` via
+ // IncrementalSourceStreamFetcher.hasEnterPureBinlogPhase) streaming
paths.
+ LsnOffset latestBoundary = LsnOffset.valueOf(COMMIT_LSN);
+ LsnOffset inCommitEvent = completeOffset(COMMIT_LSN, CHANGE_LSN, 1L);
+
+ Assertions.assertFalse(inCommitEvent.isAfter(latestBoundary));
+ Assertions.assertTrue(latestBoundary.isAfter(inCommitEvent));
+ Assertions.assertFalse(inCommitEvent.isAtOrAfter(latestBoundary));
+ Assertions.assertTrue(inCommitEvent.isAtOrBefore(latestBoundary));
+ Assertions.assertTrue(latestBoundary.isAtOrAfter(inCommitEvent));
+ Assertions.assertFalse(latestBoundary.isAtOrBefore(inCommitEvent));
+ }
+
+ @Test
+ void testEventAfterBoundaryCommitIsEmitted() {
+ LsnOffset latestBoundary = LsnOffset.valueOf(COMMIT_LSN);
+ LsnOffset laterEvent = completeOffset(NEXT_COMMIT_LSN, CHANGE_LSN, 1L);
+
+ Assertions.assertTrue(laterEvent.isAfter(latestBoundary));
+ Assertions.assertFalse(latestBoundary.isAfter(laterEvent));
+ }
+
+ @Test
+ void testCompleteEventIsAfterInitialOffset() {
+ LsnOffset event = completeOffset(COMMIT_LSN, CHANGE_LSN, 1L);
+
+ Assertions.assertTrue(event.isAfter(LsnOffset.INITIAL_OFFSET));
+ Assertions.assertFalse(LsnOffset.INITIAL_OFFSET.isAfter(event));
+ }
+
+ @Test
+ void testRestoredCheckpointOffsetKeepsInCommitPrecision() {
+ // The #10571 recovery path: a complete event position serialized into
a checkpoint is
+ // restored through the offset map and used as the incremental startup
offset. Events
+ // at or before the restored position must be skipped, later in-commit
events emitted.
+ LsnOffset checkpointed = completeOffset(COMMIT_LSN, CHANGE_LSN, 2L);
+ Map<String, String> serialized = new
HashMap<>(checkpointed.getOffset());
+
+ LsnOffset restoredStartupOffset = LsnOffset.valueOf(serialized);
+ LsnOffset replayedEvent = completeOffset(COMMIT_LSN, CHANGE_LSN, 2L);
+ LsnOffset nextEvent = completeOffset(COMMIT_LSN, CHANGE_LSN, 3L);
+
+ Assertions.assertFalse(replayedEvent.isAfter(restoredStartupOffset));
+ Assertions.assertTrue(nextEvent.isAfter(restoredStartupOffset));
+ }
+
+ @Test
+ void testTimestampBoundaryEmitsSameCommitEvents() {
+ // startup.mode=timestamp resolves the boundary via
fn_cdc_map_time_to_lsn('smallest
+ // greater than or equal', ts); the returned LSN is the COMMIT log
record of the
+ // matching transaction. Rows belonging to that transaction must be
emitted, so every
+ // real in-commit event must compare as isAfter(timestampBoundary).
+ LsnOffset boundary = LsnOffset.timestampBoundary(COMMIT_LSN);
+ LsnOffset firstInCommitEvent = completeOffset(COMMIT_LSN, CHANGE_LSN,
1L);
+ LsnOffset laterInCommitEvent = completeOffset(COMMIT_LSN,
NEXT_CHANGE_LSN, 1L);
+ LsnOffset earlierCommitEvent =
+ completeOffset("00000027:00000a80:0002",
"00000027:00000a80:0002", 1L);
+ LsnOffset laterCommitEvent = completeOffset(NEXT_COMMIT_LSN,
CHANGE_LSN, 1L);
+
+ Assertions.assertTrue(firstInCommitEvent.isAfter(boundary));
+ Assertions.assertTrue(laterInCommitEvent.isAfter(boundary));
+ Assertions.assertFalse(earlierCommitEvent.isAfter(boundary));
+ Assertions.assertTrue(laterCommitEvent.isAfter(boundary));
+ // The boundary itself carries a complete (commit, change, serial)
position so that
+ // #10571 resume-precision semantic still applies once a checkpoint
with a real
+ // change-LSN position is loaded.
+ Assertions.assertTrue(boundary.getChangeLsn().isAvailable());
+ Assertions.assertNotNull(boundary.getEventSerialNo());
+ }
+
@Test
void testNoStoppingOffsetIsNeverStop() {
Assertions.assertTrue(LsnOffset.NO_STOPPING_OFFSET.isNeverStop());
diff --git
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/utils/SqlServerUtilsTest.java
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/utils/SqlServerUtilsTest.java
index fa0c606c5e..63f23fcaf7 100644
---
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/utils/SqlServerUtilsTest.java
+++
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-sqlserver/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/sqlserver/utils/SqlServerUtilsTest.java
@@ -25,8 +25,12 @@ import
org.apache.seatunnel.connectors.seatunnel.cdc.sqlserver.source.offset.Lsn
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
+import io.debezium.connector.sqlserver.SourceInfo;
import io.debezium.relational.TableId;
+import java.util.HashMap;
+import java.util.Map;
+
public class SqlServerUtilsTest {
@Test
public void testSplitScanQuery() {
@@ -82,4 +86,21 @@ public class SqlServerUtilsTest {
Assertions.assertThrows(
RuntimeException.class, () ->
SqlServerUtils.lsnStringToOffset(invalidLsn));
}
+
+ @Test
+ public void testGetLsnPositionPreservesCompleteSqlServerOffset() {
+ Map<String, Object> sourceOffset = new HashMap<>();
+ sourceOffset.put(SourceInfo.COMMIT_LSN_KEY, "00000027:00000a80:0003");
+ sourceOffset.put(SourceInfo.CHANGE_LSN_KEY, "00000027:00000a80:0004");
+ sourceOffset.put(SourceInfo.EVENT_SERIAL_NO_KEY, 2L);
+
+ LsnOffset offset = SqlServerUtils.getLsnPosition(sourceOffset);
+
+ Assertions.assertEquals("00000027:00000a80:0003",
offset.getCommitLsn().toString());
+ Assertions.assertEquals("00000027:00000a80:0004",
offset.getChangeLsn().toString());
+ Assertions.assertEquals("2", offset.getEventSerialNo());
+
+ sourceOffset.put(SourceInfo.EVENT_SERIAL_NO_KEY, 3L);
+
Assertions.assertTrue(SqlServerUtils.getLsnPosition(sourceOffset).isAfter(offset));
+ }
}