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

Reply via email to