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]

Reply via email to