DanielLeens commented on PR #11206:
URL: https://github.com/apache/seatunnel/pull/11206#issuecomment-5379306626
# What Problem Does This PR Solve?
- User pain: a MySQL CDC job configured with
`database-pattern`/`table-pattern` (wildcards) previously required a full job
restart to pick up a table created after the job started.
- Fix approach: two opt-in switches — `scan.newly-added-table.enabled`
(restore-time diff against checkpointed tables, snapshot backfill for new
matches) and `scan.binlog.newly-added-table.enabled` (turn the `CREATE TABLE`
Debezium schema-change record into a `CreateTableEvent`, register the
deserializer schema in-flight, and let `MultiTableSinkWriter` create a
per-table JDBC writer at runtime).
- One-sentence summary: the overall design is sound and this feature is
close, but the feature's own dedicated E2E test still fails deterministically
on the current head, and there are unresolved restore/checkpoint gaps on both
the source and sink side.
Housekeeping note, for transparency: nothing has changed source-wise since
my last comment on this same head (`c97e8583e855`) — same commit, same diff,
and `gh pr view 11206 --json headRefName,statusCheckRollup,updatedAt` confirms
the head SHA and the failing `Build` check are unchanged. Rather than re-post
the same conclusion verbatim, I re-derived the key evidence myself from scratch
this round — pulled and grepped the actual CI log, re-read the actual source
files line by line, and traced one additional path (`IncrementalSource`'s
pre-PR restore behavior) that resolves an open question from the previous
round. Below is the full review with that fresh evidence.
# 1. Code Change Review
## 1.1 Core Logic Analysis
**Independent re-verification of the primary blocker (E2E not passing).** I
pulled the actual fork CI log myself for this exact head (`gh api
repos/DanielLeens/seatunnel/actions/jobs/96389079475/logs`, job
`mysql-cdc-connector-it (11, ubuntu-latest)`, ~250k lines) rather than trusting
the earlier summary:
- `grep -c "Registered newly added CDC table" job96389079475.log` → `0`.
That log line is emitted only on the success branch of
`SeaTunnelRowDebeziumDeserializeSchema.handleTableChangeStruct()` (the exact
method this PR adds to wire binlog-discovered `CREATE TABLE` records into the
live deserializer), and it never fires once in the whole run.
- The test fails deterministically, at the same line, in every container
variant I found in the log (`Mysql8_4CDCIT`, `MysqlCDCIT`, ...):
`AbstractMysqlCDCITBase.java:868`, with
`org.awaitility.core.ConditionTimeoutException: ... newly added wildcard sink
table not readable yet: java.sql.SQLSyntaxErrorException: Table
'sink.source_payments' doesn't exist within 2 minutes.`
I also re-read the actual call chain myself to confirm the mechanism, not
just the symptom:
```text
AbstractDebeziumDeserializationSchema.deserialize() [:65-78]
if isSchemaChangeEvent(record):
for each tableChangeStruct in record's TABLE_CHANGES field:
handleTableChangeStruct(tableChangeStruct) <- overridden hook
SeaTunnelRowDebeziumDeserializeSchema.handleTableChangeStruct() [:127-157]
if !scanBinlogNewlyAddedTableEnabled: return
for each CREATE-type tableChange:
catalogTable = tableChangeCatalogTableConverter.convert(tableChange)
// null => filtered out
if catalogTable == null: return
if containsTable(...): return
tables.add(catalogTable); pendingCreateTables.add(catalogTable)
tableRowConverters = createTableRowConverters(...)
log.info("Registered newly added CDC table {}", ...) <- never observed
in CI
```
Both `AbstractDebeziumDeserializationSchema.deserialize()` and
`handleTableChangeStruct()` look internally correct on static reading, and
`MySqlSourceConfigFactory.java:104-107` correctly forces
`include.schema.changes=true` when `scanBinlogNewlyAddedTableEnabled` is set,
so Debezium should be emitting schema-change records at all. Since the log line
proves the method's success path never completes, the actual defect is upstream
of this method — either the schema-change `SourceRecord` for a table that did
not exist at snapshot time never reaches `deserialize()` for the binlog phase,
or `isSchemaChangeEvent(record)` never classifies it as one. That is exactly
consistent with my prior finding, now backed by a log I pulled and grepped
myself rather than taking on faith. I'd suggest the next concrete debugging
step is adding a temporary trace at the top of
`AbstractDebeziumDeserializationSchema.deserialize()` to confirm whether the
schema-change record for the runtime table arrives
there at all.
**New finding this round — resolves an open question from my last comment
and partially corrects @SEZ9's Issue 4.** @SEZ9 flagged that
`scan.newly-added-table.enabled` defaulting to `true` "appears to be new
behavior introduced by this PR" and could silently trigger unplanned snapshot
reads for existing jobs after upgrade. I diffed `IncrementalSource.java`
against the pre-PR base (`f1a1a0abb`) myself to settle this:
```diff
- Set<TableId> newTables = Sets.difference(capturedTables,
checkpointCapturedTables);
+ Set<TableId> newTables =
+ isScanNewlyAddedTableEnabledOnRestore()
+ ? Sets.difference(capturedTables, checkpointCapturedTables)
+ : Collections.emptySet();
```
The `Sets.difference(capturedTables, checkpointCapturedTables)` restore-time
reconciliation was **already unconditional pre-PR code** — every restart of
every CDC connector already re-diffed the freshly-discovered catalog against
the checkpointed table set and folded new matches into the snapshot phase. This
PR's only change here is wrapping that pre-existing behavior in a new
`isScanNewlyAddedTableEnabledOnRestore()` guard that defaults to `true`, i.e.
it makes an already-existing behavior *optionally disable-able* rather than
introducing it. So on the specific mechanism @SEZ9 raised, the docs' claim
("This keeps the existing SeaTunnel restore behavior by default",
`docs/en/connectors/source/MySQL-CDC.md:195`) is accurate, and I'd downgrade
this from "silently changes behavior for existing jobs" to a
documentation-clarity nit at most — worth a one-line callout in the option's
description that this was already the connector's restore behavior, but not a
compatibility blocker.
**Re-verified two other source-level claims by reading the code directly,
not by re-quoting the prior round:**
- `EventType.java:22` — confirmed `SCHEMA_CHANGE_CREATE_TABLE` is the first
enum constant, ahead of the eight pre-existing `SCHEMA_CHANGE_*` constants, so
@SEZ9's Issue 2 is factually accurate about the code. On the *impact* claim, I
independently traced one of the concrete serialization paths myself:
`JobEventReportOperation.writeInternal/readInternal`
(`seatunnel-engine/.../event/JobEventReportOperation.java:55-72`) serializes
the `List<Event>` through a plain
`java.io.ObjectOutputStream`/`ObjectInputStream` — per the JLS, default Java
enum serialization writes the constant's *name*, not its ordinal, so this
specific path is unaffected by the insertion position. I have not re-audited
every connector or every serializer in `seatunnel-engine`, so I can't rule out
an ordinal-sensitive path existing elsewhere, but on the evidence I could
independently confirm, I agree with the prior round's downgrade of this from
High to Low — insert-at-end is still the right defensive practic
e for enum evolution, just not a proven active bug today.
- `MultiTableSink.java` — confirmed the class does carry a new field,
`private final ReadonlyConfig options;` (`:84`), plus the
`FactoryUtil`/`TableSinkFactory`/`ConfigValidator` imports @SEZ9 cited. I did
not trace whether `MultiTableSink` instances actually cross a
`serialVersionUID`-sensitive restore boundary in Zeta's checkpoint path, so I
can't independently confirm or refute the severity of @SEZ9's Issue 3 this
round either — passing it through as open, same as last time.
## 1.2 Compatibility Impact
**Partially incompatible only within this feature's own new surface, fully
compatible for existing users.** Both new options default to values that
preserve pre-PR behavior for jobs that don't opt in:
`scan.binlog.newly-added-table.enabled=false` by default (a genuinely new code
path, gated off), and `scan.newly-added-table.enabled=true` by default, which —
per the diff evidence above — reproduces the connector's pre-existing
unconditional restore-time table-reconciliation behavior rather than
introducing a new default-on behavior. The new
`validateBinlogNewlyAddedTableConfig` guard only rejects a config combination
(`scan.binlog.newly-added-table.enabled=true` + fixed `table-names`) that could
not have existed before this PR (the option didn't exist), so it breaks no
running job.
## 1.3 Performance / Side-Effect Analysis
The capture-pattern filter added in an earlier round
(`MySqlIncrementalSource.java:280-298`, reusing Debezium's own
`dataCollectionFilter()`) is a per-DDL-event check, not per-row, and is
strictly cheaper than the pre-filter code for out-of-scope tables since it
skips a `MySqlCatalogTableUtils.toCatalogTable()` conversion. I have no new
performance concerns beyond what's already on record; @SEZ9's Issue 5 (no
cap/throttle on runtime writer registration) remains a fair open concern I have
no counter-evidence against — a burst of matching `CREATE TABLE` statements
really would allocate one JDBC writer per table with no ceiling.
## 1.4 Error Handling and Logging
No new swallowed exceptions in this round's delta.
`validateBinlogNewlyAddedTableConfig` still fails fast with an actionable
message. The one behavior worth flagging for logging: given the E2E evidence
above, the registration path fails *silently* from the user's perspective —
there is no WARN/ERROR anywhere on the path from "schema-change record
received" to "no table registered" when the runtime table never gets picked up,
which makes this class of failure hard to diagnose in production, not just in
CI.
# 2. Code Quality Assessment
## 2.1 Coding Standards
Unchanged from prior rounds: new methods
(`validateBinlogNewlyAddedTableConfig`,
`isScanNewlyAddedTableEnabledOnRestore`) carry doc comments, no wildcard
imports, ASF headers present on new files.
## 2.2 Test Coverage and Test Stability
**Rating: High risk.** The feature's own dedicated regression test,
`testMysqlCdcByWildcardsConfigWithNewlyAddedTable`
(`AbstractMysqlCDCITBase.java:868`), is confirmed failing on live CI for the
current head — I re-pulled and re-grepped the log myself this round rather than
relying on the previous round's citation, and got the identical result (0
registrations logged, deterministic failure across every container variant
present in the log). This is not an inference from a diff; it is the actual,
reproducible CI outcome for exactly the feature this PR adds. The
awaitility-based wait in that test is itself sound (condition-based polling
through the expected transient `SQLSyntaxErrorException`, not a disguised
sleep), so the "High risk" rating is about the underlying feature not
completing end-to-end, not about test flakiness hygiene.
## 2.3 Documentation Updates
@SEZ9's Issues 6 and 7 stand and I have no counter-evidence:
`docs/en/connectors/sink/Jdbc.md` is not updated despite
`JdbcSinkFactory`/`AbstractJdbcSinkWriter` gaining runtime newly-added-table
resolution logic, and neither the EN nor ZH MySQL-CDC page has a runnable HOCON
example for the two new options. On the one doc claim I did re-verify myself
(`scan.newly-added-table.enabled` default and its restore-behavior description,
`docs/en/connectors/source/MySQL-CDC.md:195,232`), the doc text is accurate
against the code.
# 3. Architectural Soundness
## 3.1 Elegance of the Solution
The design shape is right — hooking Debezium's own schema-change stream and
routing runtime tables through the existing `SupportMultiTableSinkWriter`
extension point is a sensible reuse of existing infrastructure rather than a
parallel mechanism. What's missing isn't a redesign, it's (a) finding why the
registration path doesn't actually fire in the E2E scenario, and (b) closing
the restore/checkpoint gaps on both sides of the pipeline.
## 3.2 Maintainability
The `null`-as-"filtered-out" convention through
`TableChangeCatalogTableConverter` remains a minor implicit-contract nit (a
reader has to know `null` means "intentionally excluded," not "conversion
failed"), consistently handled at every call site traced across rounds, so it's
a style note rather than a bug.
## 3.3 Extensibility
`isScanNewlyAddedTableEnabledOnRestore()` as an overridable protected hook
on the shared `IncrementalSource` base is a reasonable extension point for
other CDC connectors that want the same restore-time compatibility switch later.
## 3.4 Historical-Version Compatibility
No serialization/checkpoint format change in this round's delta. Two real
historical-compatibility-adjacent risks remain open and are the actual
blockers, not the `EventType` ordinal question: (1) a table registered purely
via the binlog `CREATE TABLE` path is an in-memory-only addition to the
deserializer's live `tables` list — it is never written into the enumerator's
checkpointed table set (`IncrementalSource.java:432-441`), so it can silently
drop out of a restored job; (2)
`SupportMultiTableSinkWriter.createSinkWriter(CatalogTable, Context)`
(`seatunnel-api/.../SupportMultiTableSinkWriter.java:54`) has no parameter for
restored writer state and no `restoreSinkWriter` counterpart, so a
runtime-created writer's pending state (e.g. buffered JDBC batches) cannot
survive a restart either. Both need to be closed together — fixing only one
side still leaves data-loss/duplication risk on restart for the tables this
feature is meant to add.
# 4. Issue Summary
| # | Issue | Location | Severity | Raised by another reviewer |
|---|-------|----------|----------|------------------------------|
| 1 | Feature's own E2E test fails deterministically on live CI for the
current head; the success-path log line for binlog table registration never
appears in the run | `AbstractMysqlCDCITBase.java:868`;
`SeaTunnelRowDebeziumDeserializeSchema.java` (`handleTableChangeStruct`) | High
| No (carryover, re-confirmed this round with a fresh independent log pull) |
| 2 | Binlog-registered tables are not part of the source's
checkpoint/restore state; lost on restart | `IncrementalSource.java:432-441` |
High | No (carryover) |
| 3 | Sink-side runtime-writer SPI has no restore variant; runtime-created
writer state cannot survive failure recovery |
`SupportMultiTableSinkWriter.java:54` | High | Yes (@SEZ9, Issue 1) |
| 4 | `MultiTableSink` gained factory/config dependencies and a new field
(`options`); serialVersionUID/classloader risk not independently confirmed
either way | `MultiTableSink.java:23-43,84` | Medium (unconfirmed) | Yes
(@SEZ9, Issue 3) |
| 5 | No cap/throttle on runtime writer registration |
`SupportMultiTableSinkWriter.java:43` | Medium | Yes (@SEZ9, Issue 5) |
| 6 | JDBC sink docs not updated despite gaining runtime table-creation
behavior | `docs/en/connectors/sink/Jdbc.md` (not touched) | Medium | Yes
(@SEZ9, Issue 6) |
| 7 | No example HOCON config for the new options |
`docs/en/connectors/source/MySQL-CDC.md:196` | Medium | Yes (@SEZ9, Issue 7) |
| 8 | `scan.newly-added-table.enabled` default-value description could be
clearer that it reproduces pre-existing restore behavior, not a new one |
`docs/en/connectors/source/MySQL-CDC.md:195` | Low (downgraded from @SEZ9's
Medium after independently diffing pre-PR `IncrementalSource.java`) | Yes
(@SEZ9, Issue 4 — partial disagreement, resolved with evidence) |
| 9 | `EventType.SCHEMA_CHANGE_CREATE_TABLE` inserted at ordinal 0; checked
one concrete serialization path (`JobEventReportOperation`) myself and
confirmed it is name-based, not ordinal-based, but did not exhaustively audit
every connector/serializer | `EventType.java:22` | Low | Yes (@SEZ9, Issue 2 —
partial disagreement on severity, upheld from last round) |
| 10 | SPI Javadoc/terminology nits on `SupportMultiTableSinkWriter`
(missing `@param`/`@return`/`@throws`, "newly created" vs "newly added" naming)
| `SupportMultiTableSinkWriter.java:43,53` | Low | Yes (@SEZ9, Issue 8) |
# 5. Merge Recommendation
### Conclusion: Not recommended for merge
1. Blockers
- Root-cause and fix `testMysqlCdcByWildcardsConfigWithNewlyAddedTable`
(Issue 1). This is not a diff-inference this round — I pulled and grepped the
CI log myself and got the same deterministic failure, on every container
variant, with zero occurrences of the registration success log line.
- Close the checkpoint/restore gap on both sides of the pipeline: the
source-side gap (Issue 2, my own carryover) and the sink-side gap (Issue 3,
@SEZ9's finding, which I independently confirmed by reading
`SupportMultiTableSinkWriter.java` myself). These need a state-aware
`restoreSinkWriter`-style counterpart on the sink SPI plus the equivalent
treatment on the source-side enumerator checkpoint, so a runtime-registered
table doesn't just disappear or lose its writer state across a restart.
2. Recommended fixes (non-blocking)
- Issues 4-7, 9, 10 (MultiTableSink serialization/classloader risk, a
registration cap/throttle, JDBC/example doc updates, EventType append-only
style, SPI Javadoc/terminology) are all reasonable follow-ups.
- Issue 8 is now mostly a documentation-wording nit rather than an open
compatibility question, based on the pre-PR diff evidence above.
Overall assessment: this PR has made real, independently-verifiable forward
progress across its many rounds — the capture-pattern filter, the
config-validation guard, and an earlier-round snapshot-read column-indexing fix
are all genuinely correct and valuable on their own. But the feature's own
advertised primary scenario — a table created after the job starts gets picked
up and written to without a restart — still does not work on the current head,
confirmed by CI evidence I pulled and checked myself this round, and the
checkpoint/restore story for a runtime-registered table is incomplete on both
the source and the sink side. @SEZ9's review is the only committer-level
structured review on this PR so far, and my conclusion agrees with theirs that
the blocking items need to be fixed before merge; I'd add that Issue 1 (the E2E
test) should be treated as equally blocking alongside the restore-state gaps,
since right now the feature cannot be demonstrated to work at all in its
own test. As the author, this is posted as a plain comment rather than a
formal review since GitHub does not allow self-approval; it carries no approval
weight.
--
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]