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

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-12321-025d08f9bcb0523defa59566e7210a6ddff633af
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit c3f06f9a82c119e427df0fb1f440bae583148f62
Author: Daniel <[email protected]>
AuthorDate: Sun Sep 27 06:33:02 2026 +0000

    [Fix][CDC] Keep the live schema on restore when schema change propagation 
is disabled (#12321)
    
    Co-authored-by: DanielLeens <[email protected]>
    Co-authored-by: Claude Fable 5.1 <[email protected]>
---
 .../source/reader/IncrementalSourceReader.java     | 56 ++++++++++++-
 .../source/reader/IncrementalSourceReaderTest.java | 98 ++++++++++++++++++++--
 2 files changed, 146 insertions(+), 8 deletions(-)

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 2c64007b79..98b61593f6 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
@@ -24,6 +24,7 @@ import 
org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
 import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.type.MultipleRowType;
 import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.connectors.cdc.base.config.JdbcSourceConfig;
 import org.apache.seatunnel.connectors.cdc.base.config.SourceConfig;
 import org.apache.seatunnel.connectors.cdc.base.dialect.DataSourceDialect;
 import 
org.apache.seatunnel.connectors.cdc.base.source.event.CompletedSnapshotPhaseEvent;
@@ -238,7 +239,10 @@ public class IncrementalSourceReader<T, C extends 
SourceConfig>
             return new SnapshotSplitState(split.asSnapshotSplit());
         } else {
             IncrementalSplit incrementalSplit = split.asIncrementalSplit();
-            restoreCheckpointState(incrementalSplit, 
debeziumDeserializationSchema);
+            restoreCheckpointState(
+                    incrementalSplit,
+                    debeziumDeserializationSchema,
+                    isSchemaChangeEnabled(sourceConfig));
             IncrementalSplitState splitState = new 
IncrementalSplitState(incrementalSplit);
             if (splitState.autoEnterPureIncrementPhaseIfAllowed()) {
                 log.info(
@@ -255,11 +259,37 @@ public class IncrementalSourceReader<T, C extends 
SourceConfig>
         }
     }
 
+    /**
+     * Restores the deserializer's runtime schema and Debezium table history 
from a checkpointed
+     * incremental split.
+     *
+     * <p>The checkpointed runtime schema is only restored when {@code 
schemaChangeEnabled} is true.
+     * Restoring it replaces the schema discovered from the live database at 
startup, so any column
+     * added while the job was stopped disappears from the produced rows until 
the change stream
+     * widens the schema again; only a job that propagates schema changes 
(MySQL DDL events,
+     * PostgreSQL RELATION messages, all gated on {@code 
schema-changes.enabled}) can do that. A job
+     * with schema change propagation disabled would otherwise keep silently 
dropping such columns
+     * after every savepoint restore, so it keeps the live-discovered schema, 
which is the contract
+     * it always had. Debezium table history is restored in both cases: it 
only drives how the
+     * change stream itself is decoded and is never widened by SeaTunnel.
+     *
+     * @param incrementalSplit the split restored from checkpoint state
+     * @param debeziumDeserializationSchema the deserializer whose runtime 
state is restored
+     * @param schemaChangeEnabled whether this job propagates source schema 
changes downstream
+     */
     static <T> void restoreCheckpointState(
             IncrementalSplit incrementalSplit,
-            DebeziumDeserializationSchema<T> debeziumDeserializationSchema) {
+            DebeziumDeserializationSchema<T> debeziumDeserializationSchema,
+            boolean schemaChangeEnabled) {
         List<CatalogTable> checkpointTables = 
incrementalSplit.getCheckpointTables();
-        if (checkpointTables != null && !checkpointTables.isEmpty()) {
+        if (!schemaChangeEnabled) {
+            if ((checkpointTables != null && !checkpointTables.isEmpty())
+                    || incrementalSplit.getCheckpointDataType() != null) {
+                log.info(
+                        "The incremental split[{}] carries a checkpoint 
schema, but schema change propagation is disabled for this job, so the live 
discovered schema is kept instead of restoring the checkpoint schema.",
+                        incrementalSplit.splitId());
+            }
+        } else if (checkpointTables != null && !checkpointTables.isEmpty()) {
             log.info(
                     "The incremental split[{}] has {} checkpoint table(s) for 
restore: {}.",
                     incrementalSplit.splitId(),
@@ -293,6 +323,26 @@ public class IncrementalSourceReader<T, C extends 
SourceConfig>
         }
     }
 
+    /**
+     * Resolves whether the job propagates source schema changes downstream, 
i.e. the value of
+     * {@code schema-changes.enabled} as every JDBC-based CDC connector 
forwards it to Debezium's
+     * {@code include.schema.changes}. This is the same switch that gates DDL 
emission in the MySQL
+     * connector and the RELATION listener in the PostgreSQL connector, so it 
tells exactly whether
+     * a checkpoint-restored runtime schema can ever be widened again by the 
change stream. Non-JDBC
+     * sources never emit schema change events and therefore report false.
+     *
+     * @param sourceConfig the reader's source configuration
+     * @return true when schema change propagation is enabled for this job
+     */
+    static boolean isSchemaChangeEnabled(SourceConfig sourceConfig) {
+        if (sourceConfig instanceof JdbcSourceConfig) {
+            return ((JdbcSourceConfig) sourceConfig)
+                    .getDbzConnectorConfig()
+                    .isSchemaChangesHistoryEnabled();
+        }
+        return false;
+    }
+
     private static List<CatalogTable> restoreLegacyCheckpointTables(
             IncrementalSplit incrementalSplit) {
         if (incrementalSplit.getCheckpointDataType() instanceof 
MultipleRowType) {
diff --git 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReaderTest.java
 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReaderTest.java
index 6d29ddc13f..c7832751cb 100644
--- 
a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReaderTest.java
+++ 
b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReaderTest.java
@@ -29,6 +29,7 @@ import org.apache.seatunnel.api.table.type.BasicType;
 import org.apache.seatunnel.api.table.type.MultipleRowType;
 import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
 import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.connectors.cdc.base.config.JdbcSourceConfig;
 import org.apache.seatunnel.connectors.cdc.base.config.SourceConfig;
 import org.apache.seatunnel.connectors.cdc.base.dialect.DataSourceDialect;
 import org.apache.seatunnel.connectors.cdc.base.source.offset.Offset;
@@ -54,6 +55,7 @@ import org.junit.jupiter.api.Test;
 import org.mockito.ArgumentCaptor;
 import org.mockito.Mockito;
 
+import io.debezium.relational.RelationalDatabaseConnectorConfig;
 import io.debezium.relational.TableId;
 
 import java.io.IOException;
@@ -672,7 +674,7 @@ class IncrementalSourceReaderTest {
                         checkpointTables,
                         historyTableChanges);
 
-        IncrementalSourceReader.restoreCheckpointState(incrementalSplit, 
schema);
+        IncrementalSourceReader.restoreCheckpointState(incrementalSplit, 
schema, true);
 
         Mockito.verify(schema).restoreCheckpointProducedType(checkpointTables);
         
Mockito.verify(schema).restoreCheckpointHistoryTableChanges(historyTableChanges);
@@ -691,7 +693,7 @@ class IncrementalSourceReaderTest {
                         Mockito.mock(Offset.class),
                         Collections.emptyList());
 
-        IncrementalSourceReader.restoreCheckpointState(incrementalSplit, 
schema);
+        IncrementalSourceReader.restoreCheckpointState(incrementalSplit, 
schema, true);
 
         Mockito.verifyNoInteractions(schema);
     }
@@ -716,7 +718,7 @@ class IncrementalSourceReaderTest {
                         Collections.emptyList(),
                         checkpointRowType);
 
-        IncrementalSourceReader.restoreCheckpointState(split, schema);
+        IncrementalSourceReader.restoreCheckpointState(split, schema, true);
 
         @SuppressWarnings("unchecked")
         ArgumentCaptor<List<CatalogTable>> captor = 
ArgumentCaptor.forClass(List.class);
@@ -754,7 +756,7 @@ class IncrementalSourceReaderTest {
                         Collections.emptyList(),
                         checkpointRowType);
 
-        IncrementalSourceReader.restoreCheckpointState(split, schema);
+        IncrementalSourceReader.restoreCheckpointState(split, schema, true);
 
         @SuppressWarnings("unchecked")
         ArgumentCaptor<List<CatalogTable>> captor = 
ArgumentCaptor.forClass(List.class);
@@ -786,8 +788,94 @@ class IncrementalSourceReaderTest {
                         Collections.emptyList(),
                         checkpointRowType);
 
-        IncrementalSourceReader.restoreCheckpointState(split, schema);
+        IncrementalSourceReader.restoreCheckpointState(split, schema, true);
 
         Mockito.verifyNoInteractions(schema);
     }
+
+    /**
+     * Regression guard for the savepoint-restore column drop: a job that does 
not propagate schema
+     * changes (the default) must keep the schema discovered from the live 
database, because nothing
+     * in its change stream could ever widen a checkpoint schema that predates 
an ADD COLUMN
+     * executed while the job was stopped. Debezium history is still restored, 
since it only drives
+     * decoding of the stream itself.
+     */
+    @Test
+    void restoreCheckpointStateKeepsLiveSchemaWhenSchemaChangesAreDisabled() {
+        @SuppressWarnings("unchecked")
+        DebeziumDeserializationSchema<Object> schema =
+                Mockito.mock(DebeziumDeserializationSchema.class);
+        CatalogTable checkpointTable = Mockito.mock(CatalogTable.class);
+        Mockito.when(checkpointTable.getTablePath())
+                .thenReturn(TablePath.of("catalog", "database", "customers"));
+        Map<TableId, byte[]> historyTableChanges =
+                Collections.singletonMap(
+                        new TableId("catalog", "database", "customers"), new 
byte[] {1});
+        IncrementalSplit incrementalSplit =
+                new IncrementalSplit(
+                        "incremental-split-0",
+                        Collections.emptyList(),
+                        Mockito.mock(Offset.class),
+                        Mockito.mock(Offset.class),
+                        Collections.emptyList(),
+                        Collections.singletonList(checkpointTable),
+                        historyTableChanges);
+
+        IncrementalSourceReader.restoreCheckpointState(incrementalSplit, 
schema, false);
+
+        Mockito.verify(schema, 
Mockito.never()).restoreCheckpointProducedType(Mockito.anyList());
+        
Mockito.verify(schema).restoreCheckpointHistoryTableChanges(historyTableChanges);
+    }
+
+    /**
+     * The legacy checkpoint data type path is subject to the same rule: 
without schema change
+     * propagation the checkpointed row type must not replace the live schema.
+     */
+    @Test
+    void 
restoreCheckpointStateKeepsLiveSchemaForLegacyCheckpointWhenSchemaChangesAreDisabled()
 {
+        @SuppressWarnings("unchecked")
+        DebeziumDeserializationSchema<Object> schema =
+                Mockito.mock(DebeziumDeserializationSchema.class);
+        SeaTunnelRowType checkpointRowType =
+                new SeaTunnelRowType(
+                        new String[] {"id", "name"},
+                        new SeaTunnelDataType[] {BasicType.INT_TYPE, 
BasicType.STRING_TYPE});
+        IncrementalSplit split =
+                new IncrementalSplit(
+                        "incremental-split-0",
+                        Collections.singletonList(new TableId("catalog", 
"database", "customers")),
+                        Mockito.mock(Offset.class),
+                        Mockito.mock(Offset.class),
+                        Collections.emptyList(),
+                        checkpointRowType);
+
+        IncrementalSourceReader.restoreCheckpointState(split, schema, false);
+
+        Mockito.verifyNoInteractions(schema);
+    }
+
+    /**
+     * The switch must mirror what the connectors hand to Debezium as {@code
+     * include.schema.changes}, because that is the flag gating DDL emission 
(MySQL) and the
+     * RELATION listener (PostgreSQL); a non-JDBC source never emits schema 
change events.
+     */
+    @Test
+    void isSchemaChangeEnabledMirrorsDebeziumIncludeSchemaChanges() {
+        RelationalDatabaseConnectorConfig enabledConfig =
+                Mockito.mock(RelationalDatabaseConnectorConfig.class);
+        
Mockito.when(enabledConfig.isSchemaChangesHistoryEnabled()).thenReturn(true);
+        JdbcSourceConfig enabledSource = Mockito.mock(JdbcSourceConfig.class);
+        
Mockito.when(enabledSource.getDbzConnectorConfig()).thenReturn(enabledConfig);
+        
Assertions.assertTrue(IncrementalSourceReader.isSchemaChangeEnabled(enabledSource));
+
+        RelationalDatabaseConnectorConfig disabledConfig =
+                Mockito.mock(RelationalDatabaseConnectorConfig.class);
+        
Mockito.when(disabledConfig.isSchemaChangesHistoryEnabled()).thenReturn(false);
+        JdbcSourceConfig disabledSource = Mockito.mock(JdbcSourceConfig.class);
+        
Mockito.when(disabledSource.getDbzConnectorConfig()).thenReturn(disabledConfig);
+        
Assertions.assertFalse(IncrementalSourceReader.isSchemaChangeEnabled(disabledSource));
+
+        Assertions.assertFalse(
+                
IncrementalSourceReader.isSchemaChangeEnabled(Mockito.mock(SourceConfig.class)));
+    }
 }

Reply via email to