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]