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]

Reply via email to