DanielLeens commented on PR #12321:
URL: https://github.com/apache/seatunnel/pull/12321#issuecomment-5678712310

   *Posting this as a plain issue comment rather than a `gh pr review` — GitHub 
does not let an author submit a formal review on their own PR. I reviewed this 
with the same scrutiny I'd apply to anyone else's CDC-restore change, including 
re-deriving the config plumbing from scratch rather than trusting my own PR 
description.*
   
   # What Problem Does This PR Solve?
   
   - **User pain point**: `IncrementalSourceReader.restoreCheckpointState` 
unconditionally restored the checkpointed `CatalogTable`(s) (or the legacy 
checkpoint row type) into the deserializer on every incremental-split restore, 
overwriting the schema discovered live from the database at reader start-up. If 
a column was added to the source table (e.g. `ALTER TABLE ... ADD COLUMN`) 
while the job was stopped, then the job restored from a savepoint/checkpoint 
taken *before* that DDL, the restore would reinstate the old, narrower 
checkpoint schema — and if the job does not propagate schema-change events 
downstream (`schema-changes.enabled=false`, which is the documented default), 
nothing would ever widen that schema back. The new column would be silently 
dropped from every produced row for the rest of the job's life. This regressed 
`OpengaussCDCIT#testAddFieldWithRestore` on `dev` (introduced by #11503, per my 
prior root-cause note on this failure).
   - **Fix approach**: Gate checkpoint-schema restoration (both the current 
`checkpointTables` path and the legacy `checkpointDataType` path) on whether 
the job actually propagates schema changes downstream. When disabled, skip 
restoring the checkpoint schema entirely and keep the schema discovered live at 
start-up — the same contract the reader had before checkpoint-schema 
restoration was introduced. When enabled, keep the existing restore behavior 
unchanged, since a DDL/RELATION-driven change stream can widen the schema again 
after restore. Debezium's own table history is restored in both cases, since it 
only drives change-stream decoding and is never widened by SeaTunnel itself.
   - **One-sentence summary**: Checkpoint schema is now only restored when the 
job can actually keep that schema up to date afterward (schema-change 
propagation enabled); otherwise the reader keeps the live-discovered schema so 
newly added columns are never silently and permanently dropped.
   
   # 1. Code Change Review
   
   ## 1.1 Core Logic Analysis
   
   **Core changes**: 
`connector-cdc-base/.../source/reader/IncrementalSourceReader.java` — 
`restoreCheckpointState` gains a `boolean schemaChangeEnabled` parameter; new 
helper `isSchemaChangeEnabled(SourceConfig)`.
   
   Before:
   ```java
   static <T> void restoreCheckpointState(
           IncrementalSplit incrementalSplit,
           DebeziumDeserializationSchema<T> debeziumDeserializationSchema) {
       List<CatalogTable> checkpointTables = 
incrementalSplit.getCheckpointTables();
       if (checkpointTables != null && !checkpointTables.isEmpty()) {
           ...
           
debeziumDeserializationSchema.restoreCheckpointProducedType(checkpointTables);
       } else if (incrementalSplit.getCheckpointDataType() != null) {
           ... // legacy path, same unconditional restore
       }
       // history table changes restored unconditionally
   }
   ```
   After:
   ```java
   static <T> void restoreCheckpointState(
           IncrementalSplit incrementalSplit,
           DebeziumDeserializationSchema<T> debeziumDeserializationSchema,
           boolean schemaChangeEnabled) {
       List<CatalogTable> checkpointTables = 
incrementalSplit.getCheckpointTables();
       if (!schemaChangeEnabled) {
           if ((checkpointTables != null && !checkpointTables.isEmpty())
                   || incrementalSplit.getCheckpointDataType() != null) {
               log.info("... schema change propagation is disabled ... live 
discovered schema is kept ...");
           }
       } else if (checkpointTables != null && !checkpointTables.isEmpty()) {
           ...
           
debeziumDeserializationSchema.restoreCheckpointProducedType(checkpointTables);
       } else if (incrementalSplit.getCheckpointDataType() != null) {
           ... // legacy path, unchanged when schemaChangeEnabled == true
       }
       // history table changes still restored unconditionally
   }
   
   static boolean isSchemaChangeEnabled(SourceConfig sourceConfig) {
       if (sourceConfig instanceof JdbcSourceConfig) {
           return ((JdbcSourceConfig) 
sourceConfig).getDbzConnectorConfig().isSchemaChangesHistoryEnabled();
       }
       return false;
   }
   ```
   Call site (`createSplitState`): `restoreCheckpointState(incrementalSplit, 
debeziumDeserializationSchema, isSchemaChangeEnabled(sourceConfig));`
   
   **Key findings**:
   - The normal path reaches this on every incremental-split restore for every 
JDBC-relational CDC connector (MySQL, PostgreSQL/OpenGauss, Oracle, SQL Server, 
Db2, MongoDB — all extend the shared 
`IncrementalSource`/`IncrementalSourceReader`), which is the 
checkpoint/savepoint recovery path, not a rare corner case — it fires on every 
job restart that resumes an incremental split.
   - I independently re-derived, rather than trusted, that 
`isSchemaChangeEnabled` genuinely mirrors the `schema-changes.enabled` option: 
for MySQL, PostgreSQL, Oracle, and SQL Server, each connector's own 
`SourceConfigFactory` sets Debezium's `include.schema.changes` property from 
the same `schemaChangeEnabled` field that is populated from 
`SourceOptions.SCHEMA_CHANGES_ENABLED` (verified by grepping 
`props.setProperty(SCHEMA_CHANGE_KEY / "include.schema.changes", 
String.valueOf(schemaChangeEnabled))` in all four connectors' config 
factories). `isSchemaChangesHistoryEnabled()` is a real Debezium 
(`io.debezium.relational.RelationalDatabaseConnectorConfig`) API reading that 
same property back. So the new helper is reading the actual, effective switch, 
not a proxy that could drift from it.
   - `SCHEMA_CHANGES_ENABLED` defaults to `false` 
(`SourceOptions.java:120-124`) — meaning the vast majority of existing CDC jobs 
(anyone who never explicitly opted into `schema-changes.enabled=true`) were 
exposed to the pre-fix bug on every restore that crossed a DDL boundary. This 
makes the pre-fix defect high-severity in practice, and this fix a meaningful 
correctness restoration, not a narrow edge-case patch.
   - MongoDB's `SourceConfig` implementation (`MongodbSourceConfig`) does 
**not** implement `JdbcSourceConfig`, so `isSchemaChangeEnabled` correctly 
returns `false` for it unconditionally — meaning MongoDB (schemaless, with no 
DDL/RELATION-style schema-change propagation mechanism of its own) now always 
keeps the live-discovered schema on restore, which is exactly the safe default 
this PR's own principle calls for. I checked this is not accidental: 
`IncrementalSourceReader` is shared by `MongodbIncrementalSource extends 
IncrementalSource`, so MongoDB genuinely goes through this same code path, and 
the `instanceof JdbcSourceConfig` check is the correct discriminator (TiDB and 
Vitess, by contrast, do not go through this shared reader at all — they have no 
`IncrementalSource` subclass — so they are unaffected by this change either 
way).
   - This is a **precise, root-cause fix**, not a workaround: it ties the 
restore decision to the actual invariant that determines correctness (can this 
job ever widen a restored schema again?) rather than special-casing the 
specific failing test's symptom. The Javadoc's own framing — "it keeps the 
live-discovered schema, which is the contract it always had" — is accurate: 
this restores the pre-#11503 contract specifically for the case where it's 
unsafe to do otherwise, while preserving #11503's new behavior for the case 
where it's safe (propagation enabled).
   
   **In-depth correctness analysis**:
   - Debezium table-history restoration (`historyTableChanges`) is 
unconditional in both branches — verified this block sits after the `if 
(!schemaChangeEnabled) {...} else if (...) {...}` chain, not inside it. This 
matches the Javadoc's claim that history "only drives how the change stream 
itself is decoded and is never widened by SeaTunnel," and is confirmed by the 
new `restoreCheckpointStateKeepsLiveSchemaWhenSchemaChangesAreDisabled` test, 
which asserts `restoreCheckpointHistoryTableChanges` is still invoked while 
`restoreCheckpointProducedType` is not.
   - The legacy `checkpointDataType` path is subject to the identical gate 
(verified via the new 
`restoreCheckpointStateKeepsLiveSchemaForLegacyCheckpointWhenSchemaChangesAreDisabled`
 test) — this closes the same defect for jobs whose checkpoints predate the 
`checkpointTables` mechanism, not just the current-format ones.
   - No `IncrementalSplit` field or checkpoint serialization format is changed 
— the fix only changes how already-existing checkpoint content is *interpreted* 
at restore time, so old checkpoints remain fully readable; nothing needs 
migrating.
   - One thing I specifically checked and did not find a problem with: whether 
restoring the *live* schema instead of the *checkpoint-time* schema, in the 
`schemaChangeEnabled=true` case that is unchanged by this PR, could itself 
cause a decode mismatch against binlog/WAL positions that predate a schema 
change not yet replayed. That's out of scope for this PR (behavior there is 
unchanged from #11503), but I traced it far enough to be confident this PR does 
not make that scenario any better or worse — it only changes the 
`schemaChangeEnabled=false` branch.
   
   ## 1.2 Compatibility Impact
   
   **Fully compatible.** This is a bug-fix that restores the reader's original, 
pre-#11503 contract for the (default) `schema-changes.enabled=false` case; it 
does not remove, rename, or change the default of any config option, does not 
touch checkpoint/savepoint serialization format, and does not change any public 
API surface (`restoreCheckpointState`'s new parameter is a `static` 
package-private method, not a public API). For jobs that already had 
`schema-changes.enabled=true`, behavior is completely unchanged. For jobs on 
the (default) `schema-changes.enabled=false` path, this changes behavior in the 
direction of correctness — recovering columns that would otherwise have been 
silently and permanently dropped — which is squarely a bug fix, not a new 
incompatibility to document in `incompatible-changes.md`.
   
   ## 1.3 Performance / Side-Effect Analysis
   
   - `isSchemaChangeEnabled` is a single `instanceof` check plus, for the JDBC 
case, two cheap getter calls — evaluated once per split restore, not per row. 
Negligible cost.
   - The `schemaChangeEnabled=false` branch does strictly less work than before 
(skips `restoreCheckpointProducedType`/legacy-table resolution entirely, only 
logs), so if anything this is a small performance improvement on the restore 
path for the common (default) configuration, not a regression.
   - No new threading, locking, or resource-release concerns — this is 
synchronous, single-threaded restore-time logic with no I/O beyond what was 
already there (reading fields off the already-deserialized `IncrementalSplit`).
   
   ## 1.4 Error Handling and Logging
   
   The new `!schemaChangeEnabled` branch logs at `INFO` when a checkpoint 
actually carried a schema that is now being skipped (both 
`checkpointTables`-non-empty and legacy `checkpointDataType`-present cases), 
which is the right visibility for an operator trying to understand why a 
restored job's schema didn't change — it doesn't silently do nothing. No 
exceptions are swallowed; no sensitive information is logged (table/split 
identifiers only, consistent with the rest of this method's existing logging).
   
   No blocking or non-blocking issues found in the production code.
   
   # 2. Code Quality Assessment
   
   ## 2.1 Coding Standards
   
   Both new/changed methods carry thorough Javadoc explaining the *why*, not 
just the *what* — `restoreCheckpointState`'s Javadoc walks through exactly why 
the gate exists and what happens on each side of it, and 
`isSchemaChangeEnabled`'s Javadoc explicitly names the Debezium property it 
mirrors and which connectors gate on it, which is exactly the kind of 
"why/constraint" documentation this class of change needs and which I would 
have flagged as missing had it not been there. No wildcard imports, no dead 
code left behind, style consistent with the surrounding file.
   
   ## 2.2 Test Coverage and Test Stability
   
   Coverage is comprehensive for the branch this PR adds: the four pre-existing 
tests (checkpoint-tables restore, empty split, current-format restore, 
legacy-format restore) are all updated to pass `true` and continue to assert 
the unchanged (`schemaChangeEnabled=true`) behavior, and three new tests 
exercise the new `false` branch — checkpoint-tables-present-but-skipped (with 
history still restored), legacy-format-skipped, and a dedicated 
`isSchemaChangeEnabled` mirroring test using mocks for both the `true`/`false` 
Debezium-config cases and a non-`JdbcSourceConfig` case (confirming non-JDBC 
sources correctly report `false`). This directly covers the three code paths I 
traced above (current-format gate, legacy-format gate, and the config-plumbing 
helper itself), not just a single happy-path assertion.
   
   **Test-stability conclusion (Section 5.10.2)**: All tests (existing and new) 
are pure Mockito-based unit tests with no threads, no `Thread.sleep`, no 
containers, no timing or environment dependency, and no shared mutable state 
across tests. **Stability rating: Stable.**
   
   ## 2.3 Documentation Updates
   
   No `docs/en`/`docs/zh` update is included, and I don't think one is 
required: `schema-changes.enabled` is an existing, already-documented option 
whose contract is unchanged by this PR (its documented purpose — "send schema 
change events downstream" — was never "and also determines whether checkpoint 
schema is restored," so there's no existing doc claim this PR contradicts); 
this PR fixes an internal restore-time defect in how that flag's *absence* was 
handled, it does not add or rename anything user-facing. If a reviewer feels 
the CDC docs should explicitly call out this restore-time interaction as a 
documented consequence of `schema-changes.enabled`, I'm open to adding a short 
note, but I don't believe it rises to the same bar as the two other PRs in this 
batch (#12318/#12320), which changed observable output values for existing 
valid configurations.
   
   # 3. Architectural Soundness
   
   ## 3.1 Elegance of the Solution
   
   **Precise fix.** It conditions the restore decision on the actual capability 
(can the job widen a restored schema afterward?) rather than adding a special 
case for the specific regression test, and it applies uniformly to both the 
current and legacy checkpoint-schema formats and to every connector that goes 
through the shared reader, including correctly generalizing to non-JDBC 
(MongoDB) sources via the `instanceof` check's `false` default rather than 
requiring a MongoDB-specific carve-out.
   
   ## 3.2 Maintainability
   
   The gate is centralized in one method (`restoreCheckpointState`) and one 
small helper (`isSchemaChangeEnabled`); a future connector added to the 
`IncrementalSource` family automatically gets the safe default (`false`, keep 
live schema) unless it explicitly is a `JdbcSourceConfig` that forwards 
`include.schema.changes`, which is the correct fail-safe direction.
   
   ## 3.3 Extensibility
   
   If a future non-JDBC connector *does* gain its own schema-change propagation 
mechanism, `isSchemaChangeEnabled` would need a corresponding branch — the 
current `instanceof JdbcSourceConfig` check is honest about only covering the 
JDBC family today rather than pretending to be fully general, which I think is 
the right level of abstraction for what currently exists.
   
   ## 3.4 Historical-Version Compatibility
   
   Fully compatible, as discussed in 1.2 — no checkpoint format or config 
surface changes; this restores prior, correct behavior for the default 
configuration rather than introducing a new one. Jobs upgrading into this fix 
need no migration action; they simply stop losing columns on restore.
   
   # 4. Issue Summary
   
   No blocking or non-blocking issues found.
   
   # 5. Merge Recommendation
   
   ### Conclusion: Ready to merge
   
   1. **Blockers — must be fixed**: None.
   2. **Recommended fixes — non-blocking**: None.
   
   **Overall assessment**: I traced the full config chain from 
`schema-changes.enabled` through each of the four JDBC-based CDC connectors' 
`include.schema.changes` property-setting, through Debezium's 
`isSchemaChangesHistoryEnabled()`, to confirm the new gate reads the real, 
effective switch rather than a proxy — and confirmed the `false` default means 
most existing jobs were exposed to the pre-fix defect, making this a high-value 
correctness fix rather than a narrow edge case. I also specifically checked the 
non-JDBC (MongoDB) and non-`IncrementalSource` (TiDB/Vitess) connectors to make 
sure the `instanceof` discriminator generalizes safely rather than silently 
mishandling connectors outside the four I could directly verify, and found the 
fail-safe default (keep live schema) is applied correctly everywhere this 
reader is used. No checkpoint format, API, or config-option changes; test 
coverage directly exercises every branch this PR touches with deterministic, 
stable unit tests. I
  don't see a better alternative implementation for the stated scope of this 
fix — the one architectural question worth surfacing (whether 
`schemaChangeEnabled=true`'s existing restore-then-widen approach is itself 
fully safe against schema drift between the checkpoint offset and the 
live-discovered schema at restart) is explicitly out of scope for this PR and 
unchanged from #11503's already-shipped behavior.
   


-- 
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