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


##########
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:
   Thanks for the detailed catch and the concrete DB2 code — I checked this 
against the current head (`83a904f4c0`) and your concern is valid, and I'd 
actually raise the severity above where I had it tracked.
   
   Confirmed root cause: `IncrementalSplit.pruneTables()` 
(`IncrementalSplit.java:148-162`) rebuilds a Debezium `TableId` straight from 
`CatalogTable#getTablePath()`:
   
   ```java
   new TableId(
       table.getTablePath().getDatabaseName(),
       table.getTablePath().getSchemaName(),
       table.getTablePath().getTableName())
   ```
   
   For Db2 this always mismatches `capturedTables`, because 
`Db2Dialect.discoverDataCollections()` returns tables straight from 
`listOfChangeTables()`/`TableDiscoveryUtils.listTables()`, and both 
`Db2Dialect.toDb2TableId()` and `TableDiscoveryUtils.toConfiguredDb2TableId()` 
document exactly why:
   
   > "Db2 Debezium metadata always uses an empty catalog because one connector 
instance captures a single configured database. SeaTunnel catalog tables keep 
the database name, so all runtime lookups are normalized before comparing with 
Debezium metadata."
   
   So `capturedTableSet` for Db2 always carries `TableId("", schema, table)`, 
while the ad-hoc reconstruction in `pruneTables()` produces 
`TableId(realDbName, schema, table)`. `TableId.equals()` includes the catalog, 
so `capturedTableSet.contains(...)` is false for every entry — not just for 
removed tables. This is worse than the "can mismatch" framing I had for this as 
a non-blocking Medium item in my earlier reviews: for Db2 it isn't a 
partial/edge case, it deterministically wipes `checkpointTables` to an empty 
(but non-null) list on every single incremental-phase restore that carries 
checkpoint tables, since `snapshotCheckpointDataType()` populates 
`checkpointTables` unconditionally on every checkpoint 
(`IncrementalSourceReader.java:334-336`).
   
   That empty-but-non-null list then propagates further than just internal 
bookkeeping: `BaseChangeStreamTableSourceFactory.getRestoreTableStruct()` 
(`BaseChangeStreamTableSourceFactory.java:86-90`) checks 
`incrementalSplit.getCheckpointTables() != null` (not `isEmpty()`), so it 
accepts the emptied list as the authoritative restored table struct and returns 
it to `restoreSource(...)`. For Db2 CDC jobs that rely on this path, that means 
job restart can lose the checkpoint-carried table struct entirely, not just 
fail to prune a removed table correctly.
   
   I agree with your fix direction: `pruneTables()` shouldn't reconstruct a 
generic `TableId` itself. Threading a dialect-aware converter through (your 
`JdbcDataSourceDialect#toTableId(TablePath)` default + `Db2Dialect` override, 
passed into `pruneTables`/used from `IncrementalSourceReader`) is the right 
shape, since it reuses the same normalization Db2Dialect already applies 
everywhere else (`checkAllTablesEnabledCapture`, `queryTableSchema`, 
`getPrimaryKey`, `getConstraintKeys`).
   
   @hutiefang76 given the above, I'd treat this as a High-severity blocker for 
Db2 specifically (not the pre-existing non-blocking Medium item), since it 
affects every restore rather than only the removed-table case. Worth fixing 
before merge — thanks again @nzw921rx for pinning down the concrete failure 
mode with real Db2 code.



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