dingsongjie commented on code in PR #12308:
URL: https://github.com/apache/seatunnel/pull/12308#discussion_r4023475968


##########
seatunnel-connectors-v2/connector-cdc/connector-cdc-tidb/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/tidb/source/enumerator/TiDBSourceSplitEnumerator.java:
##########
@@ -102,19 +152,38 @@ public void open() {
     @Override
     public void run() throws Exception {
         Set<Integer> readers = context.registeredReaders();
+        List<TiDBSourceSplit> sourceSplits = new ArrayList<>();
         if (shouldEnumerate) {
-            List<TiDBSourceSplit> sourceSplits = getTiDBSourceSplit();
+            sourceSplits = getTiDBSourceSplit(tableIds.keySet());
             log.info(
-                    "{} Enumerated TiDB CDC splits, database={}, table={}, 
splitCount={}.",
+                    "{} Enumerated TiDB CDC splits, tables={}, splitCount={}.",
                     CDC_DIAG_PREFIX,
-                    sourceConfig.getDatabaseName(),
-                    sourceConfig.getTableName(),
+                    tableIds.keySet(),
                     sourceSplits.size());
             synchronized (stateLock) {
+                enumeratedTables.addAll(tableIds.keySet());
                 addPendingSplit(sourceSplits);
                 shouldEnumerate = false;
-                assignSplit(readers);
             }
+        } else if (enumeratedTables != null) {

Review Comment:
   Good catch, and thanks for the precise repro — confirmed exactly as you 
described: the pre-change checkpoint deserializes enumeratedTables as null, 
shouldEnumerate is already false, so the branch was skipped forever and null 
was re-saved on every later checkpoint, silently dropping tables added at 
restore time.
   
   Fixed in 756980494, implemented as your first suggestion: the legacy state 
is initialized from the restored split tables. Specifically, run() now 
reconstructs the ledger from every table that had a split in flight at 
checkpoint time — the readers' restored splits (all engines route them through 
addSplitsBack() before run(); verified in the Zeta SourceSplitEnumeratorTask 
state machine and the Flink CoordinatedSource/ParallelSource wrappers) plus the 
checkpoint state's own pending remainder — and then runs the same 
enumerate-the-missing logic against the configured set. Your exact scenario now 
works: restoring a db.table_one checkpoint with the config expanded to 
[db.table_one, db.table_two] enumerates and snapshots table_two while table_one 
keeps streaming from its checkpoint position. An empty baseline (nothing in 
flight) correctly means every configured table counts as newly added. The 
reconstruction logs at INFO and persists with the next snapshot, so null never 
survives pas
 t the first checkpoint.
   
   The unit test that previously asserted the broken behavior is deleted; the 
new tests pin your repro (restored-splits path with table expansion), plus 
zero-in-flight, pending-window, in-place-upgrade, and legacy-config-discard 
cases — module at 31/31.



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