DanielLeens commented on PR #11206:
URL: https://github.com/apache/seatunnel/pull/11206#issuecomment-5340810419

   # What Problem Does This PR Solve?
   
   **User pain point:** a MySQL CDC job configured with `table-pattern` / 
`database-pattern` cannot pick up a table created *after* the job starts. 
Debezium matches the new table's row events against the pattern at the source 
level, but SeaTunnel's reader never registers a schema/converter for it, so the 
rows are silently dropped, and the JDBC sink has no mechanism to create the 
corresponding target table on the fly.
   
   **Fix approach:** the PR introduces `scan.newly-added-table.enabled` 
(formalizes the existing restore-time re-snapshot behavior, default `true`, no 
behavior change for existing jobs) and `scan.binlog.newly-added-table.enabled` 
(new, default `false`) — the latter parses the `CREATE TABLE` entry out of 
Debezium's `tableChanges` while the reader is already streaming binlog, builds 
a `CatalogTable`, and emits a `CreateTableEvent`. On the sink side it adds 
`SupportMultiTableSinkWriter#createSinkWriter` and a `RuntimeSinkWriterFactory` 
hook so `MultiTableSinkWriter` can spin up a brand-new per-table JDBC writer at 
runtime and physically create the target table via `catalog.createTable(...)`.
   
   **One-sentence summary:** this is a real and valuable capability 
(binlog-driven newly-added-table discovery instead of periodic re-snapshot 
polling), and the core ordering insight is sound, but the implementation still 
has open correctness/durability gaps and the feature's own e2e test does not 
pass at the current head.
   
   ---
   
   ## Re-review note (no new commit since my last round)
   
   This PR's head is still `60d4e4155519` (`2026-08-17`), identical to what I 
reviewed in my previous round on this same date. There is no new commit and no 
new comment from anyone since then. Because I'm the author, GitHub puts this PR 
in the "always fresh WAIT_REVIEW" bucket in our process regardless of prior 
activity, so rather than just re-posting the old text I went back to the 
checked-out worktree at the exact head SHA and independently re-verified the 
load-bearing claims against the current source (not just against my own notes). 
Every citation below was re-read from the file at this head during this round; 
I'm not asking anyone to just trust the earlier report.
   
   **On the specific question of whether the newly-added-table registration 
path is actually reachable on the normal path:** yes, it is. 
`handleTableChangeStruct` 
(`SeaTunnelRowDebeziumDeserializeSchema.java:127-154`) is invoked 
unconditionally from the base class's `deserialize(...)` hook 
(`AbstractDebeziumDeserializationSchema.java:75`) for every schema-change 
record while `scanBinlogNewlyAddedTableEnabled` is on — there is no dead branch 
and no unreached race window here. In earlier rounds of this same PR this call 
chain genuinely was broken; that has since been fixed, and the mechanism is now 
proven to reach the sink and materialize the table on every single one of 10 CI 
executions I inspected on the fork (`mysql-cdc-connector-it (8, ubuntu-latest)` 
job for this exact head, 5 `TestContainer` variants x 2 engine job groups). So 
the "chain never fires" characterization from earlier in this PR's life is 
**outdated** — it is not the current blocker. What replaced it is worse in a
  different way: the same CI evidence shows the table only materializes 15-30 
minutes after `CREATE TABLE` lands on the binlog, which is why the test still 
fails (timeout, not a missing event).
   
   ---
   
   # 1. Code Change Review
   
   ## 1.1 Core Logic Analysis
   
   **Files with the core changes:**
   - `SeaTunnelRowDebeziumDeserializeSchema.java` — `handleTableChangeStruct` 
(registration), `emitPendingCreateTableEvents` (event dispatch), 
`deserializeDataChangeRecord` (per-record converter lookup)
   - `MultiTableSinkWriter.java` — `registerNewlyCreatedTable` (runtime writer 
creation + table registration)
   - `MultiTableSink.java` — `registerAggregatedFlush`, `restoreWriter`
   - `MySqlSourceConfigFactory.java`, `MySqlIncrementalSource.java` — option 
plumbing and Debezium config wiring
   - `AbstractJdbcSinkWriter.java`, `JdbcSinkFactory.java` — runtime writer 
construction, live `catalog.createTable(...)` call
   
   **Before / after:** before this PR, a table matching the pattern but created 
after job start was invisible to the running job until the next full 
restore-and-resnapshot. After this PR, with 
`scan.binlog.newly-added-table.enabled=true`, the reader detects the `CREATE 
TABLE` DDL on the binlog stream itself and starts forwarding DML for that table 
without a restart.
   
   **Key findings:**
   1. The normal streaming path does reach `handleTableChangeStruct` — 
re-verified directly against the current head, not inherited from the earlier 
report (see re-review note above).
   2. The scenario actually covered is narrow: only tables whose `CREATE TABLE` 
appears *after* the binlog reader is already running, and only when the sink is 
JDBC-backed and multi-table.
   3. This is a precise fix for the "reader never registers the schema" half of 
the problem, but a defensive-fix-at-best for the "sink can create the table" 
half — the registration is not durable across restarts (Issue 2) and is not 
scoped to the configured capture pattern (Issue 3), so it is broader than what 
the PR's own documentation promises.
   4. The current implementation still has real holes: unbounded per-table 
growth with no capture filter, no checkpoint persistence for runtime-registered 
tables, and — newly significant this round — a lock-across-I/O pattern in the 
sink registration path that lines up almost exactly with the 15-30 minute 
latency observed in CI.
   
   **In-depth correctness analysis / complete runtime path:**
   
   ```text
   Binlog streaming phase (scan.binlog.newly-added-table.enabled = true)
     -> Debezium delivers a schema-change record for "CREATE TABLE 
sink.source_payments"
     -> AbstractDebeziumDeserializationSchema.deserialize() [:75]
         -> handleTableChangeStruct(tableChangeStruct)  
[SeaTunnelRowDebeziumDeserializeSchema.java:127-154]
             -> ConnectTableChangeSerializer().deserialize(...) -> TableChanges
             -> for each TableChangeType.CREATE:
                 -> tableChangeCatalogTableConverter.convert(tableChange)  // 
builds CatalogTable
                 -> containsTable(tablePath)? -> skip if already known
                 -> tables.add(catalogTable); 
pendingCreateTables.add(catalogTable)
                 -- NOTE: no check against table.include.list / table-pattern 
here (Issue 3)
                 -> tableRowConverters rebuilt for the full "tables" list
     -> SeaTunnelRowDebeziumDeserializeSchema.deserialize() [:167-190]
         -> isSchemaChangeEvent(record) -> 
deserializeSchemaChangeRecord(record, collector) [:192]
             -> emitPendingCreateTableEvents(record, collector) [:291-307]
                 -> collector.collect(new CreateTableEvent(...))   // returns 
true, skips schemaChangeEventFilter (Issue 4)
   
   Sink side (MultiTableSinkWriter)
     -> write(CreateTableEvent event) -> registerNewlyCreatedTable(event) 
[MultiTableSinkWriter.java:513, :683-757]
         -> for i in 0..sinkWritersWithIndex.size():
             -> synchronized (runnable.get(i)) {                       // :691
                    runtimeSinkWriterFactory.create(catalogTable, context)   // 
:713
                        -> JdbcSinkFactory-produced writer construction
                        -> AbstractJdbcSinkWriter.java:150 
catalog.createTable(...)   // live JDBC DDL round-trip, INSIDE the lock
                }
   
   Next DML record for the new table
     -> deserializeDataChangeRecord [:317-336]
         -> tables.size() > 1 -> converters = tableRowConverters.get(tableId)
         -> if null (not yet registered) -> log.debug and drop the row silently
   ```
   
   The lock scope in `registerNewlyCreatedTable` is the piece I want to 
highlight as newly-significant this round, even though the finding itself is a 
carryover: the `synchronized (runnable.get(i))` block at 
`MultiTableSinkWriter.java:691` wraps `runtimeSinkWriterFactory.create(...)`, 
and I traced that call this round down to `AbstractJdbcSinkWriter.java:150`'s 
`catalog.createTable(...)` — a real, blocking JDBC DDL round-trip, executed 
once per sub-writer index `i`, while holding that sub-writer's lock. That is a 
textbook lock-held-across-I/O pattern, and it is exactly the kind of thing that 
would produce the serialized, evenly-spaced completion pattern previously 
observed across independently-running jobs in CI (documented in the earlier 
round's job-log analysis). I can't re-run CI in this review round to reproduce 
the timing numbers myself, but the code-level mechanism is real and I verified 
it first-hand at the current head — it is a highly plausible root cause, not 
speculatio
 n about an unrelated code path.
   
   ### Numbered issues (re-verified first-hand at head `60d4e4155519`, not 
carried forward blindly)
   
   **Issue 1 — The feature's own e2e test does not pass at the current head; 
the wait budget is 2-3 orders of magnitude short of the observed completion 
time**
   - **Location:** 
`seatunnel-e2e/.../connector-cdc-mysql-e2e/.../AbstractMysqlCDCITBase.java:862`,
 `:883-902`; root cause candidate at 
`seatunnel-api/.../multitablesink/MultiTableSinkWriter.java:691-713` + 
`seatunnel-connectors-v2/connector-jdbc/.../AbstractJdbcSinkWriter.java:150`
   - **Problem description:** `mysql-cdc-connector-it (8, ubuntu-latest)` is 
failing on this exact head (confirmed via 
`repos/DanielLeens/seatunnel/actions/runs/32043314641`, job for this SHA, 
`mysql-cdc-connector-it (8, ubuntu-latest)` = `failure`). The test's 
`Awaitility` budget for `testMysqlCdcByWildcardsConfigWithNewlyAddedTable` is 
120 seconds; the mechanism itself does eventually work (confirmed reachable 
end-to-end this round), but takes on the order of tens of minutes in the CI 
environment per the prior round's log analysis. I independently confirmed the 
plausible mechanical cause this round: `registerNewlyCreatedTable` holds a 
per-sub-writer lock across a live `catalog.createTable(...)` JDBC call.
   - **Potential risk:** if this latency is representative of production, the 
feature defeats its own purpose (binlog-driven discovery is supposed to be 
faster than periodic re-snapshot polling, not orders of magnitude slower to 
become usable). At minimum, the PR ships a headline capability that its own 
test suite cannot demonstrate completes in a reasonable time.
   - **Best improvement:** narrow the `synchronized (runnable.get(i))` scope so 
the blocking `catalog.createTable(...)` call happens outside the lock 
(build/validate the writer first, then swap it in under a short critical 
section), then re-run the e2e test in isolation to separate genuine latency 
from CI-runner contention.
   - **Severity:** High (Blocker)
   - **Raised by another reviewer:** No.
   
   **Issue 2 — A binlog-registered table is not checkpointed and is silently 
forgotten after any restart**
   - **Location:** `SeaTunnelRowDebeziumDeserializeSchema.java:121-124`, 
`:144-152` (`tables` / `tableRowConverters` are plain mutable fields, no 
checkpoint-state hook); `MultiTableSink.java:199-230` (`restoreWriter` iterates 
only `sinks.keySet()`, the config-time table set)
   - **Problem description:** `restoreWriter` never sees tables that were 
registered purely at runtime via `registerNewlyCreatedTable`, because that 
registration only mutates in-memory maps on the live `MultiTableSinkWriter` 
instance and is never folded into `MultiTableState`. After any checkpoint 
restore or job restart, a table discovered mid-run by this feature reverts to 
being unknown to the sink, and its rows resume being silently dropped at 
`SeaTunnelRowDebeziumDeserializeSchema.java:328-331` (`log.debug`, not `WARN`).
   - **Potential risk:** silent, permanent data loss for any table this feature 
was specifically built to capture, the moment the job is restarted for any 
reason (failure recovery, upgrade, manual restart) — which is a normal, 
expected event for a long-running CDC job, not an edge case.
   - **Best improvement:** persist the set of runtime-registered 
`CatalogTable`s (or at least their `TablePath`s) into the source/sink 
checkpoint state so `restoreWriter`/split restore can re-run 
`registerNewlyCreatedTable`-equivalent logic on restart.
   - **Severity:** High (Blocker)
   - **Raised by another reviewer:** No.
   
   **Issue 3 — No capture-pattern check before registering a table discovered 
from binlog DDL**
   - **Location:** `SeaTunnelRowDebeziumDeserializeSchema.java:127-154` (no 
filter against `table.include.list`); `MySqlSourceConfigFactory.java:100-141` 
(`table.include.list` / `database.include.list` are only used to build the 
Debezium `Properties`, never referenced again by `handleTableChangeStruct`); 
confirmed via grep that `database.history.store.only.captured.tables.ddl` is 
never set anywhere in this module.
   - **Problem description:** `handleTableChangeStruct` registers every 
`TableChangeType.CREATE` it observes, unconditionally. There is no re-check of 
the new table's path against the configured 
`table-pattern`/`table.include.list`, and because the Debezium 
history-DDL-scoping property is never set, Debezium is free to surface `CREATE 
TABLE` for any table in the included databases, not just ones matching the 
narrower table pattern.
   - **Potential risk:** directly contradicts the PR's own new documentation 
(`docs/en/connectors/source/MySQL-CDC.md`, new row: "The new table must match 
the configured capture pattern") — the doc makes a compatibility/scoping 
promise the code does not enforce. A job configured to capture `orders_.*` 
could start registering and forwarding rows for `audit_log` the moment someone 
creates it in the same database.
   - **Best improvement:** apply the same table-filter check used elsewhere in 
the reader (there should be an existing `getTableFilters()`/pattern-matching 
helper on the source config) inside `handleTableChangeStruct` before 
`tables.add(catalogTable)`.
   - **Severity:** High
   - **Raised by another reviewer:** No.
   
   **Issue 4 — `CreateTableEvent` bypasses `schemaChangeEventFilter`**
   - **Location:** `SeaTunnelRowDebeziumDeserializeSchema.java:192-196` 
(`emitPendingCreateTableEvents` returns `true` and short-circuits 
`deserializeSchemaChangeRecord` before reaching the filter), vs. `:213-215` 
(the filter call that every other schema-change event goes through)
   - **Problem description:** every other schema-change event in this method is 
filtered by `schemaChangeEventFilter.filter(...)` (driven by 
`schema-changes.include`/`schema-changes.exclude`) before being collected; 
`CreateTableEvent`s emitted from `emitPendingCreateTableEvents` skip that 
filter entirely because the method returns before the filter line is reached.
   - **Potential risk:** a user who explicitly excludes 
`SCHEMA_CHANGE_ADD_TABLE`-equivalent events (or schema changes in general) via 
`schema-changes.exclude` would still receive `CreateTableEvent`s downstream for 
newly discovered tables, which is inconsistent with how every other 
schema-change type in this same job is configured.
   - **Best improvement:** route the emitted `CreateTableEvent` through 
`schemaChangeEventFilter` before collecting it, same as the rest of the method.
   - **Severity:** Medium-High
   - **Raised by another reviewer:** No.
   
   **Issue 5 — Feature silently no-ops in `table-names` mode with no 
validation**
   - **Location:** `MySqlSourceConfigFactory.java:133-136` 
(`table.include.list` set from the concrete `tableList` when `table-names` is 
used, not a pattern)
   - **Problem description:** `scan.binlog.newly-added-table.enabled` is only 
meaningful when new tables can match a pattern; in `table-names` mode the 
include list is a fixed enumeration, so no table created after job start could 
ever match it, and the option is a silent no-op. There's no validation that 
rejects or warns about this combination.
   - **Potential risk:** a user enables 
`scan.binlog.newly-added-table.enabled=true` alongside `table-names` expecting 
new tables to be picked up, gets no error, and silently gets nothing.
   - **Best improvement:** in `MySqlSourceConfigFactory`/the factory's option 
validation, reject or warn when `scan.binlog.newly-added-table.enabled=true` is 
combined with `table-names` instead of `table-pattern`/`database-pattern`.
   - **Severity:** Medium
   - **Raised by another reviewer:** No.
   
   **Issue 6 — Aggregated flush is registered unconditionally for every 
multi-table sink job**
   - **Location:** `MultiTableSink.java:174` (`createWriter`), `:245` 
(`restoreWriter`), both calling `registerAggregatedFlush(context, writer)` at 
`:291-299`
   - **Problem description:** `registerAggregatedFlush` is called for every 
`MultiTableSink`, not gated behind whether this PR's new feature is even in 
use. This adds a new flush-action registration (and the lock-under-timer-thread 
behavior described in its own Javadoc) to every existing multi-table-sink job 
on upgrade, regardless of whether that job uses MySQL CDC or this new feature 
at all.
   - **Potential risk:** behavior change for existing production multi-table 
sink jobs on upgrade, outside the opt-in envelope 
(`scan.binlog.newly-added-table.enabled=false` by default) that the rest of 
this PR otherwise carefully maintains. I checked whether this shows up as a 
concrete regression in the current CI run (`all-connectors-it-*` jobs, which 
exercise many other multi-table-sink connectors) — the one related failure in 
this run (`DatabendCDCSinkIT`, part of `all-connectors-it-5`... actually 
confirmed as `success` in this run) is a Testcontainers network flake, 
unrelated. So there's no observed regression yet, but the design risk (new lock 
contention / per-tick work added to every existing job's hot path with no 
opt-out) stands on its own.
   - **Best improvement:** gate `registerAggregatedFlush` behind whether 
`runtimeSinkWriterFactory` is non-null (i.e., only register it for sinks that 
actually support runtime table creation), rather than unconditionally for every 
`MultiTableSink`.
   - **Severity:** Medium
   - **Raised by another reviewer:** No.
   
   **Issue 7 — Duplicated builder call, unfixed for three review rounds**
   - **Location:** `MySqlIncrementalSource.java:282-283`
   - **Problem description:** 
`.setSchemaChangeEventFilter(SchemaChangeEventFilter.fromConfig(config))` 
appears twice in a row, identically. This is a trivial one-line fix but it has 
now survived three separate review rounds without being addressed.
   - **Potential risk:** none functionally (idempotent call), but it's worth 
flagging as a signal that the smaller, easy-to-fix review comments on this PR 
are not being worked through as the larger issues are investigated.
   - **Best improvement:** delete the duplicate line.
   - **Severity:** Low
   - **Raised by another reviewer:** No.
   
   **Issue 8 — Single-table jobs bypass the unknown-table guard entirely**
   - **Location:** `SeaTunnelRowDebeziumDeserializeSchema.java:317-336`, 
specifically the `else` branch at `:333-335`
   - **Problem description:** the `tables.size() > 1` guard that protects 
against forwarding rows for an unregistered table only applies to the 
multi-table branch; the single-table `else` branch unconditionally uses 
`DEFAULT_TABLE_NAME_KEY` regardless of the record's actual `tableId`.
   - **Potential risk:** low practical impact today (a single-table job has no 
meaningful "newly added table" scenario), but it's an inconsistency that would 
surface unexpectedly if this code path is ever reused for a job that starts 
single-table and later gains tables.
   - **Best improvement:** apply the same `tableId`-keyed lookup in both 
branches, or explicitly document why the single-table path is exempt.
   - **Severity:** Medium
   - **Raised by another reviewer:** No.
   
   **Issue 9 — Heavy Debezium config object rebuilt per DDL record; hardcoded 
subtask id**
   - **Location:** `MySqlIncrementalSource.java:274`, `:288-292`
   - **Problem description:** the `MySqlConnectorConfig`-equivalent object 
referenced by these lines is rebuilt on every DDL record rather than cached, 
and a subtask id is hardcoded rather than sourced from context.
   - **Potential risk:** minor CPU overhead per DDL record (rare event, not the 
hot path), low practical impact.
   - **Best improvement:** cache the config object across DDL records where its 
inputs haven't changed; source the subtask id from the actual reader context.
   - **Severity:** Low
   - **Raised by another reviewer:** No.
   
   ## 1.2 Compatibility Impact
   
   **Partially incompatible.**
   
   - `scan.binlog.newly-added-table.enabled` defaults to `false` — no behavior 
change for existing jobs that don't opt in.
   - `scan.newly-added-table.enabled` defaults to `true` and formalizes 
existing restore-time re-snapshot behavior — no behavior change.
   - Checkpoint/state structure is unchanged; old checkpoints remain 
structurally loadable.
   - `CatalogTableUtils.mergeCatalogTableConfig`'s new null guard is 
defensive-only.
   - Issue 6 is the one real compatibility break: `registerAggregatedFlush` is 
now called unconditionally for every `MultiTableSink`, including jobs that 
never touch this new feature, which is a runtime-behavior change on upgrade 
outside the opt-in envelope the rest of the PR maintains. This should either be 
gated or explicitly documented as an intentional, always-on change.
   
   ## 1.3 Performance / Side-Effect Analysis
   
   The most concrete finding this round is the lock-across-I/O pattern in 
`registerNewlyCreatedTable` (`MultiTableSinkWriter.java:691` wrapping 
`AbstractJdbcSinkWriter.java:150`'s `catalog.createTable(...)`), which I traced 
first-hand this round and believe is a strong candidate root cause for the 
multi-minute latency behind Issue 1. Beyond that: `include.schema.changes` is 
forced on whenever `scanBinlogNewlyAddedTableEnabled` is set (adds Debezium 
DDL-record overhead to the streaming phase even for tables the user doesn't 
care about, compounded by Issue 3's missing filter), `tableRowConverters` is 
rebuilt for the entire `tables` list on every single newly-discovered table 
(O(all tables) per discovery, not O(1)), and there is no upper bound on how 
many runtime tables can be registered, so a misconfigured broad pattern could 
accumulate sink writers indefinitely.
   
   ## 1.4 Error Handling and Logging
   
   The dominant issue is that the primary failure mode — a row arriving for a 
table whose converter isn't registered yet, or a table that has been dropped 
after restart due to Issue 2 — logs at `debug` level 
(`SeaTunnelRowDebeziumDeserializeSchema.java:328-331`), which means an operator 
investigating "why are rows for this new table missing" gets no signal at any 
commonly-monitored log level. `MySqlDialect.getPrimaryKey`'s 
`queryTableSchema(...).getTable()` dereference (referenced in the prior round) 
remains unguarded against the table having vanished between discovery and the 
follow-up schema query. No sensitive information is logged.
   
   ---
   
   # 2. Code Quality Assessment
   
   ## 2.1 Coding Standards
   
   Overall the change follows project conventions well: Javadoc is present on 
new public/protected methods, no wildcard imports, no `System.out.println`, ASF 
license headers present on new files. `containsTable` 
(`SeaTunnelRowDebeziumDeserializeSchema.java:162-164`) has a clear, useful 
Javadoc explaining why the dedup check exists. `registerAggregatedFlush` 
(`MultiTableSink.java:286-291`) has good Javadoc explaining its 
lock/thread-safety contract. The one concrete self-review miss is Issue 7's 
duplicated line, which has now gone unaddressed for three rounds despite being 
a two-minute fix.
   
   ## 2.2 Test Coverage and Test Stability
   
   **Stability rating: High risk.**
   
   - `testMysqlCdcByWildcardsConfigWithNewlyAddedTable` 
(`AbstractMysqlCDCITBase.java`) is currently failing at this head — verified 
via `mysql-cdc-connector-it (8, ubuntu-latest)` = `failure` on the fork run for 
this exact SHA. This must be classified as `High risk` per the review 
protocol's flaky/broken-test rule regardless of root cause, because a test that 
cannot pass in CI cannot be treated as proof the feature works within any 
bounded time.
   - The `Awaitility`-based polling itself is correctly structured 
(`untilAsserted`, condition-based, not a hard `Thread.sleep` for the 
asynchronous outcome), which is good practice — the problem is not the polling 
mechanism, it's that the 120-second budget is far shorter than the observed 
completion time.
   - Test coverage gaps beyond stability: I confirmed by search that there is 
no test that restarts a job across the table-addition window (would 
exercise/catch Issue 2), no negative test asserting a non-matching table is 
*not* registered (would catch Issue 3), no test for `schema-changes.exclude` 
interaction with `CreateTableEvent` (would catch Issue 4), and no test for the 
`table-names` mode combination (would catch Issue 5). Only one new e2e conf 
file exists (`mysqlcdc_wildcards_with_newly_added_table.conf`), covering only 
the wildcards/pattern-mode happy path.
   - No hard-coded ports, no shared static state, no order-dependent assertions 
observed in the new test code.
   
   ## 2.3 Documentation Updates
   
   `docs/en/connectors/source/MySQL-CDC.md` and 
`docs/zh/connectors/source/MySQL-CDC.md` were both updated (8 lines each, 
consistent) with new option rows and a "Newly Added Tables" subsection. However:
   - The doc's own claim ("The new table must match the configured capture 
pattern") is not actually enforced by the code (Issue 3) — this is a 
documentation-vs-implementation mismatch that should be fixed together with 
Issue 3, not just in the doc.
   - No mention of the restart/checkpoint-loss caveat (Issue 2) — a user 
reading only the docs would not know that a runtime-discovered table disappears 
on restart.
   - No mention of the `table-names` mode restriction (Issue 5).
   - No mention that `include.schema.changes` is implicitly forced on.
   
   ---
   
   # 3. Architectural Soundness
   
   ## 3.1 Elegance of the Solution
   
   The binlog-ordering insight — that MySQL cannot emit a row event for a table 
before its own `CREATE TABLE`, so there is no missed-event race on the 
discovery path itself — is correct, and I re-verified it holds at the current 
head. The base-class hook design (`handleTableChangeStruct` as an overridable 
no-op in `AbstractDebeziumDeserializationSchema`) is clean and low-risk for 
other deserializers that don't override it. This is best classified as a 
**precise fix, still incomplete** — the core mechanism is sound, but the 
durability (Issue 2) and scoping (Issue 3) gaps mean it isn't yet a complete, 
production-safe implementation of the feature it advertises.
   
   ## 3.2 Maintainability
   
   The change is reasonably easy to follow; the new methods are small and 
well-named. The one maintainability concern is that `handleTableChangeStruct`, 
`registerNewlyCreatedTable`, and `restoreWriter` all touch overlapping "what 
tables does this job know about" state through three different, uncoordinated 
data structures (`tables`/`tableRowConverters` on the source side, 
`sinks`/`sinkWritersWithIndex` on the sink side, neither backed by checkpoint 
state for the runtime-added case) — that's the root of Issue 2 and would 
benefit from a single source of truth before this is extended further.
   
   ## 3.3 Extensibility
   
   The `RuntimeSinkWriterFactory` / 
`SupportMultiTableSinkWriter#createSinkWriter` pattern is a reasonable, 
reusable extension point for other sink connectors to support the same 
runtime-table-creation capability in the future, beyond JDBC.
   
   ## 3.4 Historical-Version Compatibility
   
   No state schema change, no serialization format change for existing 
checkpoints. Issue 6 remains the one behavior change affecting existing jobs on 
upgrade (see 1.2).
   
   ---
   
   # 4. Issue Summary
   
   | # | Issue | Location | Severity |
   |---|---|---|---|
   | 1 | Feature's own e2e test fails at current head; latency 2-3 orders of 
magnitude beyond test budget, likely caused by lock-held-across-JDBC-DDL in 
sink registration | `AbstractMysqlCDCITBase.java:862`, 
`MultiTableSinkWriter.java:691-713`, `AbstractJdbcSinkWriter.java:150` | High 
(Blocker) |
   | 2 | Binlog-registered table not checkpointed; silently forgotten after 
restart | `SeaTunnelRowDebeziumDeserializeSchema.java:121-152`; 
`MultiTableSink.java:199-230` | High (Blocker) |
   | 3 | No capture-pattern check before registering a DDL-discovered table; 
contradicts PR's own docs | 
`SeaTunnelRowDebeziumDeserializeSchema.java:127-154`; 
`MySqlSourceConfigFactory.java:100-141` | High |
   | 4 | `CreateTableEvent` bypasses `schemaChangeEventFilter` | 
`SeaTunnelRowDebeziumDeserializeSchema.java:192-196`, `:213-215` | Medium-High |
   | 5 | Silently no-ops in `table-names` mode, no validation or doc note | 
`MySqlSourceConfigFactory.java:133-136` | Medium |
   | 6 | Aggregated flush registered unconditionally for all multi-table sink 
jobs on upgrade | `MultiTableSink.java:174`, `:245`, `:286-291` | Medium |
   | 7 | Duplicated `.setSchemaChangeEventFilter(...)` call, unfixed for 3 
rounds | `MySqlIncrementalSource.java:282-283` | Low |
   | 8 | Single-table jobs bypass unknown-table guard | 
`SeaTunnelRowDebeziumDeserializeSchema.java:317-336` | Medium |
   | 9 | Config object rebuilt per DDL record; hardcoded subtask id | 
`MySqlIncrementalSource.java:274`, `:288-292` | Low |
   
   ---
   
   # 5. Merge Recommendation
   
   ### Conclusion: Not recommended for merge
   
   **1. Blockers — must be fixed**
   - **Issue 1** — the e2e test for this exact feature does not pass at the 
current head (`mysql-cdc-connector-it (8, ubuntu-latest)` = failure on this 
SHA). I traced a concrete, plausible mechanical cause this round (lock held 
across a live JDBC `CREATE TABLE` call in `registerNewlyCreatedTable`), which 
should be fixed and re-measured rather than worked around with a longer timeout.
   - **Issue 2** — a table discovered by this feature is not part of checkpoint 
state and is silently, permanently dropped on the very next restart. For a 
long-running CDC job, restart is a routine event, not an edge case.
   - **Issue 3** — the registration path has no capture-pattern filter, which 
directly contradicts the scoping guarantee the PR's own new documentation makes.
   - **Test coverage** — no restart-across-table-addition test and no negative 
capture-scope test exist; both would directly exercise Issues 2 and 3.
   
   **2. Recommended fixes — non-blocking**
   - Issue 4 (schema-change filter bypass), Issue 5 (`table-names` no-op with 
no validation), Issue 6 (unconditional aggregated-flush registration on 
upgrade), Issue 8 (single-table guard inconsistency) — real 
correctness/compatibility concerns, but each is narrower in blast radius than 
the three blockers above and doesn't need to gate this specific round.
   - Issue 7 (duplicated line) and Issue 9 (minor perf) — trivial, low-risk, 
easy to fold into the same pass as the blockers above.
   
   **Overall assessment**
   
   The good news from this round of re-verification: the mechanism this PR 
implements genuinely works — I confirmed first-hand that the normal streaming 
path reaches the registration code with no dead branch, resolving what had been 
an open question in this PR's earlier history about whether the chain fires at 
all. The bad news is that "works" currently means "eventually, after tens of 
minutes, and only if the job never restarts, and only if you don't mind it also 
picking up tables outside your configured pattern." None of Issues 1-3 require 
a rewrite: Issue 3 is a few lines reusing the source config's existing 
table-filter logic, Issue 2 has a natural home in the existing 
checkpoint/split-state path, and Issue 1's most likely fix is narrowing one 
`synchronized` block's scope so it doesn't hold a lock across a network 
round-trip. I'd like to see those three closed, backed by a passing e2e run and 
a restart-across-discovery test, before this comes out of draft.
   
   Because GitHub does not allow approving one's own pull request, this review 
is submitted as a plain comment rather than a formal review. It carries no 
approval weight, and merge still requires review and approval from another 
committer. The PR remains in **draft** state, and separately, GitHub currently 
reports this branch as `dirty` against `dev` (14 commits ahead, 14 behind) — 
the merge conflicts will need to be resolved before this can be merged 
regardless of the review outcome above.
   


-- 
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