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