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


##########
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:
   Resolved on the current head (merge commit ee2a1424199, and re-verified 
after my push at 4bd75ee8128).
   
   You're right that the generic `new TableId(database, schema, table)` 
reconstruction from `TablePath` doesn't work for Db2, since Db2's Debezium 
identifiers always use an empty catalog (`new TableId("", schema, table)`).
   
   This was already fixed exactly along the lines you suggested, in commit 
7df41cf5d9021a9d00ec05e4f7dda8321ebd80ea:
   - `DataSourceDialect.toTableId(TablePath)` — new default method on the 
dialect interface, generic `(database, schema, table)` mapping.
   - `Db2Dialect.toTableId(TablePath)` — overridden to return `new TableId("", 
schema, table)`, matching Db2's empty-catalog format.
   - `IncrementalSplit.pruneTables(Collection<TableId>, Function<TablePath, 
TableId> tableIdConverter)` — now takes the converter as a parameter instead of 
hardcoding the generic constructor.
   - `IncrementalSourceReader` calls 
`incrementalSplit.pruneTables(capturedTables, dataSourceDialect::toTableId)`, 
so the conversion always goes through the dialect.
   - `IncrementalSplitTest#testPruneTablesUsesDialectSpecificTableIdConverter` 
and `Db2IncrementalSourceFactoryTest` cover the Db2 empty-catalog case 
explicitly.
   
   Dismissing the review as resolved. Thanks for the detailed repro.



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