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]