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]

Reply via email to