dingsongjie commented on PR #12308:
URL: https://github.com/apache/seatunnel/pull/12308#issuecomment-5693656804

   > Thanks for this contribution, @dingsongjie — welcome to the Apache 
SeaTunnel community! This is a well-structured feature: the per-table isolation 
in the reader (`preWrites`/`commits`/`committedEvents` keyed by table), the 
enumerator's table-diff logic, and the depth of the new unit/E2E test suite are 
genuinely good engineering. I did a full local checkout and read through the 
diff end-to-end; I found one backward-compatibility bug that needs to be closed 
before this can be merged, plus a handful of smaller items.
   > 
   > # What Problem Does This PR Solve?
   > * **User pain point**: Today one `TiDB-CDC` source block can only capture 
a single table. Capturing N tables requires N separate source blocks (N TiKV 
sessions, N sets of connector overhead), which is wasteful and awkward to 
manage.
   > * **Fix approach**: Adds a `table-names` list option alongside the legacy 
`database-name`/`table-name` pair. `TiDBSource`/`TiDBSourceFactory` now resolve 
to a list of `CatalogTable`s, `TiDBSourceSplitEnumerator` tracks per-table 
enumeration state (`enumeratedTables`) so it can detect tables added across a 
restore, and `TiDBSourceReader` keys its transaction-assembly buffers and 
deserializers per table so splits from different tables can be interleaved 
safely on one reader instance.
   > * **One-sentence summary**: Multiplexes multiple TiDB tables onto a single 
`TiDB-CDC` source/reader instance while trying to preserve restore 
compatibility for existing single-table jobs.
   > 
   > # 1. Code Change Review
   > ## 1.1 Core Logic Analysis
   > **`TiDBSourceOptions.getTableFullNames()`** 
(`seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/config/TiDBSourceOptions.java:131-147`)
 prefers `table-names`, dedupes with a `LinkedHashSet`, validates each entry 
has a `.` separator, and falls back to `database-name`/`table-name` when 
`table-names` is absent — this is a clean, well-tested compatibility shim (see 
`TiDBSourceFactoryTest`).
   > 
   > **`TiDBSourceSplitEnumerator.run()`** 
(`.../enumerator/TiDBSourceSplitEnumerator.java:152-191`) is the part that 
matters most for correctness:
   > 
   > ```java
   > if (shouldEnumerate) {
   >     sourceSplits = getTiDBSourceSplit(tableIds.keySet());
   >     ...
   >     enumeratedTables.addAll(tableIds.keySet());
   >     addPendingSplit(sourceSplits);
   >     shouldEnumerate = false;
   > } else if (enumeratedTables != null) {
   >     Set<String> missingTables = new HashSet<>(tableIds.keySet());
   >     missingTables.removeAll(enumeratedTables);
   >     if (!missingTables.isEmpty()) {
   >         sourceSplits = getTiDBSourceSplit(missingTables);
   >         ...
   >     }
   > }
   > ```
   > 
   > `enumeratedTables` is populated from 
`TiDBSourceCheckpointState.getEnumeratedTablesRef()` in the constructor 
(`TiDBSourceSplitEnumerator.java:79-87`). `TiDBSourceCheckpointState` 
(`.../enumerator/TiDBSourceCheckpointState.java:41-46`) is a pre-existing 
`Serializable` class (unchanged `serialVersionUID = 6292978509042158791L`); the 
new `enumeratedTables` field is simply absent from any checkpoint/savepoint 
binary written by the connector version that is currently in `dev`/released 
today. Per standard Java serialization semantics, deserializing such a byte 
stream into the new class definition leaves the new field `null` — the class's 
own Javadoc (lines 41-45) even documents this: _"Null when deserialized from a 
legacy checkpoint, in which case newly added tables are not discovered on 
restore."_
   > 
   > Tracing the actual runtime path for that case:
   > 
   > 1. Existing production job runs today's single-table connector 
(`database-name`/`table-name`), takes a savepoint. `shouldEnumerate=false` is 
already persisted (enumeration happened long ago).
   > 2. User upgrades to this PR's connector version and restores the same job, 
now with `table-names` expanded to include a second table.
   > 3. `TiDBSourceSplitEnumerator` constructor sets `enumeratedTables = null` 
(from `getEnumeratedTablesRef()`), `shouldEnumerate = false` (restored).
   > 4. `run()`: `shouldEnumerate` is `false` → first branch skipped. 
`enumeratedTables != null` is `false` → second branch skipped entirely.
   > 5. Result: the enumerator calls `assignSplit(readers)` with only whatever 
was in the restored `pendingSplit` (i.e. nothing new), then 
`signalNoMoreSplits`. **The newly added table never gets a split, never runs a 
snapshot, and is never streamed — silently, with no warning or error logged.** 
The job reports healthy; the new table's sink simply stays empty forever.
   > 
   > This is not a hypothetical: the PR's own new unit test **proves** it — 
`TiDBSourceSplitEnumeratorTest.runShouldNotEnumerateAnythingForLegacyCheckpointWithoutEnumeratedTables()`
 (`.../enumerator/TiDBSourceSplitEnumeratorTest.java:265-284`) constructs 
exactly this scenario (`new TiDBSourceCheckpointState(false, 
Collections.emptyMap())` = a legacy-shaped state, two configured tables) and 
asserts `context.getAssignedSplits(0/1).isEmpty()` and 
`currentUnassignedSplitSize() == 0` — i.e. the test locks in "nothing gets 
assigned" as the expected outcome, rather than treating it as a bug to fix.
   > 
   > Compare this with the E2E coverage: 
`testTiDBCdcSavepointRestoreWithAddedTable` (`TiDBCDCIT.java:388`) starts its 
"initial" job from `tidbcdc_multi_table_to_tidb_single.conf`, which already 
uses the **new** `table-names` option with one entry — so that savepoint 
already contains a non-null `enumeratedTables` (the "happy path" `else if 
(enumeratedTables != null)` branch, which does work correctly per 
`runShouldEnumerateOnlyTablesAddedAfterCheckpointOnRestore`). There is no E2E 
(and no positive unit test) that restores a savepoint taken with the _actual_ 
pre-PR connector version and confirms the newly-added table gets discovered — 
because it doesn't.
   > 
   > This directly contradicts the PR's own doc update: _"When a job is 
restored from a savepoint, tables added to `table-names` run a fresh snapshot 
before joining the incremental stream"_ 
(`docs/en/connectors/source/TiDB-CDC.md`, Notes section) — that promise is 
false for any job that was running the connector before this PR.
   > 
   > **Runtime path diagram** (legacy-checkpoint restore with an added table):
   > 
   > ```
   > Old job (pre-PR, single table) --savepoint--> binary WITHOUT 
`enumeratedTables` field
   >         |
   >         v
   > New connector version, config: table-names=[old_table, new_table]
   >         |
   >         v
   > TiDBSourceSplitEnumerator(ctor) --deserialize--> shouldEnumerate=false, 
enumeratedTables=null
   >         |
   >         v
   > run(): shouldEnumerate? no --> enumeratedTables != null? no --> [both 
branches skipped]
   >         |
   >         v
   > assignSplit(readers) using only restored pendingSplit (old_table only)
   >         |
   >         v
   > signalNoMoreSplits  ==>  new_table: NEVER snapshotted, NEVER streamed, NO 
error/warning logged
   > ```
   > 
   > ## 1.2 Compatibility Impact
   > **Partially incompatible.** Config-level backward compatibility 
(`database-name`/`table-name` still works) is solid. Split-level state 
(`TiDBSourceSplit`) is unchanged and deserializes fine from old checkpoints. 
But **checkpoint/savepoint-level compatibility for the exact upgrade scenario 
this feature is meant to enable — adding tables to an existing job — is broken 
silently** for any checkpoint taken before this PR. This is a documented (in 
the Javadoc, even) but unresolved gap, not an unknown unknown.
   > 
   > ## 1.3 Performance / Side-Effect Analysis
   > * `committedEvents`/`preWrites`/`commits` are now `Map<String, ...>` keyed 
per table instead of one instance per reader (`TiDBSourceReader.java:84-89`). 
The existing in-code comment ("5000W queue size may be safe size") was written 
for a single unbounded queue per reader; with N tables multiplexed onto one 
reader, worst-case memory is now N × that size if a downstream sink stalls on 
all tables simultaneously. Not a regression in per-table behavior, but the 
aggregate resource envelope per reader instance grows with table count — worth 
a callout in docs/config guidance for users planning wide multi-table jobs, not 
a blocker.
   > * Round-robin split assignment (`getSplitOwner`) now mixes splits from 
different tables across readers, which is expected and fine given 
`TiDBSourceSplit.tableFullName()` based dispatch downstream.
   > 
   > ## 1.4 Error Handling and Logging
   > **Issue 1 — Legacy checkpoint restore silently drops newly added tables 
(data loss, no error/warning)**
   > 
   > * **Location**: 
`seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/enumerator/TiDBSourceSplitEnumerator.java:152-191`
 (specifically the `else if (enumeratedTables != null)` guard at line 168), 
backed by `TiDBSourceCheckpointState.java:41-45`.
   > * **Problem description**: When restoring from a checkpoint written by the 
connector version currently in `dev` (i.e. any real, existing single-table 
TiDB-CDC job), `enumeratedTables` deserializes to `null`. Because 
`shouldEnumerate` is also already `false` in that same legacy state, both 
branches in `run()` are skipped, so newly configured tables never get splits, 
never snapshot, and never stream — with zero log output indicating anything was 
skipped.
   > * **Potential risk**: Silent, permanent data loss for any table added 
during a savepoint-restore upgrade of an existing production TiDB-CDC job — the 
exact upgrade path this feature is supposed to support, and the one the PR's 
own docs claim works.
   > * **Best improvement**: When `enumeratedTables == null` post-restore, 
treat it as "everything currently in `pendingSplit`'s table set was already 
enumerated" (derive the legacy-enumerated set from the restored 
splits/`tableIds` that were live before this change) and run the "enumerate the 
difference" logic against the full configured table set — the same way the 
non-null branch already does. At minimum, if that derivation is judged unsafe, 
fail restore loudly with a clear exception telling the operator to do a fresh 
start instead of silently dropping tables. Either approach is acceptable; 
silently doing nothing is not.
   > * **Severity**: Blocker (P0 — silent data loss on a documented, 
in-repo-tested-as-current-behavior upgrade path).
   > * **Raised by another reviewer**: Yes — @goutamadwant flagged this in an 
inline comment on the same lines before this review; I independently reproduced 
and confirmed the root cause via the checkpoint class, the enumerator logic, 
and the PR's own unit test that asserts the broken behavior as expected.
   > 
   > # 2. Code Quality Assessment
   > ## 2.1 Coding Standards
   > New core methods/classes 
(`TiDBSourceOptions.getTableFullNames/parseDatabaseName/parseTableName/tableFullName`,
 `TiDBSourceConfig.getTableFullNames/hasConfiguredTables`, 
`TiDBSourceSplitEnumerator.discardRemovedTables`, 
`TiDBSourceCheckpointState.enumeratedTables`) all carry Javadoc explaining 
purpose and edge cases — good discipline, and the `TiDBSourceCheckpointState` 
Javadoc is honest about the exact gap this review flags (it just doesn't close 
it). Minor: the `TABLE_NAMES` option description says _"Mutually exclusive with 
database-name/table-name"_ (`TiDBSourceOptions.java:58-59`), but the actual 
behavior (and the user-facing doc table) is "takes precedence when both are 
set," not a validation failure when both are present — the option Javadoc 
should be corrected to match the real (and better) behavior.
   > 
   > ## 2.2 Test Coverage and Test Stability
   > **Risk present.** Unit coverage for the reader (`TiDBSourceReaderTest`, 
incl. `transactionBuffersShouldBeIsolatedPerTable`) and for enumerator table 
add/remove within the _new_ checkpoint format is thorough and well-targeted. 
However:
   > 
   > * 
`TiDBSourceSplitEnumeratorTest.runShouldNotEnumerateAnythingForLegacyCheckpointWithoutEnumeratedTables`
 (`enumerator/TiDBSourceSplitEnumeratorTest.java:265-284`) asserts the broken 
behavior described in Issue 1 as correct, which will need to be rewritten (not 
just left in place) once the fix lands.
   > * No E2E test restores a savepoint that mimics a genuinely pre-PR 
checkpoint (i.e. one taken with `table-names` absent, using only 
`database-name`/`table-name`) and then adds a table — every E2E "add table" 
scenario starts from a job already running the new code with `table-names` set. 
This is the actual real-world upgrade path and it isn't covered.
   > 
   > ## 2.3 Documentation Updates
   > `docs/en` and `docs/zh` are both updated in parallel with a new 
"Multi-Table Sync" example and updated option table/notes — good bilingual 
discipline. The added Notes-section compatibility claim ("tables added to 
`table-names` run a fresh snapshot [...] after restore") should either be 
scoped to "when restoring from a checkpoint taken with `table-names` already 
configured" or removed until Issue 1 is fixed, so the docs don't overstate the 
guarantee.
   > 
   > # 3. Architectural Soundness
   > ## 3.1 Elegance of the Solution
   > Precise fix for the "one source, many tables" goal — keying per-table 
state throughout the reader is the right shape rather than bolting a loop on 
top. The checkpoint-compatibility gap is a genuine gap in an otherwise 
well-thought-out design (the author clearly considered it, given the Javadoc), 
not a symptom of a wrong overall approach.
   > 
   > ## 3.2 Maintainability
   > Good — `TiDBSourceOptions` centralizes all name-parsing/validation logic 
used consistently by the factory, source, enumerator and reader, avoiding 
duplicated parsing.
   > 
   > ## 3.3 Extensibility
   > Reasonable. `List<CatalogTable>` plumbed end-to-end makes it 
straightforward to add further per-table behavior later.
   > 
   > ## 3.4 Historical-Version Compatibility
   > This is the core blocker in this PR — see 1.1/1.2/Issue 1. Everything else 
(config option compatibility, split state compatibility) checks out; 
checkpoint-state compatibility for the add-table-on-restore path does not.
   > 
   > # 4. Issue Summary
   > #  Issue   Location        Severity
   > 1  Legacy (pre-PR) checkpoint restore silently drops newly added tables — 
no snapshot, no streaming, no warning    
`TiDBSourceSplitEnumerator.java:152-191`, 
`TiDBSourceCheckpointState.java:41-45`        Blocker (P0)
   > 2  `TABLE_NAMES` option Javadoc says "mutually exclusive" but 
actual/documented behavior is "takes precedence"     
`TiDBSourceOptions.java:51-59`  Minor
   > 3  Per-table unbounded `committedEvents` queues multiply the worst-case 
memory envelope by table count vs. the single-table design the sizing comment 
was written for      `TiDBSourceReader.java:84-89`   Minor / doc-note
   > 4  New unit test locks in the Issue 1 behavior as "expected" rather than 
testing for correct behavior      `TiDBSourceSplitEnumeratorTest.java:265-284`  
  Follow-up of Issue 1
   > # 5. Merge Recommendation
   > ### Conclusion: Ready to merge after fixes
   > 1. **Blockers — must be fixed**
   >    
   >    * Issue 1: close the legacy-checkpoint (`enumeratedTables == null`) gap 
so tables added during a savepoint-restore of an _existing, pre-this-PR_ 
TiDB-CDC job are actually discovered and snapshotted — or, if that's judged 
infeasible, fail restore explicitly with a clear message instead of silently 
doing nothing. Update 
`runShouldNotEnumerateAnythingForLegacyCheckpointWithoutEnumeratedTables` to 
assert the fixed behavior, and add an E2E (or at least a more realistic unit 
test) that starts from a state shaped like a genuinely pre-PR checkpoint (no 
`table-names` used at all for the initial job) before adding a table on restore.
   >    * Scope the docs' restore-compatibility claim in 
`docs/en|zh/connectors/source/TiDB-CDC.md` to match whatever the fixed behavior 
actually guarantees.
   > 2. **Recommended fixes — non-blocking**
   >    
   >    * Fix the `TABLE_NAMES` option description (Issue 2).
   >    * Note the multi-table memory-envelope implication of per-table 
unbounded queues in the docs or as a follow-up sizing guideline (Issue 3).
   > 
   > **Overall assessment**: The multi-table design itself — per-table buffer 
isolation in the reader, table-diffing in the enumerator, config-level backward 
compatibility — is solid, well-tested where it's tested, and clearly the result 
of careful thought (the author even documented the exact gap in a Javadoc 
comment, which made this much faster to verify). The one blocker is real, 
reproducible from the PR's own test suite, and directly contradicts the PR's 
own documentation of the upgrade path — it needs to be closed (or explicitly 
fenced off with a fail-fast) before this ships, since TiDB-CDC jobs already 
running in production today are exactly the population this bug would hit. Once 
Issue 1 is addressed, I'd be glad to take another pass quickly — the rest of 
the implementation is in good shape.
   > 
   > CI note: the "Build" check on the head commit (`53ce0aa`) completed with 
`action_required` (https://github.com/apache/seatunnel/runs/103907115145) — the 
workflow run needs a maintainer to approve it before it will actually execute 
for this first-time contributor; it has not produced a real pass/fail signal 
yet.
   
   Thanks for the exceptionally thorough review — the end-to-end checkout and 
the runtime-path trace made this fast to act on. All findings are addressed in 
the four new commits (756980494, ddcf29683, 58293f25a, 1b78c6f74). Where items 
overlap with @goutamadwant's inline findings I've kept the details in that 
thread and cross-reference here.
   
   Issue 1 (P0) — Fixed in 756980494; same root cause and same fix as 
@goutamadwant's inline comment above (full mechanism in my reply there: the 
ledger is reconstructed in run() from every table that had a split in flight at 
checkpoint time — the readers' restored splits, which all engines route through 
addSplitsBack() before run(), plus the state's own pending remainder — after 
which the same enumerate-the-difference logic runs against the configured set). 
To your specific asks beyond that thread:
   
   - The test that asserted the broken behavior is deleted; the corrected 
behavior is pinned by five tests covering your repro (steady-state restore with 
expansion), zero-splits-in-flight, the checkpoint pending-window, in-place 
upgrade, and legacy-config discard — module at 31/31.
   - The E2E you asked for is in ddcf29683: 
testTiDBCdcLegacyConfigRestoredWithTableNamesExpanded starts its initial job 
from a config using only database-name/table-name (table-names entirely 
absent), takes the savepoint, then restores with table-names expanded — the 
added table runs its initial snapshot. The reverse direction (table-names job 
restored with legacy keys; the removed table stops receiving changes) is 
covered by testTiDBCdcTableNamesRestoredWithLegacyConfigReduced. One honest 
limitation: a same-build e2e always writes the ledger field, so the null-ledger 
reconstruction itself is exercised at UT level via genuinely field-less states 
(legacy constructors); shipping a real pre-PR binary in CI felt out of scope 
for this PR.
   
   Docs claim (2.3) — Rather than scoping the claim down, 756980494 makes it 
true: with the reconstruction, "tables added to table-names run a fresh 
snapshot after restore" now holds for pre-PR checkpoints as well, so the Notes 
wording stands as a universal guarantee. 58293f25a additionally adds the en/zh 
callouts from Issue 3 (per-table assembly buffers → worst-case per-reader 
memory grows linearly with table count; round-robin split assignment with 
per-table isolation).
   
   Issue 2 — Fixed in 1b78c6f74: the option description now reads "Takes 
precedence over database-name/table-name when both are set", matching 
getTableFullNames() and the doc tables in both languages.
   
   Issue 4 — Covered by the test rewrite in 756980494 described above.
   
   CI note — Understood, thanks for flagging; the workflow still needs 
maintainer approval to run for a first-time contribution. The new head is 
1b78c6f7 — happy for anyone with permissions to approve the run; locally the 
connector UT suite (31/31) and the new Zeta e2e scenarios are green on my side.
   
   Whenever you have time for the second pass you kindly offered — everything 
except Issue 1's follow-ups was non-blocking, so hopefully it's a quick one.


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