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]