DanielLeens opened a new pull request, #12321:
URL: https://github.com/apache/seatunnel/pull/12321

   ## Summary
   
   Since #11503 merged (2026-09-13 04:01Z), `dev`'s own `all-connectors-it-2` 
job has been red on every completed run with 
`OpengaussCDCIT#testAddFieldWithRestore` (`actual iterable was <null> at index 
[1][4]`); the same failure shows up on unrelated PRs (#11077, #12216, ...). The 
last green dev run for that job is `34581442965` (2026-09-11). Earlier red runs 
on 09-12/09-13 were masked by the minio Docker Hub outage (`IcebergSourceIT` 
failed first and banned the rest of the reactor), which is why this surfaced 
only after #12287.
   
   This is a behavioural regression in CDC savepoint restore, not a flaky test: 
a column added to the source table while the job was stopped is silently 
dropped from every row written after the restore.
   
   ## Root cause (log-confirmed on dev run `34821874806`, job `Run / 
all-connectors-it-2 (8)`)
   
   The test takes a savepoint, runs `ALTER TABLE ... ADD COLUMN f_big` on 
source and sink, inserts a row, and restores the job. The reader log shows what 
#11503's restore path does now:
   
   ```
   SeaTunnelRowDebeziumDeserializeSchema - Table[...opengauss_cdc_table_3] 
restore before: ROW<id INT,f_bytea BYTES,f_small SMALLINT,f_int INT,f_big 
BIGINT>
   SeaTunnelRowDebeziumDeserializeSchema - Table[...opengauss_cdc_table_3] 
restore after:  ROW<id INT,f_bytea BYTES,f_small SMALLINT,f_int INT>
   SeaTunnelRowDebeziumDeserializeSchema - Emit restored schema for table[...]: 
ROW<id INT,f_bytea BYTES,f_small SMALLINT,f_int INT>
   AbstractJdbcSinkWriter - Restore runtime schema for table 
...sink_opengauss_cdc_table_3 without applying physical DDL
   ```
   
   Execution chain:
   
   1. `IncrementalSourceReader#initializedState` -> `restoreCheckpointState` 
now calls 
`DebeziumDeserializationSchema#restoreCheckpointProducedType(checkpointTables)` 
whenever the split carries checkpoint tables (before #11503 this was gated on 
the deprecated `checkpointDataType`, which was never set, so the 
live-discovered schema was always kept).
   2. `SeaTunnelRowDebeziumDeserializeSchema#restoreCheckpointProducedType` 
replaces the schema discovered from the live database at startup (`restore 
before`, which already contains `f_big`) with the checkpoint schema (`restore 
after`, which predates the DDL) and queues a `RestoreTableSchemaEvent`.
   3. The JDBC sink adopts the narrowed runtime schema; every subsequent row is 
converted without `f_big`, so the sink column stays `NULL` and the test's 
`assertIterableEquals(source, sink)` fails at the new column of the inserted 
row.
   4. Nothing can widen the schema again: the job runs with the default 
`schema-changes.enabled = false`, and that same flag is what gates DDL emission 
in the MySQL connector (`include.schema.changes`) and the RELATION listener in 
the PostgreSQL connector (`PostgresSourceFetchTaskContext#configure`, 
`connectorConfig.isSchemaChangesHistoryEnabled()`). With propagation disabled 
the narrowed schema is permanent until the job is restarted without state.
   
   #11503 targeted the opposite situation: a failover of a job that does 
propagate schema changes, where events between the checkpoint offset and a 
later DDL must be decoded with the checkpointed schema and the DDL in the 
stream widens it afterwards. That reasoning only holds when the stream is 
allowed to carry schema changes.
   
   ## Fix
   
   `IncrementalSourceReader` (connector-cdc-base only):
   
   - `restoreCheckpointState` gains a `schemaChangeEnabled` argument. When it 
is false, the checkpoint schema (new `checkpointTables` format and the legacy 
`checkpointDataType`) is not restored and the live-discovered schema stays 
authoritative, which is exactly the pre-#11503 contract for these jobs; an INFO 
line records the decision. When it is true, #11503's behaviour is unchanged.
   - Debezium table history (`restoreCheckpointHistoryTableChanges`) is 
restored in both cases: it only drives decoding of the change stream itself and 
SeaTunnel never widens it, so it stays correct for events between the 
checkpoint offset and a DDL.
   - New package-private `isSchemaChangeEnabled(SourceConfig)` resolves the 
flag from 
`JdbcSourceConfig#getDbzConnectorConfig().isSchemaChangesHistoryEnabled()`, the 
very switch every JDBC CDC connector derives from `schema-changes.enabled`; 
non-JDBC sources (no schema change events at all) report false.
   
   Scope: one production file, one test file. No connector, option, checkpoint 
format, sink, or E2E change. `OpengaussCDCIT#testAddFieldWithRestore` is left 
exactly as it is; it is the E2E regression proof.
   
   ### Semantics
   
   | job configuration | before #11503 | after #11503 (dev today) | this PR |
   | --- | --- | --- | --- |
   | `schema-changes.enabled = false` (default), column added while stopped | 
live schema, column kept | checkpoint schema, column silently dropped | live 
schema, column kept |
   | `schema-changes.enabled = true`, failover / restore | live schema (bug 
#10461) | checkpoint schema, widened by the stream | unchanged from #11503 |
   
   ## Tests
   
   - `IncrementalSourceReaderTest`:
     - existing restore cases now pass `schemaChangeEnabled = true` and keep 
asserting #11503's behaviour;
     - `restoreCheckpointStateKeepsLiveSchemaWhenSchemaChangesAreDisabled`: 
checkpoint tables present, propagation disabled -> 
`restoreCheckpointProducedType` never called, history still restored;
     - 
`restoreCheckpointStateKeepsLiveSchemaForLegacyCheckpointWhenSchemaChangesAreDisabled`:
 same rule for the legacy `checkpointDataType` path;
     - `isSchemaChangeEnabledMirrorsDebeziumIncludeSchemaChanges`: JDBC config 
true/false and non-JDBC config.
   - E2E: `OpengaussCDCIT#testAddFieldWithRestore` (unchanged) covers the 
regression end to end; the MySQL CDC restore ITs that enable 
`schema-changes.enabled` keep exercising the #11503 path.
   - Local: `./mvnw spotless:apply -pl 
seatunnel-connectors-v2/connector-cdc/connector-cdc-base -nsu 
-Dmaven.gitcommitid.skip=true` only. Compilation, unit tests and E2E are 
verified by this PR's GitHub CI, per this initiative's no-local-build policy.
   
   ## Files
   
   - 
`seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java`
   - 
`seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReaderTest.java`
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to