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]