det101 commented on code in PR #11271:
URL: https://github.com/apache/seatunnel/pull/11271#discussion_r3811132359
##########
seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java:
##########
@@ -146,7 +147,16 @@ public void addSplits(List<SourceSplitBase> splits) {
unfinishedSplits.add(split);
}
} else {
- unfinishedSplits.add(split.asIncrementalSplit());
+ IncrementalSplit incrementalSplit =
+
pruneRestoredIncrementalSplit(split.asIncrementalSplit());
+ if (incrementalSplit.getTableIds().isEmpty()) {
+ log.info(
+ "subtask {} skip restored incremental split {}
because all tables have been removed from current configuration.",
+ subtaskId,
+ incrementalSplit.splitId());
+ } else {
+ unfinishedSplits.add(incrementalSplit);
+ }
Review Comment:
When a restored incremental split is pruned to empty `tableIds`, this branch
skips the split. If that leaves `unfinishedSplits` empty, the existing `else`
below sets `needSendSplitRequest = true`, so the reader will ask the enumerator
for a new split.
On that path, `IncrementalSplitAssigner.getNext()` →
`createIncrementalSplits()` builds a `List<TableId>[incrementalParallelism]`
and always calls `createIncrementalSplit(capturedTable, …)` for every slot. If
remaining captured tables are fewer than `incrementalParallelism` (including
zero), some slots stay `null`, and `createIncrementalSplit` NPEs on
`capturedTables.contains(...)`.
Also, restore does not populate `assignedSplits` from reader-held
incremental splits, so re-creating splits here may re-assign tables that
another restored reader is already consuming.
Typical issue #11260 case (`parallelism = 1`, remove some but not all
tables) usually keeps a non-empty split and avoids this path; the risk shows up
more when incremental parallelism > remaining tables, or when all tables on one
reader’s split are removed.
Worth either skipping null partitions in `createIncrementalSplits`, or
signaling no-more-splits instead of re-requesting when prune empties the
restored split.
--
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]