DanielLeens commented on PR #11206:
URL: https://github.com/apache/seatunnel/pull/11206#issuecomment-5473713721
Update on this exact head (`9577677bb86c`, `2026-08-30`), two new commits
since my last comment (`5459248432`, `2026-08-29T00:41:51Z`). Both are real
code changes, so I re-traced them fully rather than trusting the commit
messages.
# What Problem Does This PR Solve?
A MySQL CDC job using `database-pattern`/`table-pattern` (wildcards)
previously could not pick up a table created after the job started. This PR
adds `scan.binlog.newly-added-table.enabled` to convert a Debezium `CREATE
TABLE` schema record into a `CreateTableEvent` in-flight, and gives
`MultiTableSinkWriter`/`MultiTableSink` a way to create a per-table JDBC writer
(and physical table) at runtime.
# 1. Code Change Review
## 1.1 Core Logic Analysis — this is the important part, so I'm being direct
about it
**The feature's own E2E test now passes, and I traced exactly why.**
`mysql-cdc-connector-it` is green on both JDK 8 (`99275592981`) and JDK 11
(`99275593075`) in this head's fork run (`33316056455`) — the first time in
this PR's history that's been true. I did not take that at face value; I read
the actual fix.
Commit `9577677bb86c` ("[Fix][Zeta] Forward runtime create table events")
touches `SeaTunnelSourceCollector.java:279-297` in `seatunnel-engine-server`,
and it is the real root cause of every prior failure I reported on this PR
(including my `2026-08-20` finding that `CreateTableEvent` "never appears
anywhere in the job log for this head", and the `2026-08-29` finding of a
sink-side `Table ... doesn't exist` timeout). The chain is:
```
SeaTunnelRowDebeziumDeserializeSchema.handleTableChangeStruct()
(connector-cdc-base, unchanged)
-> registers the new CatalogTable, logs "Registered newly added CDC table"
(line 155)
-> emitPendingCreateTableEvents() calls
collector.collect(createTableEvent) (line 317)
-> AT RUNTIME this collector is SeaTunnelSourceCollector (Zeta engine)
-> collect(SchemaChangeEvent event), rowType is MultipleRowType,
tableId not yet
in rowTypeMap (true for every newly-added table by definition)
BEFORE this commit: log.warn("Ignore schema change event for unknown
table..."), return
-> sendRecordToNext() is NEVER called, event
dies here
AFTER this commit: if (event instanceof CreateTableEvent) seed
rowTypeMap from
createTableEvent.getChangeAfter(), fall through
to sendRecordToNext()
```
So the connector-side logging was correct the whole time ("Registered newly
added CDC table" firing was real progress, as I said on 08-29) — the event was
being built and handed to the engine correctly, then silently dropped one layer
up, inside the Zeta engine's own collector, before it could ever reach
`MultiTableSinkWriter.applySchemaChangeEvent()` or trigger a save-mode table
creation. That fully explains why the sink-side table was never created within
the 2-minute Awaitility window: the sink never received the event in the first
place. This was a real engine-level bug, not a connector-level one, and not
something the connector-side fixes in `82165d5d18a1` could have addressed on
their own.
The new unit test
(`SeaTunnelSourceCollectorSchemaChangeTest.shouldForwardCreateTableEventForUnknownTableInMultipleRowType`)
directly exercises this: `collect(CreateTableEvent)` for an unknown table
followed by `collect(row)` for that same table, asserting `output.received()`
fires twice. I checked it against a real `SeaTunnelSourceCollector` (not a
partial mock of the method under test), and it fails against the pre-fix code
path (the event would never reach `sendRecordToNext`). This is the right
regression test for this bug.
Commit `496a3253694c` ("[Fix][API] Copy primary key columns for
serialization") is a second, independent fix in the same problem space, this
time for the Flink engine: `PrimaryKey`'s constructor previously aliased the
caller-supplied `columnNames` list directly. When an immutable/unmodifiable
list is passed in — which happens along the CDC schema-conversion path —
Flink's Kryo serializer attempts to populate that list in place during
(de)serialization and throws. The fix defensively copies into a new
`ArrayList`. `PrimaryKeyTest.copiesImmutableColumnNames` proves both the
copy-on-construct behavior and that mutations to the primary key's own list
don't leak back into the caller's list. Correctly scoped, and the 2-arg legacy
constructor is preserved by delegation, so no observable behavior change for
existing callers passing mutable lists.
**Both fixes are genuine, targeted, and test-backed. I'm not going to
relitigate my own past over-optimistic and then over-corrected readings on this
PR (08-18, 08-19) — I looked at the actual before/after source and the actual
current CI job logs for this exact head, not the commit message, before writing
this.**
## 1.2 Compatibility Impact
Fully compatible. `SeaTunnelSourceCollector`'s new branch only activates for
`event instanceof CreateTableEvent` with a `tableId` not already in
`rowTypeMap` — a case that, absent this PR's own
`scan.binlog.newly-added-table.enabled` feature, should not occur for any
pre-existing job (any table a pre-existing job knows about is already seeded
into `rowTypeMap` before the first event for it arrives). No existing
schema-change handling path is touched; the `else` branch for known tables and
the fallback `log.warn(...)`+`return` for any other unknown-table event type
are byte-for-byte unchanged. The `PrimaryKey` change is a pure
internal-representation fix (copy vs. alias) with the legacy 2-arg constructor
preserved.
## 1.3 Performance / Side-Effect Analysis
Negligible. One `instanceof` check plus one `HashMap.put` per
newly-discovered table (not per row), and one `ArrayList` allocation per
`PrimaryKey` construction (already happens once per table/schema-change event,
not per row).
## 1.4 Error Handling and Logging
Carrying forward, unresolved, from my `2026-08-28` review (`5050333217`) —
neither of this round's two commits touches either of these:
**Issue 1 (carried over, unresolved, High): uncaught
`IllegalArgumentException` from runtime identifier validation crashes the
source reader instead of skipping the offending table**
- Location: `MySqlCatalogTableUtils.java:109` (`validateIdentifier`),
reached with no surrounding try/catch from `handleTableChangeStruct`
(`SeaTunnelRowDebeziumDeserializeSchema.java:127-153`).
- I re-checked this against the current head: `handleTableChangeStruct`
still calls `tableChangeCatalogTableConverter.convert(tableChange)` directly
inside the `tableChanges.forEach` lambda, no try/catch anywhere in that call
chain.
- Risk unchanged: a source-DB principal with only `CREATE TABLE` privilege
on a table matching the capture pattern can crash the whole running job by
creating one table/column whose name fails identifier validation — under this
feature's own enabled path.
- Severity: High.
**Issue 2 (carried over, unresolved, High): `hasSourceMatchedWriter` accepts
every `CreateTableEvent` whenever a runtime factory exists, with no per-sink
capability gate on Zeta**
- Location: `MultiTableSinkWriter.java:506-508` — `if
(runtimeSinkWriterFactory != null) { return true; }` still runs unconditionally
before the `supportsNewlyCreatedTable()` check just below it, and I confirmed
`runtimeSinkWriterFactory` is non-null for every `MultiTableSink` instance.
- On the Zeta engine there is still no upstream
`supports()`/`SchemaChangeType` gate before this, unlike Flink's
`SchemaOperator`. This still contradicts the PR's own documented JDBC-only
scope.
- Severity: High.
Both are exactly as I described them on 08-28/08-29; I'm not adding new
instances of these, just confirming they're still live at the current head.
# 2. Code Quality Assessment
## 2.1 Coding Standards
Both new commits are cleanly scoped (one production file + one test file
each), have doc comments explaining the *why* (not just what) at the exact
lines that need it, and match the surrounding code style.
## 2.2 Test Coverage and Test Stability
-
`SeaTunnelSourceCollectorSchemaChangeTest.shouldForwardCreateTableEventForUnknownTableInMultipleRowType`
and `PrimaryKeyTest.copiesImmutableColumnNames` are both deterministic,
single-threaded, no `Thread.sleep`, no shared static state, no timing
dependency. Stable.
- The E2E level: `mysql-cdc-connector-it` passing 2/2 (JDK 8 + JDK 11) on
this head is real signal, but it's one run. I'm not calling this "proven
stable" off a single green run given this test's specific history of flipping
between failure modes on this PR; I'd want to see it stay green through at
least one more CI cycle (e.g. after the two Issues above are fixed) before
calling the E2E coverage solid.
## 2.3 Documentation Updates
No doc changes in either of this round's two commits, and none were needed —
neither touches user-facing config or contract surface.
# 3. Architectural Soundness
## 3.1 / 3.2 / 3.3
No change from my `08-28` assessment: the design direction (binlog-driven
discovery over restart-based re-snapshot) is sound, and the dead
`SupportMultiTableSinkWriter#createSinkWriter` SPI question (Issue 4, Medium,
carried over unchanged) is still the main maintainability wart.
## 3.4 Historical-Version Compatibility
Both new fields/behaviors are new to this unmerged PR, so there is no
released-version checkpoint/restore compatibility surface being touched by
either of this round's commits.
# 4. Issue Summary
| # | Issue | Location | Severity | Status |
|---|---|---|---|---|
| 1 | Identifier-validation `IllegalArgumentException` crashes job instead
of skipping the table | `MySqlCatalogTableUtils.java:109` | High | Open,
unchanged |
| 2 | `hasSourceMatchedWriter` skips the sink-capability gate on Zeta |
`MultiTableSinkWriter.java:506-508` | High | Open, unchanged |
| 3 | `SeaTunnelSourceCollector` silently dropped `CreateTableEvent` for
unknown tables | `SeaTunnelSourceCollector.java:279-297` | High | **Fixed this
round, verified** |
| 4 | `SupportMultiTableSinkWriter#createSinkWriter` SPI still dead code,
not deprecated | `SupportMultiTableSinkWriter.java` | Medium | Open, unchanged |
| 5 | No unit test exercises `MultiTableSink`'s real `FactoryUtil`-based
runtime path | `MultiTableSinkWriterTest.java` | Medium | Open, unchanged |
| 6 | Docs describe a JDBC-only scope not enforced in code (same root cause
as #2) | `docs/en/connectors/source/MySQL-CDC.md` | Medium | Open, unchanged |
# 5. Merge Recommendation
**Not recommended for merge yet — but this is genuine, verified progress,
not a repeat of my prior status quo.** The headline capability (a newly-created
wildcard table gets discovered from binlog and picked up end-to-end without a
job restart) now actually completes in CI for the first time in this PR's
history, and I traced the fix to a real, previously-misdiagnosed-by-me engine
bug rather than accepting the commit message. That said, two High-severity
issues from my last round (Issues 1 and 2 above) are untouched by this round's
commits and remain the reason I'm not calling this ready: one is an
availability bug reachable by the same low-privileged source actor the
feature's own identifier-validation hardening was meant to constrain, the other
means the runtime-writer path silently applies to sink types the docs say it
shouldn't, on the primary engine.
Also flagging for transparency since I'm the author here and this carries no
approval weight either way: this PR is `diverged` from `dev` (`ahead_by=21`,
`behind_by=89`) and still in draft. CI on this head has exactly one unrelated
failure
(`PostgresCDCIT.testPostgresCdcSnapshotOnlyAndCommittedOffsetStartupModes`, a
`ConditionTimeout` on JDK 8 only, passing on JDK 11 in the same run —
consistent with known pre-existing Postgres CDC E2E flakiness, not this PR's
diff) and one unrelated cancellation (`paimon-connector-it`, a known
pre-existing classloader-hang timeout, also not this PR's diff). Neither blocks
the conclusion above; the real blockers are Issues 1 and 2.
As before, this is a plain comment rather than a formal review since GitHub
does not allow self-approval on this PR; it carries no approval weight, and a
maintainer with write access still needs to do the actual review/merge step
once Issues 1 and 2 are addressed.
--
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]