DanielLeens commented on code in PR #11271:
URL: https://github.com/apache/seatunnel/pull/11271#discussion_r3899617577


##########
seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/split/IncrementalSplit.java:
##########
@@ -131,4 +135,48 @@ public IncrementalSplit(
         this.checkpointTables = checkpointTables;
         this.historyTableChanges = historyTableChanges;
     }
+
+    /** Returns restored checkpoint state limited to the tables captured by 
the current job. */
+    public IncrementalSplit pruneTables(Collection<TableId> capturedTables) {
+        Set<TableId> capturedTableSet = new HashSet<>(capturedTables);
+        List<TableId> filteredTableIds =
+                
tableIds.stream().filter(capturedTableSet::contains).collect(Collectors.toList());
+        List<CompletedSnapshotSplitInfo> filteredCompletedSnapshotSplitInfos =
+                completedSnapshotSplitInfos.stream()
+                        .filter(info -> 
capturedTableSet.contains(info.getTableId()))
+                        .collect(Collectors.toList());
+        List<CatalogTable> filteredCheckpointTables =
+                checkpointTables == null
+                        ? null
+                        : checkpointTables.stream()
+                                .filter(
+                                        table ->
+                                                capturedTableSet.contains(

Review Comment:
   @davidzollo Sounds good, thanks for looping @hutiefang76 in directly.
   
   For the record so nothing gets lost: I already confirmed @nzw921rx's Db2 
concern is real in my reply above and I've raised it as a High-severity blocker 
in my latest full review (`pruneTables()` rebuilds a generic `TableId` straight 
from `CatalogTable#getTablePath()` instead of going through Db2's 
dialect-specific normalization, so `capturedTableSet.contains(...)` never 
matches for Db2 and every incremental-phase restore with checkpoint tables gets 
wiped). No fix has landed on the current head yet, so this is still open.
   
   @hutiefang76 the fix direction @nzw921rx sketched (a 
`JdbcDataSourceDialect#toTableId(TablePath)` default overridden by 
`Db2Dialect`, threaded into `pruneTables()`/`IncrementalSourceReader`) looks 
like the right shape to me as well - it reuses the same normalization 
`Db2Dialect` already applies elsewhere rather than adding a Db2-only special 
case inside `pruneTables()`. Happy to take another pass on this as soon as it's 
pushed.
   



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