ryanmeowy commented on issue #10203:
URL: https://github.com/apache/seatunnel/issues/10203#issuecomment-6097595906

   @SEZ9 @DanielLeens — here is the design write-up for the live-runtime 
discovery gap, following your direction. It is a design only: no PR is opened 
from it, and nothing in it changes table admission or the `table_pattern` 
configuration while #12567 is under review. The items you asked for are mapped 
to sections in the header; §8 holds the open questions, and OQ3 (reader-side 
backfill mechanics) is the one I would most like direction on.
   
   ---
   
   # Live table discovery while the CDC stream is running (MySQL / OceanBase)
   
   Design draft for [#10203](https://github.com/apache/seatunnel/issues/10203), 
written for the request in
   [comment 
6096357599](https://github.com/apache/seatunnel/issues/10203#issuecomment-6096357599):
 the
   missing piece is **live discovery while the binlog stream is running**, not 
another gate on the existing
   restore path.
   
   **Status: design only.** No PR is opened here, and nothing below changes 
table admission or the
   `table_pattern` configuration while the pending per-table restore-watermark 
work (#12567 / #12565) is
   under review; implementation starts after that work lands, as asked.
   
   **Where each requested item is answered:** trigger cadence, poll vs. DDL 
event → §3.2 · enumerator
   ownership of newly admitted tables → §3.1, §3.3 · split reassignment → §3.3 
step 4 · what is persisted
   in checkpoint state → §3.5 · checkpoint failure / recovery → §3.6 · 
duplicate-admission suppression →
   §3.7 · exact MySQL vs. OceanBase scope → §3.8 · `table_pattern` backward 
compatibility → §5. Scope is
   kept to what the issue asks for: the sink-schema consequence is documented 
as a limitation (§4.2) and
   raised as a question (OQ2), not designed here.
   
   Everything below is written against `apache/seatunnel` `dev` @ `0278a6a74` 
(the base that also carries
   the pending per-table restore-watermark work, PR #12567).
   
   ---
   
   ## 1. Current behaviour (verified)
   
   | # | Fact | Evidence |
   |---|------|----------|
   | 1 | The captured table set is fixed when the enumerator is created: 
`discoverDataCollections()` runs once in `createEnumerator()` and once in 
`restoreEnumerator()`. | 
`connector-cdc-base/.../source/IncrementalSource.java:303-358`, `:360-425` |
   | 2 | The incremental-phase split hard-codes its table set. 
`createIncrementalSplits()` distributes the tables captured **at that moment** 
into `incrementalParallelism` buckets, and each `IncrementalSplit` carries that 
list. `IncrementalSplitAssigner.snapshotState()` persists only `new 
IncrementalPhaseState(startupOffset)`. | 
`.../source/enumerator/IncrementalSplitAssigner.java` 
(`createIncrementalSplits`, `createIncrementalSplit`, `snapshotState`); 
`.../source/enumerator/state/IncrementalPhaseState.java` |
   | 3 | A reader that is streaming can never take another split: 
`checkSplitOrStartNext()` returns immediately while `currentFetcher instanceof 
IncrementalSourceStreamFetcher`. Splits added to the queue are never consumed. 
| `.../source/reader/IncrementalSourceSplitReader.java:152-156`, `:190-192` |
   | 4 | Change records for a table that is in neither 
`maxSplitHighWatermarkMap` nor `finishedSplitsInfo` are silently dropped 
(`shouldEmit` returns `false`). | 
`.../source/reader/external/IncrementalSourceStreamFetcher.java:265-300`, 
`:301-313` |
   | 5 | With `table_pattern`, the Debezium runtime filter **is** the regex, so 
a table created later that matches still produces change events on the existing 
stream — they are just discarded by (4). With explicit `table_names` nothing 
can ever match. | 
`connector-cdc-mysql/.../config/MySqlSourceConfigFactory.java:104-113` |
   | 6 | `CREATE TABLE` reaches the reader (schema-change records bypass 
`shouldEmit`, `:298`) and the Debezium record carries `tableChanges`, but the 
repository's resolver only treats `ALTER TABLE` as a schema change, so today a 
`CREATE TABLE` is not turned into a `SchemaChangeEvent`. | 
`IncrementalSourceStreamFetcher.java:298`; 
`connector-cdc-base/.../schema/AbstractSchemaChangeResolver.java:56`, `:68-80` |
   | 7 | The enumerator has no timer/loop: the engine calls `run()` exactly 
once when the task goes RUNNING; afterwards only callbacks occur 
(`handleSplitRequest`, `handleSourceEvent`, `notifyCheckpointComplete`, 
`addSplitsBack`). | 
`seatunnel-engine/.../task/SourceSplitEnumeratorTask.java:334`; 
`seatunnel-api/.../source/SourceSplitEnumerator.java:55`; 
`.../enumerator/IncrementalSourceEnumerator.java` |
   | 8 | The snapshot phase already pairs a **lock-free** split snapshot with a 
**bounded per-table binlog backfill over `[low, high]`** — but only when 
`exactly-once` is enabled, and only because no live stream client exists at 
that point (both clients would need the same `database.server.id`). The 
`CHANGE_EVENT_QUEUE` is bounded for a stream split but explicitly unbounded 
(`Integer.MAX_VALUE`) for a snapshot split in exactly-once mode. | 
`connector-cdc-mysql/.../source/reader/fetch/scan/MySqlSnapshotFetchTask.java:57-119`,
 `:121-171`; `.../scan/MySqlSnapshotSplitReadTask.java:127-159`; 
`.../source/reader/fetch/MySqlSourceFetchTaskContext.java:142-161` |
   
   Net effect today: a table created after the job starts gets **no snapshot 
backfill**; whatever stream
   records it produces are dropped; it only becomes visible on restore — and 
the restore path itself is
   exactly the boundary under review in #12567.
   
   ## 2. Goal / non-goals
   
   **Goal.** When a table that matches the configured filter appears while the 
job is running:
   1. its pre-existing rows are snapshotted exactly once, with a consistent 
binlog watermark;
   2. its changes from that watermark on are streamed by the same reader, in 
order;
   3. the admission is durable across checkpoints and safe on restore/failure;
   4. no behaviour change for jobs that do not opt in, and no change to 
`table_pattern` restore semantics.
   
   **Non-goals (explicitly out of scope).**
   * Table **deletion** while streaming — stays on the existing restore-time 
diff (`IncrementalSource.restore()`).
   * Any change to how the *startup/restore* table set is computed 
(`createEnumerator()`,
     `restoreEnumerator()`, `restore()`), and any change to `table_names` mode.
   * Schema evolution for the discovered table (source `ALTER`s, sink catalog 
updates). The one consequence
     that matters — a runtime-discovered table has no `CatalogTable` from 
submission time — is stated as a
     limitation in §4.2 and raised as OQ2; it is **not** designed here.
   * No OptionRule **DSL** change and no unrelated `MySqlIncrementalSource` 
work. The two options in §7 are
     the only configuration surface this design adds, and they are the 
enable/disable knob the issue asked
     for, extended with the poll-vs-DDL choice this review asked to settle.
   * Non-MySQL dialects at runtime. The base changes are dialect-agnostic, but 
only MySQL/OceanBase opt in.
   
   ## 3. Design
   
   **Single admission authority in the enumerator; two trigger channels; the 
backfill is executed by the
   reader that already owns the table's bucket.**
   
   ### 3.1 Components
   
   1. **`TableDiscoveryEvent`** (new `SourceEvent`, base module)
      `TableId` + the offset of the DDL record. Reader → enumerator via
      `SourceReader.Context.sendSourceEventToEnumerator()` 
(`seatunnel-api/.../source/SourceReader.java:111`),
      handled in `IncrementalSourceEnumerator.handleSourceEvent()`
      (`.../enumerator/IncrementalSourceEnumerator.java:103`).
   2. **`TableAdmission`** (enumerator-side state machine) — owns eligibility, 
dedup, bucket-owner
      selection, snapshot-split creation and the persisted admission registry.
   3. **`TableAdmissionEvent`** (new `SourceEvent`) — enumerator → owner reader 
via
      `SourceSplitEnumerator.Context.sendEventToSourceReader()`
      (`seatunnel-api/.../source/SourceSplitEnumerator.java:139`), carrying the 
`SnapshotSplit` and the
      snapshot-split id it must acknowledge.
   4. **Reader admission sub-phase** — in `IncrementalSourceSplitReader`, the 
ability to interleave a
      snapshot read with a live stream fetcher (this is the core code change, 
see 3.4).
   5. **`IncrementalSourceStreamFetcher.addAdmittedTable(TableId, Offset 
highWatermark)`** — extends the
      running filter so post-watermark records of the new table are emitted.
   
   ### 3.2 Trigger cadence (question 1)
   
   Two channels, both feeding the same admission path:
   
   * **DDL-event driven (primary, low latency).** The streaming reader inspects 
raw `SourceRecord`s:
     a record whose `tableChanges` introduces a table that (a) is not in its 
split's table set / checkpoint
     state, and (b) matches the configured filter, becomes a 
`TableDiscoveryEvent`. This is deliberately
     done *before* `AbstractSchemaChangeResolver.support()`, so it does not 
change what counts as a sink
     schema change (`CREATE TABLE` is still not forwarded as a 
`SchemaChangeEvent`).
     Note: schema-change records are not subject to the per-table watermark 
filter
     (`IncrementalSourceStreamFetcher.java:298`), and with `table_pattern` the 
Debezium include list is the
     regex (fact 5), so these records do arrive.
   * **Polling (safety net).** The enumerator re-runs 
`dialect.discoverDataCollections(sourceConfig)`
     (cheap: `SHOW DATABASES` + `SHOW FULL TABLES ... BASE TABLE`,
     `connector-cdc-mysql/.../utils/TableDiscoveryUtils.java:37-91`) and diffs 
the result against the
     registry. Because the enumerator has no timer (fact 7), the tick is
     `notifyCheckpointComplete()`, throttled by a minimum interval
     (`table.discovery.interval`, default 60s; `0` disables). This covers: 
`schema.changes.enabled=false`,
     tables created before the stream started but after discovery, and DDL 
events that were filtered.
     **Dependency to state explicitly:** this channel only ticks when the job 
actually checkpoints; on a
     Zeta job without `checkpoint.interval` no checkpoint ever completes, so 
the DDL channel becomes the
     only trigger. Either document "polling requires checkpointing enabled" or 
give the enumerator a
     dedicated scheduled thread (see OQ5).
   * **Option `table.discovery.mode`** = `DISABLED` (default) | `DDL` | 
`POLLING` | `DDL_AND_POLLING`.
     `DISABLED` is byte-for-byte today's behaviour; recommended setting is 
`DDL_AND_POLLING`.
   
   ### 3.3 Enumerator admission state machine (question 2)
   
   On a candidate `TableId` `T` (from either channel):
   
   1. **Eligibility / dedup.** Reject if `T` is in `capturedTables` 
(`SplitAssigner.Context`), in
      `assignedSnapshotSplit`, in `splitCompletedOffsets`, or in the persisted 
admission registry
      (`admittedTables` / `pendingAdmissions`, 3.5). Re-check that `T` still 
matches the filter (the table
      may have been dropped between discovery and admission).
   2. **Schema / PK resolution at runtime.** `MySqlDialect.getPrimaryKey()` 
NPEs for tables missing from the
      submission-time `tableMap` 
(`connector-cdc-mysql/.../source/MySqlDialect.java:60-63`, `:129-131`), and
      sinks need the table's `CatalogTable`. The design resolves the table's 
schema/PK via JDBC on
      admission (sub-fix, §4.1).
   3. **Snapshot split creation.** Reuse the existing machinery:
      `dialect.createChunkSplitter(sourceConfig).generateSplits(T)` → 
`SnapshotSplit`(s) with the same
      low/high-watermark semantics as startup snapshots (`MySqlChunkSplitter`).
   4. **Owner selection (split reassignment, question 4).** Keep the bucket 
layout of
      `IncrementalSplitAssigner.createIncrementalSplits()` (round-robin over 
`incrementalParallelism`), so
      the table always goes to the reader that owns its bucket. This preserves 
the invariant *one table is
      streamed by exactly one reader* and keeps per-table ordering. If the 
owner is not registered yet
      (`Context.registeredReaders()`), the admission stays pending and is 
retried on the next
      `handleSplitRequest` / `notifyCheckpointComplete`. The admission snapshot 
split is **not** delivered
      through `getNext()` (it would be handed to an arbitrary reader and break 
the one-stream-fetcher rule);
      it is delivered together with the `TableAdmissionEvent`.
   5. **Record the decision** in the registry before the reader starts reading 
(so a restore never starts a
      second snapshot for the same table).
   
   ### 3.4 Reader-side backfill (the core change)
   
   `IncrementalSourceSplitReader` gains an admission sub-phase on top of the 
existing
   snapshot → incremental transition:
   
   1. The owner reader receives `TableAdmissionEvent(T, snapshotSplit)`. The 
split reader keeps it in a
      dedicated `pendingAdmission` slot instead of the normal `splits` queue.
   2. At the next `fetch()` boundary it **stops draining the stream fetcher** 
(the Debezium thread keeps
      running; the change-event queue absorbs the backlog — no offsets are 
committed, so nothing is lost).
      Two facts make this safe, and one cost must be stated, with the queue 
rule quoted exactly: for a
      stream split the queue is bounded (`queueSize = 
connectorConfig.getMaxQueueSize()`, plus
      `maxQueueSizeInBytes`, in `MySqlSourceFetchTaskContext`, defaults from 
Debezium `max.queue.size`) and
      Debezium offset flush is disabled (`offset.flush.interval.ms = 
Long.MAX_VALUE`), so the stream
      position only advances when records are emitted — pausing the drain 
freezes the position, it does not
      drop data. The cost: once the queue fills, the producer thread blocks and 
**all tables' streams
      stall** for the duration of the snapshot, so a very large new table 
pauses the whole job's intake.
      That is a real trade-off the reader-side design must surface (mitigation: 
keep draining while the
      scan runs, or use option C below). One correction to the naive reading of 
"bounded": the same file
      sets `queueSize = Integer.MAX_VALUE` when the fetched split **is** a 
snapshot split in exactly-once
      mode (`:142-153`, verbatim comment: "the queue needs to be set to a 
maximum size of `Integer.MAX_VALUE`
      (buffered a current snapshot all data)"). The admission read reuses 
exactly that snapshot-split path,
      so it gets **no** backpressure from its own queue — the reader must keep 
draining it eagerly, as the
      startup snapshot does today; the bounded-queue stall above applies to the 
*stream* split only.
   3. It runs `IncrementalSourceScanFetcher` on the admission snapshot split to 
completion, obtaining
      `CompletedSnapshotSplitInfo` with a high watermark `W_high`.
   4. It folds `T` into the running split: 
`IncrementalSourceStreamFetcher.addAdmittedTable(T, W_high)`
      refreshes the filter (`configureFilter()`), i.e. 
`maxSplitHighWatermarkMap[T] = W_high`, and the
      mutated `IncrementalSplit` now contains `T` with `W_high` as its 
per-table startup offset.
   5. Streaming resumes. For `T`: positions `< W_high` are dropped (already 
covered by the snapshot),
      positions `>= W_high` are emitted exactly once. This is the same 
exactly-once boundary the startup
      snapshot phase already uses.
   6. The reader emits the existing `CompletedSnapshotSplitsReportEvent` so the 
split assigner records
      `splitCompletedOffsets[T]` — 
`IncrementalSourceReader.reportFinishedSnapshotSplitsIfNeed()`
      (`:212-234`) already builds this event (today it is called from 
`addSplits()` `:178` and from
      `onSplitFinished()` `:208`) and 
`IncrementalSplitAssigner.onCompletedSplits()` (`:158-163`) already
      stores it — and the enumerator acknowledges it with 
`CompletedSnapshotSplitsAckEvent`.
   
   **Scope: `exactly-once` is required for v1.** The `>= W_high` guarantee in 
step 5 comes from the
   exactly-once branch of `IncrementalSourceStreamFetcher.shouldEmit()` 
(`:265-300`); with
   `exactly-once = false` the filter only compares against the split's 
`splitStartWatermark`, and the
   snapshot unit's backfill read is skipped entirely 
(`MySqlSnapshotFetchTask.java:86-92`). V1 therefore
   validates `exactly-once` at admission time and fails with a clear 
configuration error otherwise;
   defining at-least-once admission semantics is a follow-up.
   
   **Backfill mechanics — three concrete options, with the facts that decide 
between them.** How the new
   table's existing rows get their consistent binlog watermark:
   
   * **Option A — reuse the existing snapshot unit as-is, with the live client 
kept connected (buffering
     variant described above).** The unit is `MySqlSnapshotFetchTask` 
(`:57-119`):
     `MySqlSnapshotSplitReadTask.doExecute` (`:127-159`) takes `low = 
currentBinlogOffset(jdbcConnection)`
     (`:137`) → `createDataEvents(...)` (`:147`) → `high = 
currentBinlogOffset(jdbcConnection)` (`:149`)
     **with no read lock at all** (it extends Debezium's 
`AbstractSnapshotChangeEventSource` directly; the
     forked `MySqlSnapshotChangeEventSource` that carries the `FLUSH TABLES ... 
WITH READ LOCK` code is dead
     code — nothing references it), and then, in exactly-once mode only, 
`createBackfillBinlogReadTask()`
     (`:143-171`) runs a **bounded per-table binlog read of exactly `[low, 
high]`** (`table.include.list` =
     the one table, `include.schema.changes=false`). So "consistent watermark, 
no gap" is already
     implemented — but only for the startup snapshot phase, and the 
*concurrency* does not carry over: the
     backfill opens its own binlog client 
(`MySqlSourceFetchTaskContext.getBinaryLogClient()` `:216`,
     `MySqlConnectionUtils.java:61`), whose `server_id` is `database.server.id =
     ServerIdRange.getServerId(subtaskId)` 
(`MySqlSourceConfigFactory.java:99-102`), and
     `MySqlStreamingChangeEventSource.java:237` sets 
`client.setServerId(connectorConfig.serverId())`. Two
     clients in one subtask therefore collide and MySQL drops the older 
connection, so reusing this unit
     mid-stream **does** require reserving a second `server_id` slot per 
subtask (e.g. `2*id` / `2*id+1`).
     Cost: that reservation, no read lock, one extra connection briefly, plus 
the queue stall of step 2.
   * **Option B — Debezium's read-only incremental snapshot (lock-free).** The 
fork already contains
     `MySqlReadOnlyIncrementalSnapshotChangeEventSource`
     
(`connector-cdc-mysql/.../io/debezium/connector/mysql/MySqlReadOnlyIncrementalSnapshotChangeEventSource.java:134`)
     but it is **not wired in anywhere** (only self-references). It avoids 
locks by replaying binlog in
     `[low, high]` for the snapshotted table — which requires a **second binlog 
connection** in the same
     reader, and the same `server_id` (`ServerIdRange.getServerId(subtaskId)` 
is already held by the live
     stream) would make MySQL drop one of the two connections. So option B 
additionally requires reserving
     a second server-id slot per subtask (e.g. `2*id` / `2*id+1`) or stopping 
the main stream first —
     real work, but it is the only chunked/lock-free variant.
   * **Option C — bounded-stop the stream, reuse the same unit, restart.** 
`MySqlBinlogSplitReadTask`
     already supports a bounded stop (`MySqlBinlogFetchTask.java:77`, `:289`). 
Bounded-stop the running
     stream split at `P` = the position of the last record it emitted (the 
restart must be able to take
     `P` as its new start position, which is part of the work this option 
needs), run the *same*
     `[low, high]` snapshot+backfill unit for `T`, then start a fresh stream 
fetcher for the extended
     table set from `P`. No second `server_id`
     is needed (the live client is gone while the backfill client runs) and 
there is no queue pause;
     correctness holds because records of `T` in `[P, W_high)` are dropped by 
the filter while every other
     table resumes from `P` with no gap. Cost: the stream fetcher needs a 
bounded-stop/restart path it does
     not exercise today (it is created once per incremental split), and each 
admission costs one stop/start,
     so a `CREATE TABLE` burst could thrash — admissions should be batched per 
checkpoint window.
   
   Steps 1–6 above are written for Option A; under Option C step 2 becomes a 
bounded stop and step 5 a
   fresh stream fetcher. With the corrected facts, **A and B both need a 
reserved second `server_id`**
   (A: a backfill client while the live client stays connected; B: a chunked 
read-only-snapshot client),
   while **C needs no extra `server_id`** but does need a stream restart path. 
My default is therefore
   **Option C for v1** (the existing unit is reused verbatim, no new connection 
semantics), with Option A
   as the alternative if reserving server ids is preferred over adding a 
restart path, and B as the
   documented scaling path for very large tables — maintainer direction 
welcome, see OQ3.
   
   **The most fragile point in this design (called out honestly):** the 
hand-off where the snapshot
   finishes and the stream filter starts emitting `T` at exactly `W_high`. The 
emission order to the sink
   must be *snapshot rows first, then buffered/streamed binlog records for 
`T`*, and any record of `T`
   with position `< W_high` must be dropped exactly once — not zero, not twice. 
The machinery is the same
   as the startup snapshot boundary (`CompletedSnapshotSplitInfo` intervals + 
high-watermark filtering +
   `finishedUnackedSplits`/ack), but here it fires **mid-stream** instead of at 
phase transition, so it is
   the single easiest place to get wrong. Unit tests 3 and 4 in the test plan 
exist specifically to pin
   this boundary down (including "record arrives for `T` while the snapshot is 
still running" and
   "restore between snapshot completion and the admission checkpoint").
   
   **Rejected: a dedicated snapshot-only reader task.** Parallelism is fixed by 
the job and
   `incrementalParallelism` buckets must stay stable; adding tasks mid-run is 
not expressible in the
   current enumerator contract.
   
   ### 3.5 What is persisted in checkpoint state (question 3)
   
   * **Enumerator side** — extend `IncrementalPhaseState`
     (`.../source/enumerator/state/IncrementalPhaseState.java`, today a single 
`Offset startupOffset`,
     `serialVersionUID = -6809026812298443356L`) with:
     * `Set<TableId> admittedTables` — tables whose admission has been 
committed (dedup + restore guard);
     * `Set<TableId> pendingAdmissions` — assigned but not yet snapshotted (so 
a restore re-assigns the
       same split id instead of creating a new one);
     * `Map<TableId, Offset> tableHighWatermarks` — `W_high` per admitted 
table, needed to rebuild the
       stream filter on restore *before* the reader re-acknowledges.
   * **Reader side** — unchanged mechanism: 
`IncrementalSourceReader.snapshotState()`
     (`.../reader/IncrementalSourceReader.java:437-457`) already checkpoints 
unfinished splits and rewrites
     the incremental split with `checkpointTables` / `historyTableChanges`; the 
admitted table simply
     appears in the mutated `IncrementalSplit` (table set + per-table offsets). 
After #12567 lands, the
     per-table startup offset is already carried by 
`IncrementalSplit`/`IncrementalSplitState`.
   * **Compatibility.** New fields default to empty; a checkpoint/savepoint 
written by an older build
     restores with exactly today's behaviour (no admitted tables), and a build 
with the feature disabled
     never writes them.
   
   ### 3.6 Failure / recovery behaviour (question 5)
   
   * **Failure during backfill.** The admission snapshot split is in 
`pendingAdmissions` and in the
     reader's unfinished splits. On restore the enumerator re-assigns it to the 
same bucket owner and the
     standard snapshot replay path applies 
(`IncrementalSourceReader.addSplits()` /
     `finishedUnackedSplits`); nothing is committed to the stream filter until 
the admission is
     checkpointed.
   * **Failure after snapshot completion, before the admission checkpoint.** 
`capturedTables` from
     discovery minus `checkpointCapturedTables` (`IncrementalSource.restore()`, 
`:427-482`) still reports
     `T` as new, so it is snapshotted again. Records already emitted after 
`W_high` may therefore be
     re-emitted once — the same guarantee the existing startup snapshot gives, 
and the reason the sink must
     be idempotent (`PRIMARY KEY` upsert / exactly-once sink) for this feature. 
This window is documented,
     not silently accepted.
   * **Restore with stream-only state.** The stream filter for `T` is rebuilt 
from the persisted
     `tableHighWatermarks`; the stream itself replays from the last committed 
offset and drops everything
     below `W_high`.
   * **Old state / disabled feature.** No-op; `table_pattern` restore semantics 
are untouched (section 5).
   
   ### 3.7 Duplicate-admission suppression (question 6)
   
   * Only the enumerator creates admissions; readers cannot mutate their 
captured set except through
     `TableAdmissionEvent`.
   * One eligibility check against the union of `capturedTables`, 
`assignedSnapshotSplit`,
     `splitCompletedOffsets`, `admittedTables`, `pendingAdmissions` — a 
`TableId` that is already known in
     any form is rejected.
   * Keyed by `TableId`, so a DDL event and a poll discovering the same table 
in the same checkpoint window
     collapse into one admission.
   * The registry is written before the reader starts; after a restore the 
newly-discovered table set is
     diffed against it (existing `restore()` semantics), so no second snapshot 
is started for a table whose
     admission was already checkpointed.
   
   ### 3.8 MySQL vs OceanBase scope (question 7)
   
   * **base (`connector-cdc-base`)**: event types, `TableAdmission` coordinator 
in the enumerator,
     `IncrementalSourceSplitReader` admission sub-phase,
     `IncrementalSourceStreamFetcher.addAdmittedTable()`, 
`IncrementalPhaseState` fields. Dialect-agnostic.
   * **MySQL (`connector-cdc-mysql`)**: (a) raw-record inspection for `CREATE 
TABLE` discovery and
     `TableId` extraction; (b) runtime PK/schema lookup fix in 
`MySqlDialect.getPrimaryKey()`; (c) reuse
     `MySqlChunkSplitter` unchanged; (d) ITs.
   * **OceanBase**: inherited — `OceanBaseIncrementalSource extends 
MySqlIncrementalSource`
     (`connector-cdc-oceanbase/.../source/OceanBaseIncrementalSource.java:35`) 
and
     `OceanBaseIncrementalSourceFactory extends MySqlIncrementalSourceFactory` 
(`:46`). The only OceanBase
     work is an IT proving parity (DDL capture and `table_pattern` filter 
behave as on MySQL).
   * **`table_names` mode**: nothing can match a new table, so the feature is 
inert — no behaviour change.
   * **Other dialects** (Postgres/SqlServer/DB2/Vitess/Mongo/TiDB): untouched; 
the option defaults to
     `DISABLED`.
   
   ## 4. Required sub-fixes (called out so they are not surprises in review)
   
   1. **`MySqlDialect.getPrimaryKey()` cannot take an unknown table.** It reads 
from a `tableMap` built from
      the submission-time catalog (`:60-63`, `:129-131`). Admission must 
resolve PK/schema live (e.g. via
      `MySqlSchema`/`information_schema`); otherwise chunk splitting for the 
new table NPEs.
   2. **Sink-side schema for a table that was never in the job config.** 
`getProducedCatalogTables()` is
      fixed at submission, so a runtime-discovered table has no row 
type/`CatalogTable` for the sink. For
      this design the limit is accepted and documented: live discovery targets 
sinks that resolve schema
      dynamically or are append-only, and the case fails fast with a clear 
error otherwise. Turning the
      `CREATE TABLE` record into a `CatalogTable` through the existing 
schema-change plumbing
      (`SchemaChangeType.CREATE_TABLE` is currently not in 
`MySqlIncrementalSource.supports()`, `:320-329`)
      is **out of scope for this design** and only raised as OQ2, so the 
maintainers can decide whether a
      follow-up is wanted at all.
   3. `IncrementalSourceSplitReader` must be changed at all (fact 3): today a 
live stream fetcher makes
      every queued split unreachable. This is unavoidable for any in-place 
admission.
   
   ## 5. Backward compatibility of `table_pattern` (as requested)
   
   * Discovery is still evaluated exactly where it is today — 
`createEnumerator()` (`:303-358`) and
     `restoreEnumerator()` (`:360-425`) — and the restore diff logic
     (`IncrementalSource.restore()`, `:427-482`) is **not modified**. This 
design adds tables that appear
     *later*; it never changes how the startup table set is computed or 
restored.
   * No change to `table_names` mode; no change when `table.discovery.mode = 
DISABLED` (default).
   * Restore state written by a previous version restores unchanged (new fields 
default-empty).
   * Per-table ordering and "one table ↔ one reader" ownership are preserved by 
reusing the existing bucket
     layout.
   * Test items (a), (c) and (d) from the earlier proposal are preserved and 
extended rather than dropped:
     (a) pattern-compile errors still come from Debezium's 
`RelationalTableFilters` (untouched here);
     (c) exact-table-name behaviour is test item 7; (d) restore-path 
re-discovery is unchanged and covered
     by test items 2 and 7.
   
   ## 6. Test plan (starting from the (c)/(d) items, extended)
   
   1. **Unit, enumerator dedup**: DDL event and poll in the same checkpoint 
window produce one admission;
      admission for an already-captured table is ignored; with 
`incrementalParallelism > 1` the admission
      lands on the bucket owner.
   2. **Unit, state round-trip**: `IncrementalPhaseState` with/without the new 
fields; an old-serialised
      instance deserialises and behaves as `DISABLED`; admission survives 
snapshot → restore → no second
      snapshot.
   3. **Unit, split reader interleaving**: live stream fetcher + 
`TableAdmissionEvent` → snapshot read runs
      to completion, then the filter gains `T` with `W_high`; records `< 
W_high` are dropped and
      `>= W_high` emitted.
   4. **Unit, `shouldEmit`**: exactly-once path for an admitted table (`< 
W_high` dropped, `>= W_high`
      emitted); with `exactly-once = false` an admission attempt is rejected 
with a clear error (v1 scope,
      section 3.4).
   5. **MySQL IT**: job with `table_pattern` running; create a new matching 
table mid-run and insert rows;
      assert all pre-existing rows land exactly once **and** post-creation 
changes are applied; then a
      failure injection between snapshot completion and the first checkpoint, 
asserting no loss and (with a
      PK sink) no duplicates.
   6. **OceanBase IT**: same scenario, proving the inheritance needs no 
MySQL-only assumptions.
   7. **Compat IT**: `table.discovery.mode=DISABLED` behaves identically to 
today; `table_names` mode with a
      new table created mid-run is unaffected; a checkpoint produced without 
the feature restores.
   
   ## 7. Rollout
   
   * New options on the MySQL factory (inherited by OceanBase): 
`table.discovery.mode`
     (`DISABLED` default) and `table.discovery.interval` (polling, default 60s).
   * Suggested PR split: **PR1** base machinery + MySQL PK/schema lookup fix 
(no trigger wired in yet);
     **PR2** DDL channel; **PR3** polling channel + option docs; **PR4** ITs. 
Landing after #12567 keeps the
     per-table offset work single-threaded, as requested.
   
   ## 8. Open questions for maintainers (OQ1–OQ5; referenced by label from §3)
   
   1. **OQ1** `table.discovery.mode` default: `DISABLED` (strictly backward 
compatible) or `DDL` (feature visible
      but still no polling)? I lean `DISABLED` for the first PR.
   2. **OQ2** Should the new table's `CREATE TABLE` become a `CREATE_TABLE` 
schema change (adding `CatalogTable`
      to the produced catalog so ordinary sinks work), or is a documented sink 
limitation acceptable for
      the first iteration?
   3. **OQ3** Reader-side backfill mechanics. All three options reuse the 
existing unit — lock-free
      `low → SELECT → high` plus a bounded table-only binlog backfill `[low, 
high]`
      (`MySqlSnapshotFetchTask.java:57-119`, `:143-171`) — which today exists 
only in exactly-once mode;
      they differ in concurrency:
      **(A)** keep the live client connected and pause the drain while the 
owner reader runs that unit for
      `T`: needs a **reserved second `server_id`** per subtask (the backfill 
client and the live client
      would otherwise share `database.server.id`, and MySQL drops the older 
connection), plus the queue
      stall; **(B)** Debezium's read-only incremental snapshot — lock-free and 
chunked, so better for very
      large tables, but the class is unwired in this tree and it needs the same 
second `server_id`;
      **(C)** bounded-stop the live stream client, run the same unit unchanged, 
restart the stream from the
      last emitted position — **no** second `server_id` and no queue pause, at 
the price of a restart path
      the stream fetcher does not have today and one stop/start per admitted 
table (batched per checkpoint
      window). My default is (C); (A) if you would rather reserve server ids 
than add the restart path.
      Maintainer direction is welcome here, in particular from the #12567 
review perspective.
   4. **OQ4** `pendingAdmissions` as persisted state vs. derived from reader 
state — the former is cheaper to
      reason about but adds fields to `IncrementalPhaseState`; the latter 
avoids the extra set but makes the
      enumerator depend on reader state.
   5. **OQ5** Polling needs a tick; the enumerator only gets 
`notifyCheckpointComplete`. Is "at most once per
      checkpoint, throttled by `table.discovery.interval`" acceptable, or would 
you prefer a dedicated
      scheduled thread inside the enumerator?
   
   ---
   
   ### Appendix: relevant evidence index
   
   * `connector-cdc-base/.../source/IncrementalSource.java` — 
`createEnumerator()` `:303-358`,
     `restoreEnumerator()` `:360-425`, `restore()` `:427-482`.
   * 
`connector-cdc-base/.../source/enumerator/IncrementalSourceEnumerator.java` — 
`handleSourceEvent()`
     `:103`, `notifyCheckpointComplete()`, `assignSplits()`.
   * `connector-cdc-base/.../source/enumerator/IncrementalSplitAssigner.java` — 
`createIncrementalSplits()`,
     `createIncrementalSplit()`, `snapshotState()` → 
`IncrementalPhaseState(startupOffset)`.
   * `connector-cdc-base/.../source/reader/IncrementalSourceSplitReader.java` — 
`:152-156`, `:190-192`.
   * 
`connector-cdc-base/.../source/reader/external/IncrementalSourceStreamFetcher.java`
 — `shouldEmit()`
     `:265-300`, `hasEnterPureBinlogPhase()` `:301-313`, `configureFilter()` 
`:315+`.
   * `connector-cdc-base/.../source/reader/IncrementalSourceReader.java` — 
`discoverCapturedTables()`
     `:393-403`, `snapshotState()` `:437-457`.
   * `connector-cdc-mysql/.../config/MySqlSourceConfigFactory.java:104-113` — 
`table.include.list`.
   * `connector-cdc-mysql/.../source/MySqlDialect.java` — `:60-63`, `:94-102`, 
`:129-131`.
   * `connector-cdc-mysql/.../utils/TableDiscoveryUtils.java:37-91` — 
`listTables()`.
   * `connector-cdc-base/.../schema/AbstractSchemaChangeResolver.java:56`, 
`:68-80` — `SUPPORT_DDL`.
   * `connector-cdc-mysql/.../source/MySqlIncrementalSource.java:320-329` — 
`supports()`.
   * `seatunnel-engine/.../task/SourceSplitEnumeratorTask.java:334` — `run()` 
called once.
   * 
`connector-cdc-mysql/.../source/reader/fetch/scan/MySqlSnapshotFetchTask.java:57-119`
 (snapshot then
     `createBackfillBinlogSplit` `:121-129` / `createBackfillBinlogReadTask` 
`:143-171`, backfill only when
     `isExactlyOnce()`) and `.../scan/MySqlSnapshotSplitReadTask.java:127-159` 
(no read lock).
   * 
`connector-cdc-mysql/.../source/reader/fetch/MySqlSourceFetchTaskContext.java:142-161`
 — queue size
     rule (`Integer.MAX_VALUE` for snapshot splits in exactly-once); `:216` 
`getBinaryLogClient()`.
   * `connector-cdc-mysql/.../utils/MySqlConnectionUtils.java:61` and
     `io/debezium/connector/mysql/MySqlStreamingChangeEventSource.java:227`, 
`:237` — binlog client
     creation and `server_id` assignment (the collision behind OQ3 A/B).
   * Dead fork residue (zero external references): 
`io/debezium/connector/mysql/MySqlSnapshotChangeEventSource`
     (the `FLUSH TABLES ... WITH READ LOCK` class) and
     `MySqlReadOnlyIncrementalSnapshotChangeEventSource`.
   * PR #12567 (Asthenia0412, `e18f31cfc`, +435/−22) / issue #12565 — pending 
per-table restore watermarks.
   


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