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]
