nzw921rx commented on code in PR #11271:
URL: https://github.com/apache/seatunnel/pull/11271#discussion_r3891512246
##########
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:
Thank you for your contribution.
1. In the context of DB2, it seems that filtering is not possible here.
DB2 tableId:
```java
private static Map<TableId, CatalogTable> createDb2TableMap(
List<CatalogTable> catalogTables) {
Map<TableId, CatalogTable> tables = new HashMap<>();
CatalogTableUtils.convertTables(catalogTables)
.forEach(
(tableId, catalogTable) -> {
tables.put(tableId, catalogTable);
tables.put(toDb2TableId(tableId), catalogTable);
});
return tables;
}
private static TableId toDb2TableId(TableId tableId) {
return new TableId("", tableId.schema(), tableId.table());
}
```
Some other case codes processed by DB2:
```java
/**
* Db2 Debezium metadata always uses an empty catalog because one
connector instance captures a
* single configured database. Explicit SeaTunnel table names still
carry the database segment,
* so startup validation needs to drop that catalog part before
comparing with capture tables.
*/
static TableId toConfiguredDb2TableId(String configuredTable) {
int firstDot = configuredTable.indexOf('.');
int secondDot = configuredTable.indexOf('.', firstDot + 1);
if (firstDot < 0 || secondDot < 0 || secondDot ==
configuredTable.length() - 1) {
throw new SeaTunnelException(
"DB2 CDC table-names must use database.schema.table
format, but found: "
+ configuredTable);
}
return new TableId(
"",
configuredTable.substring(firstDot + 1, secondDot),
configuredTable.substring(secondDot + 1));
}
```
2. `JdbcDataSourceDialect` I think the logic for generating tableId should
use a unified interface here, and different databases can have personalized
processing methods
Suggested change:
```java
// JdbcDataSourceDialect
default TableId toTableId(TablePath tablePath) {
return new TableId(
tablePath.getDatabaseName(),
tablePath.getSchemaName(),
tablePath.getTableName());
}
```
```java
// Db2Dialect
@Override
public TableId toTableId(TablePath tablePath) {
return new TableId(
"",
tablePath.getSchemaName(),
tablePath.getTableName());
}
```
```java
// IncrementalSplit
public IncrementalSplit pruneTables(
Collection<TableId> capturedTables,
Function<TablePath, TableId> tableIdConverter) {
Set<TableId> capturedTableSet = new HashSet<>(capturedTables);
List<CatalogTable> filteredCheckpointTables =
checkpointTables == null
? null
: checkpointTables.stream()
.filter(
table ->
capturedTableSet.contains(
tableIdConverter.apply(
table.getTablePath())))
.collect(Collectors.toList());
// Keep the remaining filtering logic unchanged.
}
```
```java
// IncrementalSourceReader
IncrementalSplit prunedSplit =
incrementalSplit.pruneTables(
capturedTables,
dataSourceDialect::toTableId);
```
--
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]