This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git


The following commit(s) were added to refs/heads/master by this push:
     new fa57a21eba [flink][cdc] Track schema event progress in CDC source 
reader (#8311)
fa57a21eba is described below

commit fa57a21eba31684f25543a6b32c167482bc296ff
Author: QuakeWang <[email protected]>
AuthorDate: Mon Jun 22 18:08:24 2026 +0800

    [flink][cdc] Track schema event progress in CDC source reader (#8311)
    
    The CDC source reader restored split progress only from data row
    `recordsToSkip`. Schema change events can be emitted before any data
    rows are consumed, but that progress was not checkpointed independently.
    After recovery, the reader could emit the same schema change events
    again.
    
    This PR adds explicit schema event progress to
    `TableAwareFileStoreSourceSplit` and tracks it with a CDC-specific split
    state. The split reader skips already emitted schema change events on
    restore, while preserving V1 checkpoint compatibility through a legacy
    progress marker.
---
 .../cdc/source/TableAwareFileStoreSourceSplit.java | 104 ++++++++++++++-
 .../cdc/source/reader/CDCRecordsWithSplitIds.java  |   3 +-
 .../cdc/source/reader/CDCSourceReader.java         |  16 +--
 .../cdc/source/reader/CDCSourceSplitReader.java    |  22 ++-
 .../cdc/source/reader/CDCSourceSplitState.java     |  78 +++++++++++
 ...bleAwareFileStoreSourceSplitSerializerTest.java |  63 ++++++++-
 .../enumerator/CDCCheckpointSerializerTest.java    |  28 +++-
 .../source/reader/CDCSourceSplitReaderTest.java    | 147 ++++++++++++++++++++-
 .../cdc/source/reader/CDCSourceSplitStateTest.java | 118 +++++++++++++++++
 9 files changed, 552 insertions(+), 27 deletions(-)

diff --git 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/TableAwareFileStoreSourceSplit.java
 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/TableAwareFileStoreSourceSplit.java
index aa05922520..500ae676c8 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/TableAwareFileStoreSourceSplit.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/TableAwareFileStoreSourceSplit.java
@@ -43,6 +43,8 @@ public class TableAwareFileStoreSourceSplit extends 
FileStoreSourceSplit {
     private final Identifier identifier;
     private final @Nullable Long lastSchemaId;
     private final long schemaId;
+    private final long schemaChangeEventsToSkip;
+    private final boolean legacySchemaProgress;
 
     public TableAwareFileStoreSourceSplit(
             String id,
@@ -51,10 +53,43 @@ public class TableAwareFileStoreSourceSplit extends 
FileStoreSourceSplit {
             Identifier identifier,
             @Nullable Long lastSchemaId,
             long schemaId) {
+        this(id, split, recordsToSkip, identifier, lastSchemaId, schemaId, 0L);
+    }
+
+    public TableAwareFileStoreSourceSplit(
+            String id,
+            Split split,
+            long recordsToSkip,
+            Identifier identifier,
+            @Nullable Long lastSchemaId,
+            long schemaId,
+            long schemaChangeEventsToSkip) {
+        this(
+                id,
+                split,
+                recordsToSkip,
+                identifier,
+                lastSchemaId,
+                schemaId,
+                schemaChangeEventsToSkip,
+                false);
+    }
+
+    private TableAwareFileStoreSourceSplit(
+            String id,
+            Split split,
+            long recordsToSkip,
+            Identifier identifier,
+            @Nullable Long lastSchemaId,
+            long schemaId,
+            long schemaChangeEventsToSkip,
+            boolean legacySchemaProgress) {
         super(id, split, recordsToSkip);
         this.identifier = identifier;
         this.lastSchemaId = lastSchemaId;
         this.schemaId = schemaId;
+        this.schemaChangeEventsToSkip = schemaChangeEventsToSkip;
+        this.legacySchemaProgress = legacySchemaProgress;
     }
 
     public Identifier getIdentifier() {
@@ -69,10 +104,38 @@ public class TableAwareFileStoreSourceSplit extends 
FileStoreSourceSplit {
         return schemaId;
     }
 
+    public long schemaChangeEventsToSkip() {
+        return schemaChangeEventsToSkip;
+    }
+
+    public TableAwareFileStoreSourceSplit updateWithProgress(
+            long recordsToSkip, long schemaChangeEventsToSkip) {
+        return new TableAwareFileStoreSourceSplit(
+                splitId(),
+                split(),
+                recordsToSkip,
+                identifier,
+                lastSchemaId,
+                schemaId,
+                schemaChangeEventsToSkip,
+                legacySchemaProgress);
+    }
+
+    public boolean isLegacySchemaProgress() {
+        return legacySchemaProgress;
+    }
+
     @Override
     public TableAwareFileStoreSourceSplit updateWithRecordsToSkip(long 
recordsToSkip) {
         return new TableAwareFileStoreSourceSplit(
-                splitId(), split(), recordsToSkip, identifier, lastSchemaId, 
schemaId);
+                splitId(),
+                split(),
+                recordsToSkip,
+                identifier,
+                lastSchemaId,
+                schemaId,
+                schemaChangeEventsToSkip,
+                legacySchemaProgress);
     }
 
     @Override
@@ -87,13 +150,22 @@ public class TableAwareFileStoreSourceSplit extends 
FileStoreSourceSplit {
                 && recordsToSkip() == other.recordsToSkip()
                 && identifier.equals(other.identifier)
                 && Objects.equals(lastSchemaId, other.lastSchemaId)
-                && schemaId == other.schemaId;
+                && schemaId == other.schemaId
+                && schemaChangeEventsToSkip == other.schemaChangeEventsToSkip
+                && legacySchemaProgress == other.legacySchemaProgress;
     }
 
     @Override
     public int hashCode() {
         return Objects.hash(
-                splitId(), split(), recordsToSkip(), identifier, lastSchemaId, 
schemaId);
+                splitId(),
+                split(),
+                recordsToSkip(),
+                identifier,
+                lastSchemaId,
+                schemaId,
+                schemaChangeEventsToSkip,
+                legacySchemaProgress);
     }
 
     @Override
@@ -112,17 +184,23 @@ public class TableAwareFileStoreSourceSplit extends 
FileStoreSourceSplit {
                 + lastSchemaId
                 + ", schemaId="
                 + schemaId
+                + ", schemaChangeEventsToSkip="
+                + schemaChangeEventsToSkip
+                + ", legacySchemaProgress="
+                + legacySchemaProgress
                 + '}';
     }
 
     /** The serializer for {@link TableAwareFileStoreSourceSplit}. */
     public static class Serializer
             implements 
SimpleVersionedSerializer<TableAwareFileStoreSourceSplit> {
+        private static final int VERSION_1 = 1;
+        private static final int VERSION_2 = 2;
         private static final Long NULL_SCHEMA_ID = -1L;
 
         @Override
         public int getVersion() {
-            return 1;
+            return VERSION_2;
         }
 
         @Override
@@ -136,12 +214,18 @@ public class TableAwareFileStoreSourceSplit extends 
FileStoreSourceSplit {
             view.writeLong(
                     split.getLastSchemaId() == null ? NULL_SCHEMA_ID : 
split.getLastSchemaId());
             view.writeLong(split.getSchemaId());
+            view.writeLong(split.schemaChangeEventsToSkip());
+            view.writeBoolean(split.isLegacySchemaProgress());
             return out.toByteArray();
         }
 
         @Override
         public TableAwareFileStoreSourceSplit deserialize(int version, byte[] 
serialized)
                 throws IOException {
+            if (version != VERSION_1 && version != VERSION_2) {
+                throw new IOException(
+                        "Unsupported TableAwareFileStoreSourceSplit version: " 
+ version);
+            }
             ByteArrayInputStream in = new ByteArrayInputStream(serialized);
             DataInputViewStreamWrapper view = new 
DataInputViewStreamWrapper(in);
             String splitId = view.readUTF();
@@ -158,8 +242,18 @@ public class TableAwareFileStoreSourceSplit extends 
FileStoreSourceSplit {
                 lastSchemaId = null;
             }
             long schemaId = view.readLong();
+            long schemaChangeEventsToSkip = version == VERSION_2 ? 
view.readLong() : 0L;
+            boolean legacySchemaProgress =
+                    version == VERSION_2 ? view.readBoolean() : version == 
VERSION_1;
             return new TableAwareFileStoreSourceSplit(
-                    splitId, split, recordsToSkip, identifier, lastSchemaId, 
schemaId);
+                    splitId,
+                    split,
+                    recordsToSkip,
+                    identifier,
+                    lastSchemaId,
+                    schemaId,
+                    schemaChangeEventsToSkip,
+                    legacySchemaProgress);
         }
     }
 }
diff --git 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCRecordsWithSplitIds.java
 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCRecordsWithSplitIds.java
index 5c95c3279f..b3a7e725d4 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCRecordsWithSplitIds.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCRecordsWithSplitIds.java
@@ -18,7 +18,6 @@
 
 package org.apache.paimon.flink.pipeline.cdc.source.reader;
 
-import org.apache.paimon.flink.source.FileStoreSourceSplitState;
 import org.apache.paimon.flink.source.metrics.FileStoreSourceReaderMetrics;
 import org.apache.paimon.utils.Reference;
 
@@ -109,7 +108,7 @@ public class CDCRecordsWithSplitIds implements 
RecordsWithSplitIds<RecordIterato
             SourceReaderContext context,
             RecordIterator<Event> element,
             SourceOutput<Event> output,
-            FileStoreSourceSplitState state,
+            CDCSourceSplitState state,
             FileStoreSourceReaderMetrics metrics) {
         long timestamp = TimestampAssigner.NO_TIMESTAMP;
         if (metrics.getLatestFileCreationTime() != 
FileStoreSourceReaderMetrics.UNDEFINED) {
diff --git 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceReader.java
 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceReader.java
index 05cf7880cd..ba7c4618b0 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceReader.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceReader.java
@@ -21,7 +21,6 @@ package org.apache.paimon.flink.pipeline.cdc.source.reader;
 import org.apache.paimon.disk.IOManager;
 import org.apache.paimon.flink.pipeline.cdc.source.CDCSource;
 import 
org.apache.paimon.flink.pipeline.cdc.source.TableAwareFileStoreSourceSplit;
-import org.apache.paimon.flink.source.FileStoreSourceSplitState;
 import org.apache.paimon.flink.source.metrics.FileStoreSourceReaderMetrics;
 
 import org.apache.flink.api.connector.source.SourceReader;
@@ -37,10 +36,7 @@ import java.util.Map;
 /** A {@link SourceReader} that read records from {@link 
TableAwareFileStoreSourceSplit}. */
 public class CDCSourceReader
         extends SingleThreadMultiplexSourceReaderBase<
-                RecordIterator<Event>,
-                Event,
-                TableAwareFileStoreSourceSplit,
-                FileStoreSourceSplitState> {
+                RecordIterator<Event>, Event, TableAwareFileStoreSourceSplit, 
CDCSourceSplitState> {
     private static final Logger LOG = 
LoggerFactory.getLogger(CDCSourceReader.class);
 
     private final IOManager ioManager;
@@ -69,7 +65,7 @@ public class CDCSourceReader
     }
 
     @Override
-    protected void onSplitFinished(Map<String, FileStoreSourceSplitState> 
finishedSplitIds) {
+    protected void onSplitFinished(Map<String, CDCSourceSplitState> 
finishedSplitIds) {
         // this method is called each time when we consume one split
         // it is possible that one response from the coordinator contains 
multiple splits
         // we should only require for more splits after we've consumed all 
given splits
@@ -79,16 +75,16 @@ public class CDCSourceReader
     }
 
     @Override
-    protected FileStoreSourceSplitState 
initializedState(TableAwareFileStoreSourceSplit split) {
+    protected CDCSourceSplitState 
initializedState(TableAwareFileStoreSourceSplit split) {
         LOG.info("Initializing split {}", split);
-        return new FileStoreSourceSplitState(split);
+        return new CDCSourceSplitState(split);
     }
 
     @Override
     protected TableAwareFileStoreSourceSplit toSplitType(
-            String splitId, FileStoreSourceSplitState splitState) {
+            String splitId, CDCSourceSplitState splitState) {
         LOG.info("Converting split state {} with id {} to split", splitState, 
splitId);
-        return (TableAwareFileStoreSourceSplit) splitState.toSourceSplit();
+        return splitState.toSourceSplit();
     }
 
     @Override
diff --git 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitReader.java
 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitReader.java
index 59d3f5799c..960665bf36 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitReader.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitReader.java
@@ -231,8 +231,10 @@ public class CDCSourceSplitReader
         currentReader = createLazyRecordReader(nextSplit.split());
         currentDataRowsRead = nextSplit.recordsToSkip();
         currentSchemaChangeEvents.clear();
-        if (currentDataRowsRead == 0) {
-            currentSchemaChangeEvents.addAll(schemaChangeEvents);
+        long schemaChangeEventsToSkip =
+                schemaChangeEventsToSkip(nextSplit, schemaChangeEvents.size());
+        for (int i = (int) schemaChangeEventsToSkip; i < 
schemaChangeEvents.size(); i++) {
+            currentSchemaChangeEvents.add(schemaChangeEvents.get(i));
         }
 
         if (currentDataRowsRead > 0) {
@@ -240,6 +242,22 @@ public class CDCSourceSplitReader
         }
     }
 
+    private long schemaChangeEventsToSkip(
+            TableAwareFileStoreSourceSplit split, int schemaChangeEventCount) 
throws IOException {
+        long schemaChangeEventsToSkip =
+                split.isLegacySchemaProgress() && split.recordsToSkip() > 0
+                        ? schemaChangeEventCount
+                        : split.schemaChangeEventsToSkip();
+        if (schemaChangeEventsToSkip < 0 || schemaChangeEventsToSkip > 
schemaChangeEventCount) {
+            throw new IOException(
+                    String.format(
+                            "Invalid schema change event skip count %s for 
split %s. "
+                                    + "The split has only %s schema change 
events.",
+                            schemaChangeEventsToSkip, split.splitId(), 
schemaChangeEventCount));
+        }
+        return schemaChangeEventsToSkip;
+    }
+
     @VisibleForTesting
     protected LazyRecordReader createLazyRecordReader(Split split) {
         return new LazyRecordReader(split, currentTableReaderInfo, 
tableManager);
diff --git 
a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitState.java
 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitState.java
new file mode 100644
index 0000000000..9abb365549
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitState.java
@@ -0,0 +1,78 @@
+/*
+ * 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.paimon.flink.pipeline.cdc.source.reader;
+
+import 
org.apache.paimon.flink.pipeline.cdc.source.TableAwareFileStoreSourceSplit;
+
+import org.apache.flink.cdc.common.event.Event;
+import org.apache.flink.cdc.common.event.SchemaChangeEvent;
+import org.apache.flink.connector.file.src.util.CheckpointedPosition;
+import org.apache.flink.connector.file.src.util.RecordAndPosition;
+
+import static org.apache.paimon.utils.Preconditions.checkArgument;
+import static org.apache.paimon.utils.Preconditions.checkNotNull;
+
+/** State of the CDC source reader with independent data row and schema event 
progress. */
+public class CDCSourceSplitState {
+
+    private final TableAwareFileStoreSourceSplit split;
+
+    private long recordsToSkip;
+
+    private long schemaChangeEventsToSkip;
+
+    public CDCSourceSplitState(TableAwareFileStoreSourceSplit split) {
+        this.split = checkNotNull(split);
+        this.recordsToSkip = split.recordsToSkip();
+        this.schemaChangeEventsToSkip = split.schemaChangeEventsToSkip();
+    }
+
+    public void setPosition(RecordAndPosition<Event> position) {
+        checkArgument(position.getOffset() == CheckpointedPosition.NO_OFFSET);
+        if (position.getRecord() instanceof SchemaChangeEvent) {
+            schemaChangeEventsToSkip++;
+        } else {
+            recordsToSkip = position.getRecordSkipCount();
+        }
+    }
+
+    public long recordsToSkip() {
+        return recordsToSkip;
+    }
+
+    public long schemaChangeEventsToSkip() {
+        return schemaChangeEventsToSkip;
+    }
+
+    public TableAwareFileStoreSourceSplit toSourceSplit() {
+        return split.updateWithProgress(recordsToSkip, 
schemaChangeEventsToSkip);
+    }
+
+    @Override
+    public String toString() {
+        return "CDCSourceSplitState{"
+                + "split="
+                + split
+                + ", recordsToSkip="
+                + recordsToSkip
+                + ", schemaChangeEventsToSkip="
+                + schemaChangeEventsToSkip
+                + '}';
+    }
+}
diff --git 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/TableAwareFileStoreSourceSplitSerializerTest.java
 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/TableAwareFileStoreSourceSplitSerializerTest.java
index 79ccabf518..1da49aaf8f 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/TableAwareFileStoreSourceSplitSerializerTest.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/TableAwareFileStoreSourceSplitSerializerTest.java
@@ -22,10 +22,14 @@ import org.apache.paimon.catalog.Identifier;
 import org.apache.paimon.flink.source.FileStoreSourceSplit;
 import org.apache.paimon.flink.source.FileStoreSourceSplitState;
 import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.utils.InstantiationUtil;
+import org.apache.paimon.utils.JsonSerdeUtil;
 
 import org.apache.flink.connector.file.src.util.RecordAndPosition;
+import org.apache.flink.core.memory.DataOutputViewStreamWrapper;
 import org.junit.jupiter.api.Test;
 
+import java.io.ByteArrayOutputStream;
 import java.util.Arrays;
 
 import static 
org.apache.paimon.flink.source.FileStoreSourceSplitSerializerTest.newFile;
@@ -42,7 +46,7 @@ public class TableAwareFileStoreSourceSplitSerializerTest {
         Identifier identifier = Identifier.create("test_database", 
"test_table");
         TableAwareFileStoreSourceSplit split =
                 new TableAwareFileStoreSourceSplit(
-                        "split-1", newDataSplit(), 0L, identifier, null, 1L);
+                        "split-1", newDataSplit(), 0L, identifier, null, 1L, 
2L);
 
         TableAwareFileStoreSourceSplit.Serializer serializer =
                 new TableAwareFileStoreSourceSplit.Serializer();
@@ -52,6 +56,27 @@ public class TableAwareFileStoreSourceSplitSerializerTest {
         assertThat(deserialized).isEqualTo(split);
     }
 
+    @Test
+    public void testDeserializeVersion1() throws Exception {
+        Identifier identifier = Identifier.create("test_database", 
"test_table");
+        DataSplit dataSplit = newDataSplit();
+
+        TableAwareFileStoreSourceSplit.Serializer serializer =
+                new TableAwareFileStoreSourceSplit.Serializer();
+        TableAwareFileStoreSourceSplit deserialized =
+                serializer.deserialize(
+                        1, serializeVersion1("split-1", dataSplit, 3L, 
identifier, null, 1L));
+
+        assertThat(deserialized.splitId()).isEqualTo("split-1");
+        assertThat(deserialized.split()).isEqualTo(dataSplit);
+        assertThat(deserialized.recordsToSkip()).isEqualTo(3L);
+        assertThat(deserialized.getIdentifier()).isEqualTo(identifier);
+        assertThat(deserialized.getLastSchemaId()).isNull();
+        assertThat(deserialized.getSchemaId()).isEqualTo(1L);
+        assertThat(deserialized.schemaChangeEventsToSkip()).isEqualTo(0L);
+        assertThat(deserialized.isLegacySchemaProgress()).isTrue();
+    }
+
     @Test
     public void testUpdateWithRecordsToSkipKeepsTableAwareSplit() {
         Identifier identifier = Identifier.create("test_database", 
"test_table");
@@ -74,6 +99,23 @@ public class TableAwareFileStoreSourceSplitSerializerTest {
         assertThat(tableAwareRestored.getSchemaId()).isEqualTo(2L);
     }
 
+    @Test
+    public void testUpdateWithRecordsToSkipPreservesSchemaProgress() throws 
Exception {
+        Identifier identifier = Identifier.create("test_database", 
"test_table");
+        DataSplit dataSplit = newDataSplit();
+        TableAwareFileStoreSourceSplit.Serializer serializer =
+                new TableAwareFileStoreSourceSplit.Serializer();
+        TableAwareFileStoreSourceSplit split =
+                serializer.deserialize(
+                        1, serializeVersion1("split-1", dataSplit, 3L, 
identifier, 1L, 2L));
+
+        TableAwareFileStoreSourceSplit updated = 
split.updateWithRecordsToSkip(10L);
+
+        assertThat(updated.recordsToSkip()).isEqualTo(10L);
+        assertThat(updated.schemaChangeEventsToSkip()).isEqualTo(0L);
+        assertThat(updated.isLegacySchemaProgress()).isTrue();
+    }
+
     private static DataSplit newDataSplit() {
         return DataSplit.builder()
                 .withSnapshot(1)
@@ -85,4 +127,23 @@ public class TableAwareFileStoreSourceSplitSerializerTest {
                 .withBucketPath("/temp/2") // not used
                 .build();
     }
+
+    private static byte[] serializeVersion1(
+            String splitId,
+            DataSplit split,
+            long recordsToSkip,
+            Identifier identifier,
+            Long lastSchemaId,
+            long schemaId)
+            throws Exception {
+        ByteArrayOutputStream out = new ByteArrayOutputStream();
+        DataOutputViewStreamWrapper view = new 
DataOutputViewStreamWrapper(out);
+        view.writeUTF(splitId);
+        InstantiationUtil.serializeObject(view, split);
+        view.writeLong(recordsToSkip);
+        view.writeUTF(JsonSerdeUtil.toJson(identifier));
+        view.writeLong(lastSchemaId == null ? -1L : lastSchemaId);
+        view.writeLong(schemaId);
+        return out.toByteArray();
+    }
 }
diff --git 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/enumerator/CDCCheckpointSerializerTest.java
 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/enumerator/CDCCheckpointSerializerTest.java
index 38388b12cf..4a70858cd1 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/enumerator/CDCCheckpointSerializerTest.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/enumerator/CDCCheckpointSerializerTest.java
@@ -22,9 +22,9 @@ import org.apache.paimon.catalog.Identifier;
 import 
org.apache.paimon.flink.pipeline.cdc.source.TableAwareFileStoreSourceSplit;
 import 
org.apache.paimon.flink.pipeline.cdc.source.enumerator.CDCCheckpoint.TableProgress;
 import org.apache.paimon.table.source.DataSplit;
+import org.apache.paimon.utils.InstantiationUtil;
 import org.apache.paimon.utils.JsonSerdeUtil;
 
-import org.apache.flink.core.io.SimpleVersionedSerializer;
 import org.apache.flink.core.memory.DataOutputViewStreamWrapper;
 import org.junit.jupiter.api.Test;
 
@@ -94,7 +94,15 @@ public class CDCCheckpointSerializerTest {
         CDCCheckpoint.Serializer serializer = new CDCCheckpoint.Serializer();
         CDCCheckpoint checkpoint = serializer.deserialize(1, version1Bytes);
 
-        assertThat(checkpoint.getSplits()).containsExactly(split);
+        TableAwareFileStoreSourceSplit restoredSplit = 
checkpoint.getSplits().iterator().next();
+        assertThat(restoredSplit.splitId()).isEqualTo(split.splitId());
+        assertThat(restoredSplit.split()).isEqualTo(split.split());
+        
assertThat(restoredSplit.recordsToSkip()).isEqualTo(split.recordsToSkip());
+        
assertThat(restoredSplit.getIdentifier()).isEqualTo(split.getIdentifier());
+        
assertThat(restoredSplit.getLastSchemaId()).isEqualTo(split.getLastSchemaId());
+        assertThat(restoredSplit.getSchemaId()).isEqualTo(split.getSchemaId());
+        assertThat(restoredSplit.schemaChangeEventsToSkip()).isEqualTo(0L);
+        assertThat(restoredSplit.isLegacySchemaProgress()).isTrue();
         assertThat(checkpoint.getTableProgressMap())
                 .containsEntry(identifier, new TableProgress(3L, null));
     }
@@ -104,12 +112,10 @@ public class CDCCheckpointSerializerTest {
             throws Exception {
         ByteArrayOutputStream out = new ByteArrayOutputStream();
         DataOutputViewStreamWrapper view = new 
DataOutputViewStreamWrapper(out);
-        SimpleVersionedSerializer<TableAwareFileStoreSourceSplit> 
splitSerializer =
-                new TableAwareFileStoreSourceSplit.Serializer();
 
         view.writeInt(splits.size());
         for (TableAwareFileStoreSourceSplit split : splits) {
-            byte[] bytes = splitSerializer.serialize(split);
+            byte[] bytes = serializeSplitVersion1(split);
             view.writeInt(bytes.length);
             view.write(bytes);
         }
@@ -119,4 +125,16 @@ public class CDCCheckpointSerializerTest {
         view.writeLong(nextSnapshotId);
         return out.toByteArray();
     }
+
+    private byte[] serializeSplitVersion1(TableAwareFileStoreSourceSplit 
split) throws Exception {
+        ByteArrayOutputStream out = new ByteArrayOutputStream();
+        DataOutputViewStreamWrapper view = new 
DataOutputViewStreamWrapper(out);
+        view.writeUTF(split.splitId());
+        InstantiationUtil.serializeObject(view, split.split());
+        view.writeLong(split.recordsToSkip());
+        view.writeUTF(JsonSerdeUtil.toJson(split.getIdentifier()));
+        view.writeLong(split.getLastSchemaId() == null ? -1L : 
split.getLastSchemaId());
+        view.writeLong(split.getSchemaId());
+        return out.toByteArray();
+    }
 }
diff --git 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitReaderTest.java
 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitReaderTest.java
index 8461ed2a66..27ae504f21 100644
--- 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitReaderTest.java
+++ 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitReaderTest.java
@@ -41,6 +41,8 @@ import org.apache.paimon.schema.TableSchema;
 import org.apache.paimon.table.source.DataSplit;
 import org.apache.paimon.table.source.Split;
 import org.apache.paimon.table.source.TableRead;
+import org.apache.paimon.utils.InstantiationUtil;
+import org.apache.paimon.utils.JsonSerdeUtil;
 import org.apache.paimon.utils.RecordWriter;
 
 import org.apache.flink.api.java.tuple.Tuple2;
@@ -57,6 +59,7 @@ import 
org.apache.flink.connector.base.source.reader.splitreader.SplitsChange;
 import org.apache.flink.connector.file.src.reader.BulkFormat;
 import org.apache.flink.connector.file.src.reader.BulkFormat.RecordIterator;
 import org.apache.flink.connector.file.src.util.RecordAndPosition;
+import org.apache.flink.core.memory.DataOutputViewStreamWrapper;
 import org.apache.flink.table.types.logical.BigIntType;
 import org.apache.flink.table.types.logical.RowType;
 import org.apache.flink.types.RowKind;
@@ -66,6 +69,7 @@ import org.junit.jupiter.api.io.TempDir;
 
 import javax.annotation.Nullable;
 
+import java.io.ByteArrayOutputStream;
 import java.io.IOException;
 import java.util.ArrayList;
 import java.util.Arrays;
@@ -342,7 +346,7 @@ public class CDCSourceSplitReaderTest {
         List<DataFileMeta> files2 = rw.writeFiles(row(1), 0, input2);
         files.addAll(files2);
 
-        assignSplit(reader, newSourceSplit("id1", row(1), 0, files, 
input1.size()));
+        assignSplit(reader, newSourceSplit("id1", row(1), 0, files, false, 
input1.size(), 1L));
 
         RecordsWithSplitIds<BulkFormat.RecordIterator<Event>> records = 
reader.fetch();
         assertRecords(records, null, "id1", input1.size(), 
Collections.emptyList());
@@ -358,6 +362,95 @@ public class CDCSourceSplitReaderTest {
         reader.close();
     }
 
+    @Test
+    public void testRestoreWithPartialSchemaChangeEventsSkipsEmittedEvents() 
throws Exception {
+        TestChangelogDataReadWrite rw = new 
TestChangelogDataReadWrite(tablePath);
+        CDCSourceSplitReader reader =
+                createReader(rw.createReadWithKey(), 
multipleSchemaChangeEvents());
+
+        List<Tuple2<Long, Long>> input = kvs();
+        List<DataFileMeta> files = rw.writeFiles(row(1), 0, input);
+
+        assignSplit(reader, newSourceSplit("id1", row(1), 0, files, false, 0L, 
1L));
+
+        RecordsWithSplitIds<BulkFormat.RecordIterator<Event>> records = 
reader.fetch();
+        assertThat(readEventTypes(records, "id1"))
+                .containsExactly(
+                        SchemaChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class);
+
+        reader.close();
+    }
+
+    @Test
+    public void testRestoreWithAllSchemaChangeEventsSkippedReadsDataRows() 
throws Exception {
+        TestChangelogDataReadWrite rw = new 
TestChangelogDataReadWrite(tablePath);
+        CDCSourceSplitReader reader =
+                createReader(rw.createReadWithKey(), 
multipleSchemaChangeEvents());
+
+        List<Tuple2<Long, Long>> input = kvs();
+        List<DataFileMeta> files = rw.writeFiles(row(1), 0, input);
+
+        assignSplit(reader, newSourceSplit("id1", row(1), 0, files, false, 0L, 
2L));
+
+        RecordsWithSplitIds<BulkFormat.RecordIterator<Event>> records = 
reader.fetch();
+        assertThat(readEventTypes(records, "id1"))
+                .containsExactly(
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class);
+
+        reader.close();
+    }
+
+    @Test
+    public void testRestoreWithInvalidSchemaChangeEventSkipCountFails() throws 
Exception {
+        TestChangelogDataReadWrite rw = new 
TestChangelogDataReadWrite(tablePath);
+        CDCSourceSplitReader reader = createReader(rw.createReadWithKey(), 
schemaChangeEvents());
+
+        List<Tuple2<Long, Long>> input = kvs();
+        List<DataFileMeta> files = rw.writeFiles(row(1), 0, input);
+
+        assignSplit(reader, newSourceSplit("id1", row(1), 0, files, false, 0L, 
2L));
+
+        assertThatThrownBy(reader::fetch)
+                .hasMessageContaining("Invalid schema change event skip count 
2");
+
+        reader.close();
+    }
+
+    @Test
+    public void 
testRestoreFromLegacySplitWithDataProgressSkipsSchemaChangeEvents()
+            throws Exception {
+        TestChangelogDataReadWrite rw = new 
TestChangelogDataReadWrite(tablePath);
+        CDCSourceSplitReader reader = createReader(rw.createReadWithKey(), 
schemaChangeEvents());
+
+        List<Tuple2<Long, Long>> input = kvs();
+        List<DataFileMeta> files = rw.writeFiles(row(1), 0, input);
+
+        TableAwareFileStoreSourceSplit split = newSourceSplit("id1", row(1), 
0, files, 1L);
+        assignSplit(reader, legacySplit(split));
+
+        RecordsWithSplitIds<BulkFormat.RecordIterator<Event>> records = 
reader.fetch();
+        assertThat(readEventTypes(records, "id1"))
+                .containsExactly(
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class,
+                        DataChangeEvent.class);
+
+        reader.close();
+    }
+
     @Test
     public void testRestore() throws Exception {
         TestChangelogDataReadWrite rw = new 
TestChangelogDataReadWrite(tablePath);
@@ -672,6 +765,26 @@ public class CDCSourceSplitReaderTest {
                                                         .BIGINT())))));
     }
 
+    private List<SchemaChangeEvent> multipleSchemaChangeEvents() {
+        return Arrays.asList(
+                new AddColumnEvent(
+                        TableId.tableId(DATABASE, TABLE),
+                        Collections.singletonList(
+                                AddColumnEvent.last(
+                                        Column.physicalColumn(
+                                                "extra_1",
+                                                
org.apache.flink.cdc.common.types.DataTypes
+                                                        .BIGINT())))),
+                new AddColumnEvent(
+                        TableId.tableId(DATABASE, TABLE),
+                        Collections.singletonList(
+                                AddColumnEvent.last(
+                                        Column.physicalColumn(
+                                                "extra_2",
+                                                
org.apache.flink.cdc.common.types.DataTypes
+                                                        .BIGINT())))));
+    }
+
     private List<Tuple2<Long, Long>> kvs() {
         return kvs(0);
     }
@@ -739,6 +852,17 @@ public class CDCSourceSplitReaderTest {
             List<DataFileMeta> files,
             boolean isIncremental,
             long recordsToSkip) {
+        return newSourceSplit(id, partition, bucket, files, isIncremental, 
recordsToSkip, 0L);
+    }
+
+    public static TableAwareFileStoreSourceSplit newSourceSplit(
+            String id,
+            BinaryRow partition,
+            int bucket,
+            List<DataFileMeta> files,
+            boolean isIncremental,
+            long recordsToSkip,
+            long schemaChangeEventsToSkip) {
         DataSplit split =
                 DataSplit.builder()
                         .withSnapshot(1)
@@ -750,7 +874,26 @@ public class CDCSourceSplitReaderTest {
                         .withBucketPath("/temp/" + bucket) // no used
                         .build();
         return new TableAwareFileStoreSourceSplit(
-                id, split, recordsToSkip, Identifier.create(DATABASE, TABLE), 
1L, 1L);
+                id,
+                split,
+                recordsToSkip,
+                Identifier.create(DATABASE, TABLE),
+                1L,
+                1L,
+                schemaChangeEventsToSkip);
+    }
+
+    private TableAwareFileStoreSourceSplit 
legacySplit(TableAwareFileStoreSourceSplit split)
+            throws Exception {
+        ByteArrayOutputStream out = new ByteArrayOutputStream();
+        DataOutputViewStreamWrapper view = new 
DataOutputViewStreamWrapper(out);
+        view.writeUTF(split.splitId());
+        InstantiationUtil.serializeObject(view, split.split());
+        view.writeLong(split.recordsToSkip());
+        view.writeUTF(JsonSerdeUtil.toJson(split.getIdentifier()));
+        view.writeLong(split.getLastSchemaId() == null ? -1L : 
split.getLastSchemaId());
+        view.writeLong(split.getSchemaId());
+        return new TableAwareFileStoreSourceSplit.Serializer().deserialize(1, 
out.toByteArray());
     }
 
     private static class TestCDCSourceSplitReader extends CDCSourceSplitReader 
{
diff --git 
a/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitStateTest.java
 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitStateTest.java
new file mode 100644
index 0000000000..cd545637eb
--- /dev/null
+++ 
b/paimon-flink/paimon-flink-cdc/src/test/java/org/apache/paimon/flink/pipeline/cdc/source/reader/CDCSourceSplitStateTest.java
@@ -0,0 +1,118 @@
+/*
+ * 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.paimon.flink.pipeline.cdc.source.reader;
+
+import org.apache.paimon.catalog.Identifier;
+import 
org.apache.paimon.flink.pipeline.cdc.source.TableAwareFileStoreSourceSplit;
+import org.apache.paimon.table.source.DataSplit;
+
+import org.apache.flink.cdc.common.event.AddColumnEvent;
+import org.apache.flink.cdc.common.event.DataChangeEvent;
+import org.apache.flink.cdc.common.event.Event;
+import org.apache.flink.cdc.common.event.TableId;
+import org.apache.flink.cdc.common.schema.Column;
+import org.apache.flink.cdc.common.types.DataTypes;
+import org.apache.flink.cdc.runtime.typeutils.BinaryRecordDataGenerator;
+import org.apache.flink.connector.file.src.util.RecordAndPosition;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+
+import static 
org.apache.paimon.flink.source.FileStoreSourceSplitSerializerTest.newFile;
+import static org.apache.paimon.io.DataFileTestUtils.row;
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link CDCSourceSplitState}. */
+public class CDCSourceSplitStateTest {
+
+    @Test
+    public void testSchemaChangeEventAdvancesOnlySchemaProgress() {
+        CDCSourceSplitState state = new CDCSourceSplitState(newSplit(3L, 1L));
+
+        state.setPosition(
+                new RecordAndPosition<Event>(schemaChangeEvent(), 
RecordAndPosition.NO_OFFSET, 3L));
+
+        assertThat(state.recordsToSkip()).isEqualTo(3L);
+        assertThat(state.schemaChangeEventsToSkip()).isEqualTo(2L);
+    }
+
+    @Test
+    public void testDataChangeEventAdvancesOnlyDataProgress() {
+        CDCSourceSplitState state = new CDCSourceSplitState(newSplit(3L, 1L));
+
+        state.setPosition(
+                new RecordAndPosition<Event>(dataChangeEvent(), 
RecordAndPosition.NO_OFFSET, 4L));
+
+        assertThat(state.recordsToSkip()).isEqualTo(4L);
+        assertThat(state.schemaChangeEventsToSkip()).isEqualTo(1L);
+    }
+
+    @Test
+    public void testToSourceSplitPreservesProgress() {
+        CDCSourceSplitState state = new CDCSourceSplitState(newSplit(3L, 1L));
+
+        state.setPosition(
+                new RecordAndPosition<Event>(schemaChangeEvent(), 
RecordAndPosition.NO_OFFSET, 3L));
+        state.setPosition(
+                new RecordAndPosition<Event>(dataChangeEvent(), 
RecordAndPosition.NO_OFFSET, 4L));
+
+        TableAwareFileStoreSourceSplit split = state.toSourceSplit();
+        assertThat(split.recordsToSkip()).isEqualTo(4L);
+        assertThat(split.schemaChangeEventsToSkip()).isEqualTo(2L);
+    }
+
+    private static TableAwareFileStoreSourceSplit newSplit(
+            long recordsToSkip, long schemaChangeEventsToSkip) {
+        DataSplit dataSplit =
+                DataSplit.builder()
+                        .withSnapshot(1)
+                        .withPartition(row(1))
+                        .withBucket(2)
+                        .withDataFiles(Arrays.asList(newFile(0), newFile(1)))
+                        .isStreaming(false)
+                        .rawConvertible(false)
+                        .withBucketPath("/temp/2") // not used
+                        .build();
+        return new TableAwareFileStoreSourceSplit(
+                "split-1",
+                dataSplit,
+                recordsToSkip,
+                Identifier.create("test_database", "test_table"),
+                1L,
+                2L,
+                schemaChangeEventsToSkip);
+    }
+
+    private static AddColumnEvent schemaChangeEvent() {
+        return new AddColumnEvent(
+                TableId.tableId("test_database", "test_table"),
+                Collections.singletonList(
+                        AddColumnEvent.last(Column.physicalColumn("extra", 
DataTypes.BIGINT()))));
+    }
+
+    private static DataChangeEvent dataChangeEvent() {
+        BinaryRecordDataGenerator generator =
+                new BinaryRecordDataGenerator(
+                        new org.apache.flink.cdc.common.types.DataType[] 
{DataTypes.BIGINT()});
+        return DataChangeEvent.insertEvent(
+                TableId.tableId("test_database", "test_table"),
+                generator.generate(new Object[] {1L}));
+    }
+}


Reply via email to