This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-11271-04fabece50f757e8fbe9967d071dbcf0cba8a414 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit af0a647d27bcc97b1e185d9a1292729622973df4 Author: hutiefang76 <[email protected]> AuthorDate: Sat Sep 5 15:50:45 2026 +0000 [Bug][Connector-V2][CDC] Prune removed tables from restored incremental splits (#11271) Signed-off-by: David Zollo <[email protected]> Co-authored-by: hutiefang <[email protected]> Co-authored-by: DanielLeens <[email protected]> Co-authored-by: davidzollo <[email protected]> Co-authored-by: David Zollo <[email protected]> Co-authored-by: zhiweiniu <[email protected]> Co-authored-by: Claude Sonnet 5 <[email protected]> --- .../introduction/concepts/incompatible-changes.md | 5 + .../introduction/concepts/incompatible-changes.md | 5 + .../cdc/base/dialect/DataSourceDialect.java | 15 ++ .../enumerator/IncrementalSplitAssigner.java | 6 + .../source/reader/IncrementalSourceReader.java | 65 +++++- .../cdc/base/source/split/IncrementalSplit.java | 66 ++++++ .../enumerator/IncrementalSplitAssignerTest.java | 90 ++++++++- .../source/reader/IncrementalSourceReaderTest.java | 225 +++++++++++++++++++++ .../base/source/split/IncrementalSplitTest.java | 189 +++++++++++++++++ .../seatunnel/cdc/db2/source/Db2Dialect.java | 12 ++ .../source/Db2IncrementalSourceFactoryTest.java | 21 ++ 11 files changed, 697 insertions(+), 2 deletions(-) diff --git a/docs/en/introduction/concepts/incompatible-changes.md b/docs/en/introduction/concepts/incompatible-changes.md index 412d577756..50561c2ebb 100644 --- a/docs/en/introduction/concepts/incompatible-changes.md +++ b/docs/en/introduction/concepts/incompatible-changes.md @@ -131,6 +131,11 @@ You need to check this document before you upgrade to related version. - **Description**: The enumerator now partitions a table (or the configured `start_rowkey` / `end_rowkey` range) into tablet-sized splits via `sampleRowKeys`. `scan_row_limit` is still applied with `query.limit(...)` once per split in the reader. Before this change the source always produced exactly one split, so `scan_row_limit` acted as a table-wide row cap. After this change a table with multiple tablets yields multiple splits even when `parallelism = 1` (the single reader is assign [...] - **Impact**: Existing jobs that set `scan_row_limit` to bound total output (sampling, testing, cost control, or downstream capacity) can read far more rows after upgrade with no config change. - **Migration Guide**: If you need a table-wide cap, narrow the scan with `start_rowkey` / `end_rowkey`, or lower `scan_row_limit` so that `scan_row_limit × expected split count` stays within the previous budget. To keep the previous single-split behavior, the connector still falls back to one split when sampling fails, returns no keys, or the intersection is empty — that is not a supported way to pin the old cap. (#11876) +- **CDC Connector: restored state for tables removed from the capture set is no longer reused** + - **Affected component**: `seatunnel-connectors-v2/connector-cdc/connector-cdc-base` and CDC connectors built on it. + - **Description**: When a CDC job restores from a checkpoint or savepoint, SeaTunnel now filters per-table incremental state against the currently captured table set before assigning the restored split. State for tables that have been removed from the job's capture configuration is not reused. If table discovery is unavailable or returns no tables, SeaTunnel keeps the restored state unchanged to avoid discarding checkpoint metadata during a transient source-database problem. + - **Impact**: A job that removes captured tables and then restores from an older checkpoint no longer attempts to resume incremental state for those removed tables. This avoids restore failures caused by stale table metadata. The behavior applies only during checkpoint/savepoint restore; newly started jobs are unchanged. + - **Migration Guide**: No configuration change is required. Before restoring an existing CDC job after changing its capture table set, verify that the removed tables are intentionally no longer part of the job. - **Breaking Change: Iceberg Connector — source table primary key is no longer silently inherited** - **Affected component**: `seatunnel-connectors-v2/connector-iceberg` diff --git a/docs/zh/introduction/concepts/incompatible-changes.md b/docs/zh/introduction/concepts/incompatible-changes.md index fc9e804ca1..887647be93 100644 --- a/docs/zh/introduction/concepts/incompatible-changes.md +++ b/docs/zh/introduction/concepts/incompatible-changes.md @@ -119,6 +119,11 @@ - **变更说明**:Enumerator 现在通过 `sampleRowKeys` 按 tablet 边界把表(或配置的 `start_rowkey` / `end_rowkey` 区间)切成多个 split。Reader 仍对每个 split 调用一次 `query.limit(...)`。此前 Source 始终只产生 1 个 split,因此 `scan_row_limit` 等价于整表行数上限。升级后,只要表有多个 tablet,即使 `parallelism = 1`(唯一 reader 会拿到全部 split),作业级上限约为 `scan_row_limit × split 数`。详见 [Google Bigtable Source](../../connectors/source/GoogleBigtable.md#scan_row_limit-int)。 - **影响**:依赖 `scan_row_limit` 限制总输出量的存量作业(抽样、测试、成本控制、下游容量)在升级后、配置不变的情况下,可能读出远超以前的行数。 - **迁移指南**:若仍需要整表级上限,请用 `start_rowkey` / `end_rowkey` 收窄扫描范围,或下调 `scan_row_limit`,使 `scan_row_limit × 预期 split 数` 不超过原预算。采样失败、无采样点或求交为空时仍会回退为单个 split,但这不是用来锁定旧语义的受支持方式。(#11876) +- **CDC Connector:已从捕获集合移除的表不再复用恢复状态** + - **影响范围**:`seatunnel-connectors-v2/connector-cdc/connector-cdc-base` 及其构建的 CDC 连接器。 + - **变更说明**:CDC 任务从 checkpoint 或 savepoint 恢复时,SeaTunnel 现在会在分配恢复后的 split 前,按照当前捕获表集合过滤表级增量状态。已从任务捕获配置中移除的表,其状态不会再被复用;如果表发现不可用或返回空集合,为避免源数据库短暂异常时丢弃 checkpoint 元数据,SeaTunnel 会保持恢复状态不变。 + - **影响**:任务移除捕获表后再从旧 checkpoint 恢复时,不再尝试恢复这些已移除表的增量状态,从而避免陈旧表元数据导致恢复失败。该行为仅作用于 checkpoint/savepoint 恢复;新启动的任务不受影响。 + - **迁移指南**:无需修改配置。变更捕获表集合后恢复现有 CDC 任务前,请确认被移除的表确实不应继续参与该任务。 - **破坏性变更:Iceberg 连接器 — 不再自动继承源表主键** - **影响范围**:`seatunnel-connectors-v2/connector-iceberg` diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/dialect/DataSourceDialect.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/dialect/DataSourceDialect.java index 58f59f4c61..d69769bc1d 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/dialect/DataSourceDialect.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/dialect/DataSourceDialect.java @@ -17,6 +17,7 @@ package org.apache.seatunnel.connectors.cdc.base.dialect; +import org.apache.seatunnel.api.table.catalog.TablePath; import org.apache.seatunnel.connectors.cdc.base.config.SourceConfig; import org.apache.seatunnel.connectors.cdc.base.source.enumerator.splitter.ChunkSplitter; import org.apache.seatunnel.connectors.cdc.base.source.offset.Offset; @@ -41,6 +42,20 @@ public interface DataSourceDialect<C extends SourceConfig> extends Serializable /** Discovers the list of data collection to capture. */ List<TableId> discoverDataCollections(C sourceConfig); + /** + * Converts a checkpoint table path to the dialect-specific Debezium identifier. + * + * <p>Dialects whose discovered table identifiers use a different namespace must override this + * method so restored checkpoint table structures are compared in the same identifier format. + * + * @param tablePath table path stored in checkpoint schema state + * @return identifier in the same format used by discovered table identifiers + */ + default TableId toTableId(TablePath tablePath) { + return new TableId( + tablePath.getDatabaseName(), tablePath.getSchemaName(), tablePath.getTableName()); + } + /** Check if the CollectionId is case-sensitive or not. */ boolean isDataCollectionIdCaseSensitive(C sourceConfig); diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/enumerator/IncrementalSplitAssigner.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/enumerator/IncrementalSplitAssigner.java index e47e08ecf8..a6188e9183 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/enumerator/IncrementalSplitAssigner.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/enumerator/IncrementalSplitAssigner.java @@ -226,6 +226,12 @@ public class IncrementalSplitAssigner<C extends SourceConfig> implements SplitAs i = 0; List<IncrementalSplit> incrementalSplits = new ArrayList<>(); for (List<TableId> capturedTable : capturedTables) { + // Restored split pruning can shrink the captured-table set below the configured + // parallelism, so some round-robin buckets may legitimately stay empty. + if (capturedTable == null || capturedTable.isEmpty()) { + i++; + continue; + } incrementalSplits.add( createIncrementalSplit(capturedTable, i++, startWithSnapshotMinimumOffset)); } diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java index a8899d3065..e4cd1a1da8 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReader.java @@ -40,6 +40,7 @@ import org.apache.seatunnel.connectors.seatunnel.common.source.reader.SingleThre import org.apache.seatunnel.connectors.seatunnel.common.source.reader.SourceReaderOptions; import org.apache.seatunnel.connectors.seatunnel.common.source.reader.fetcher.SingleThreadFetcherManager; +import io.debezium.relational.TableId; import lombok.extern.slf4j.Slf4j; import java.util.ArrayList; @@ -129,6 +130,8 @@ public class IncrementalSourceReader<T, C extends SourceConfig> public void addSplits(List<SourceSplitBase> splits) { // restore for finishedUnackedSplits List<SourceSplitBase> unfinishedSplits = new ArrayList<>(); + List<TableId> capturedTables = null; + boolean capturedTablesDiscovered = false; log.info( "subtask {} add splits: {}", subtaskId, @@ -146,7 +149,23 @@ public class IncrementalSourceReader<T, C extends SourceConfig> unfinishedSplits.add(split); } } else { - unfinishedSplits.add(split.asIncrementalSplit()); + IncrementalSplit incrementalSplit = split.asIncrementalSplit(); + if (hasRestoredCheckpointMetadata(incrementalSplit)) { + if (!capturedTablesDiscovered) { + capturedTables = discoverCapturedTables(); + capturedTablesDiscovered = true; + } + incrementalSplit = + pruneRestoredIncrementalSplit(incrementalSplit, capturedTables); + } + if (incrementalSplit.getTableIds().isEmpty()) { + log.info( + "subtask {} skip restored incremental split {} because all tables have been removed from current configuration.", + subtaskId, + incrementalSplit.splitId()); + } else { + unfinishedSplits.add(incrementalSplit); + } } } // notify split enumerator again about the finished unacked snapshot splits @@ -238,6 +257,50 @@ public class IncrementalSourceReader<T, C extends SourceConfig> } } + private List<TableId> discoverCapturedTables() { + try { + return dataSourceDialect.discoverDataCollections(sourceConfig); + } catch (Exception e) { + log.warn( + "Failed to discover captured tables while restoring CDC split. " + + "Keeping restored checkpoint state unchanged.", + e); + return null; + } + } + + private IncrementalSplit pruneRestoredIncrementalSplit( + IncrementalSplit incrementalSplit, List<TableId> capturedTables) { + if (capturedTables == null) { + return incrementalSplit; + } + if (capturedTables.isEmpty() && !incrementalSplit.getTableIds().isEmpty()) { + log.warn( + "Skip pruning restored incremental split {} because captured table discovery returned " + + "an empty result. Keeping restored checkpoint state unchanged.", + incrementalSplit.splitId()); + return incrementalSplit; + } + IncrementalSplit prunedSplit = + incrementalSplit.pruneTables(capturedTables, dataSourceDialect::toTableId); + if (prunedSplit.getTableIds().size() != incrementalSplit.getTableIds().size()) { + log.info( + "Pruned restored incremental split {} tables from {} to {} based on current captured tables.", + incrementalSplit.splitId(), + incrementalSplit.getTableIds(), + prunedSplit.getTableIds()); + } + return prunedSplit; + } + + private boolean hasRestoredCheckpointMetadata(IncrementalSplit incrementalSplit) { + return incrementalSplit.getCheckpointDataType() != null + || (incrementalSplit.getCheckpointTables() != null + && !incrementalSplit.getCheckpointTables().isEmpty()) + || (incrementalSplit.getHistoryTableChanges() != null + && !incrementalSplit.getHistoryTableChanges().isEmpty()); + } + @Override public List<SourceSplitBase> snapshotState(long checkpointId) { List<SourceSplitBase> stateSplits = super.snapshotState(checkpointId); diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/split/IncrementalSplit.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/split/IncrementalSplit.java index e9052da461..39cb117c68 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/split/IncrementalSplit.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/main/java/org/apache/seatunnel/connectors/cdc/base/source/split/IncrementalSplit.java @@ -18,6 +18,7 @@ package org.apache.seatunnel.connectors.cdc.base.source.split; import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.catalog.TablePath; import org.apache.seatunnel.api.table.type.SeaTunnelDataType; import org.apache.seatunnel.connectors.cdc.base.source.offset.Offset; @@ -26,9 +27,14 @@ import lombok.Getter; import lombok.ToString; import java.util.ArrayList; +import java.util.Collection; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; +import java.util.Set; +import java.util.function.Function; +import java.util.stream.Collectors; @ToString @Getter @@ -131,4 +137,64 @@ public class IncrementalSplit extends SourceSplitBase { this.checkpointTables = checkpointTables; this.historyTableChanges = historyTableChanges; } + + /** + * Returns restored checkpoint state limited to the tables captured by the current job. + * + * <p>The checkpoint schema stores {@link TablePath}s while split state uses Debezium {@link + * TableId}s. The caller supplies the dialect-specific conversion so both forms use the same + * namespace during restore. + * + * @param capturedTables table identifiers discovered for the current job configuration + * @param tableIdConverter converts checkpoint table paths to discovered table identifiers + * @return a copy of this split without state for tables no longer captured + */ + public IncrementalSplit pruneTables( + Collection<TableId> capturedTables, Function<TablePath, TableId> tableIdConverter) { + Set<TableId> capturedTableSet = new HashSet<>(capturedTables); + // Guard tableIds/completedSnapshotSplitInfos the same way checkpointTables and + // historyTableChanges are guarded below: the 7-arg constructor accepts null for every + // field here, so a restored split with a null list must be pruned to an empty list + // instead of throwing an NPE during checkpoint recovery. + List<TableId> filteredTableIds = + tableIds == null + ? new ArrayList<>() + : tableIds.stream() + .filter(capturedTableSet::contains) + .collect(Collectors.toList()); + List<CompletedSnapshotSplitInfo> filteredCompletedSnapshotSplitInfos = + completedSnapshotSplitInfos == null + ? new ArrayList<>() + : completedSnapshotSplitInfos.stream() + .filter(info -> capturedTableSet.contains(info.getTableId())) + .collect(Collectors.toList()); + List<CatalogTable> filteredCheckpointTables = + checkpointTables == null + ? null + : checkpointTables.stream() + .filter( + table -> + capturedTableSet.contains( + tableIdConverter.apply( + table.getTablePath()))) + .collect(Collectors.toList()); + Map<TableId, byte[]> filteredHistoryTableChanges = + historyTableChanges == null + ? null + : historyTableChanges.entrySet().stream() + .filter(entry -> capturedTableSet.contains(entry.getKey())) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + IncrementalSplit prunedSplit = + new IncrementalSplit( + splitId(), + filteredTableIds, + startupOffset, + stopOffset, + filteredCompletedSnapshotSplitInfos, + filteredCheckpointTables, + filteredHistoryTableChanges); + // Keep compatibility with checkpoints created before table-level schema history. + prunedSplit.checkpointDataType = checkpointDataType; + return prunedSplit; + } } diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/enumerator/IncrementalSplitAssignerTest.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/enumerator/IncrementalSplitAssignerTest.java index 0229f3b663..c96248e498 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/enumerator/IncrementalSplitAssignerTest.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/enumerator/IncrementalSplitAssignerTest.java @@ -25,6 +25,7 @@ import org.apache.seatunnel.connectors.cdc.base.option.StopMode; import org.apache.seatunnel.connectors.cdc.base.source.enumerator.state.IncrementalPhaseState; import org.apache.seatunnel.connectors.cdc.base.source.offset.Offset; import org.apache.seatunnel.connectors.cdc.base.source.offset.OffsetFactory; +import org.apache.seatunnel.connectors.cdc.base.source.split.IncrementalSplit; import org.junit.jupiter.api.Test; @@ -32,8 +33,14 @@ import io.debezium.relational.TableId; import java.util.Collections; import java.util.HashMap; +import java.util.LinkedHashSet; +import java.util.Optional; +import java.util.Set; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @@ -70,10 +77,91 @@ class IncrementalSplitAssignerTest { verify(offsetFactory).committedOffset(); } + @Test + void shouldSkipEmptyBucketsWhenCapturedTablesAreLessThanParallelism() { + SourceConfig sourceConfig = mock(SourceConfig.class); + OffsetFactory offsetFactory = mock(OffsetFactory.class); + when(sourceConfig.getStartupConfig()) + .thenReturn(new StartupConfig(StartupMode.INITIAL, null, null, null)); + when(sourceConfig.getStopConfig()) + .thenReturn(new StopConfig(StopMode.NEVER, null, null, null)); + when(offsetFactory.neverStop()).thenReturn(mock(Offset.class)); + + IncrementalSplitAssigner<SourceConfig> assigner = + new IncrementalSplitAssigner<>( + createContext( + sourceConfig, + Collections.singleton(TableId.parse("database.schema.table"))), + 2, + offsetFactory); + + Optional<IncrementalSplit> firstSplit = + assigner.getNext().map(split -> split.asIncrementalSplit()); + + assertTrue(firstSplit.isPresent()); + assertEquals(1, firstSplit.get().getTableIds().size()); + assertEquals("incremental-split-0", firstSplit.get().splitId()); + assertFalse(assigner.getNext().isPresent()); + } + + @Test + void shouldReturnEmptyWhenNoCapturedTablesRemain() { + SourceConfig sourceConfig = mock(SourceConfig.class); + OffsetFactory offsetFactory = mock(OffsetFactory.class); + when(sourceConfig.getStartupConfig()) + .thenReturn(new StartupConfig(StartupMode.INITIAL, null, null, null)); + when(sourceConfig.getStopConfig()) + .thenReturn(new StopConfig(StopMode.NEVER, null, null, null)); + + IncrementalSplitAssigner<SourceConfig> assigner = + new IncrementalSplitAssigner<>( + createContext(sourceConfig, Collections.emptySet()), 2, offsetFactory); + + assertFalse(assigner.getNext().isPresent()); + assertTrue(assigner.noMoreSplits()); + } + + @Test + void shouldReassignTablesRestoredFromCheckpointWithTheirWatermark() { + SourceConfig sourceConfig = mock(SourceConfig.class); + OffsetFactory offsetFactory = mock(OffsetFactory.class); + Offset startupOffset = mock(Offset.class); + Offset stoppingOffset = mock(Offset.class); + when(sourceConfig.getStartupConfig()) + .thenReturn(new StartupConfig(StartupMode.INITIAL, null, null, null)); + when(sourceConfig.getStopConfig()) + .thenReturn(new StopConfig(StopMode.NEVER, null, null, null)); + + TableId tableId = TableId.parse("database.schema.table"); + IncrementalSplit restoredSplit = + new IncrementalSplit( + "incremental-split-0", + Collections.singletonList(tableId), + startupOffset, + stoppingOffset, + Collections.emptyList()); + IncrementalSplitAssigner<SourceConfig> assigner = + new IncrementalSplitAssigner<>(createContext(sourceConfig), 1, offsetFactory); + assigner.addSplits(Collections.singletonList(restoredSplit)); + + Optional<IncrementalSplit> reassignedSplit = + assigner.getNext().map(split -> split.asIncrementalSplit()); + + assertTrue(reassignedSplit.isPresent()); + assertEquals(Collections.singletonList(tableId), reassignedSplit.get().getTableIds()); + assertSame(startupOffset, reassignedSplit.get().getStartupOffset()); + } + private SplitAssigner.Context<SourceConfig> createContext(SourceConfig sourceConfig) { + return createContext( + sourceConfig, Collections.singleton(TableId.parse("database.schema.table"))); + } + + private SplitAssigner.Context<SourceConfig> createContext( + SourceConfig sourceConfig, Set<TableId> capturedTables) { return new SplitAssigner.Context<>( sourceConfig, - Collections.singleton(TableId.parse("database.schema.table")), + new LinkedHashSet<>(capturedTables), new HashMap<>(), new HashMap<>()); } diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReaderTest.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReaderTest.java new file mode 100644 index 0000000000..1f0c85cdcd --- /dev/null +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/reader/IncrementalSourceReaderTest.java @@ -0,0 +1,225 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.connectors.cdc.base.source.reader; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.source.SourceReader; +import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.catalog.TableIdentifier; +import org.apache.seatunnel.api.table.catalog.TablePath; +import org.apache.seatunnel.api.table.catalog.TableSchema; +import org.apache.seatunnel.connectors.cdc.base.config.SourceConfig; +import org.apache.seatunnel.connectors.cdc.base.dialect.DataSourceDialect; +import org.apache.seatunnel.connectors.cdc.base.source.split.IncrementalSplit; +import org.apache.seatunnel.connectors.cdc.base.source.split.SourceRecords; +import org.apache.seatunnel.connectors.cdc.base.source.split.SourceSplitBase; +import org.apache.seatunnel.connectors.cdc.base.source.split.state.SourceSplitStateBase; +import org.apache.seatunnel.connectors.cdc.debezium.DebeziumDeserializationSchema; +import org.apache.seatunnel.connectors.seatunnel.common.source.reader.RecordEmitter; +import org.apache.seatunnel.connectors.seatunnel.common.source.reader.RecordsWithSplitIds; +import org.apache.seatunnel.connectors.seatunnel.common.source.reader.SourceReaderOptions; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import io.debezium.relational.TableId; + +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.CountDownLatch; + +class IncrementalSourceReaderTest { + + private static final TableId KEPT_TABLE = + new TableId("alpha_online", null, "account_histories"); + private static final TableId REMOVED_TABLE = + new TableId("alpha_online", null, "account_interests"); + + @Test + void testAddSplitsEnqueuesPrunedRestoredIncrementalSplit() { + SourceConfig sourceConfig = Mockito.mock(SourceConfig.class); + DataSourceDialect<SourceConfig> dialect = + Mockito.mock(DataSourceDialect.class, Mockito.CALLS_REAL_METHODS); + Mockito.when(dialect.discoverDataCollections(sourceConfig)) + .thenReturn(Collections.singletonList(KEPT_TABLE)); + SourceReader.Context context = Mockito.mock(SourceReader.Context.class); + IncrementalSourceReader<Object, SourceConfig> reader = + createReader(dialect, sourceConfig, context); + + try { + reader.addSplits(Collections.singletonList(restoredIncrementalSplit())); + + Assertions.assertEquals(1, reader.getNumberOfCurrentlyAssignedSplits()); + List<SourceSplitBase> state = reader.snapshotState(1L); + Assertions.assertEquals(1, state.size()); + Assertions.assertEquals( + Collections.singletonList(KEPT_TABLE), + state.get(0).asIncrementalSplit().getTableIds()); + } finally { + reader.close(); + } + } + + @Test + void testAddSplitsKeepsRestoredSplitWhenDiscoveryReturnsEmpty() { + SourceConfig sourceConfig = Mockito.mock(SourceConfig.class); + DataSourceDialect<SourceConfig> dialect = + Mockito.mock(DataSourceDialect.class, Mockito.CALLS_REAL_METHODS); + Mockito.when(dialect.discoverDataCollections(sourceConfig)) + .thenReturn(Collections.emptyList()); + SourceReader.Context context = Mockito.mock(SourceReader.Context.class); + IncrementalSourceReader<Object, SourceConfig> reader = + createReader(dialect, sourceConfig, context); + + try { + reader.addSplits(Collections.singletonList(restoredIncrementalSplit())); + + Assertions.assertEquals(1, reader.getNumberOfCurrentlyAssignedSplits()); + List<SourceSplitBase> state = reader.snapshotState(1L); + Assertions.assertEquals( + Arrays.asList(KEPT_TABLE, REMOVED_TABLE), + state.get(0).asIncrementalSplit().getTableIds()); + } finally { + reader.close(); + } + } + + @Test + void testAddSplitsKeepsRestoredSplitWhenDiscoveryFails() { + SourceConfig sourceConfig = Mockito.mock(SourceConfig.class); + DataSourceDialect<SourceConfig> dialect = + Mockito.mock(DataSourceDialect.class, Mockito.CALLS_REAL_METHODS); + Mockito.when(dialect.discoverDataCollections(sourceConfig)) + .thenThrow(new RuntimeException("database unavailable")); + SourceReader.Context context = Mockito.mock(SourceReader.Context.class); + IncrementalSourceReader<Object, SourceConfig> reader = + createReader(dialect, sourceConfig, context); + + try { + reader.addSplits(Collections.singletonList(restoredIncrementalSplit())); + + Assertions.assertEquals(1, reader.getNumberOfCurrentlyAssignedSplits()); + List<SourceSplitBase> state = reader.snapshotState(1L); + Assertions.assertEquals( + Arrays.asList(KEPT_TABLE, REMOVED_TABLE), + state.get(0).asIncrementalSplit().getTableIds()); + } finally { + reader.close(); + } + } + + @Test + void testAddSplitsDiscoversCapturedTablesOnlyOncePerBatch() { + SourceConfig sourceConfig = Mockito.mock(SourceConfig.class); + DataSourceDialect<SourceConfig> dialect = + Mockito.mock(DataSourceDialect.class, Mockito.CALLS_REAL_METHODS); + Mockito.when(dialect.discoverDataCollections(sourceConfig)) + .thenReturn(Collections.singletonList(KEPT_TABLE)); + SourceReader.Context context = Mockito.mock(SourceReader.Context.class); + IncrementalSourceReader<Object, SourceConfig> reader = + createReader(dialect, sourceConfig, context); + + try { + reader.addSplits( + Arrays.asList( + restoredIncrementalSplit(), + restoredIncrementalSplit("incremental-split-1"))); + + Mockito.verify(dialect, Mockito.times(1)).discoverDataCollections(sourceConfig); + } finally { + reader.close(); + } + } + + private static IncrementalSourceReader<Object, SourceConfig> createReader( + DataSourceDialect<SourceConfig> dialect, + SourceConfig sourceConfig, + SourceReader.Context context) { + @SuppressWarnings("unchecked") + IncrementalSourceSplitReader<SourceConfig> splitReader = + Mockito.mock(IncrementalSourceSplitReader.class); + CountDownLatch wakeUp = new CountDownLatch(1); + try { + Mockito.when(splitReader.fetch()) + .thenAnswer( + invocation -> { + wakeUp.await(); + return null; + }); + } catch (Exception e) { + throw new RuntimeException(e); + } + Mockito.doAnswer( + invocation -> { + wakeUp.countDown(); + return null; + }) + .when(splitReader) + .wakeUp(); + + @SuppressWarnings("unchecked") + RecordEmitter<SourceRecords, Object, SourceSplitStateBase> recordEmitter = + Mockito.mock(RecordEmitter.class); + @SuppressWarnings("unchecked") + DebeziumDeserializationSchema<Object> deserializationSchema = + Mockito.mock(DebeziumDeserializationSchema.class); + + return new IncrementalSourceReader<>( + dialect, + new ArrayBlockingQueue<RecordsWithSplitIds<SourceRecords>>(2), + () -> splitReader, + recordEmitter, + new SourceReaderOptions(ReadonlyConfig.fromMap(Collections.emptyMap())), + context, + sourceConfig, + deserializationSchema); + } + + private static IncrementalSplit restoredIncrementalSplit() { + return restoredIncrementalSplit("incremental-split-0"); + } + + private static IncrementalSplit restoredIncrementalSplit(String splitId) { + Map<TableId, byte[]> historyTableChanges = new HashMap<>(); + historyTableChanges.put(KEPT_TABLE, new byte[] {1}); + historyTableChanges.put(REMOVED_TABLE, new byte[] {2}); + return new IncrementalSplit( + splitId, + Arrays.asList(KEPT_TABLE, REMOVED_TABLE), + null, + null, + Collections.emptyList(), + Arrays.asList(catalogTable(KEPT_TABLE), catalogTable(REMOVED_TABLE)), + historyTableChanges); + } + + private static CatalogTable catalogTable(TableId tableId) { + TablePath tablePath = TablePath.of(tableId.catalog(), tableId.table()); + return CatalogTable.of( + TableIdentifier.of("test", tablePath), + TableSchema.builder().build(), + Collections.emptyMap(), + Collections.emptyList(), + ""); + } +} diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/split/IncrementalSplitTest.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/split/IncrementalSplitTest.java new file mode 100644 index 0000000000..97160c6c55 --- /dev/null +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-base/src/test/java/org/apache/seatunnel/connectors/cdc/base/source/split/IncrementalSplitTest.java @@ -0,0 +1,189 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.connectors.cdc.base.source.split; + +import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.catalog.TableIdentifier; +import org.apache.seatunnel.api.table.catalog.TablePath; +import org.apache.seatunnel.api.table.catalog.TableSchema; +import org.apache.seatunnel.api.table.type.BasicType; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +import io.debezium.relational.TableId; + +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.function.Function; + +/** Verifies restored incremental split state is retained only for currently captured tables. */ +public class IncrementalSplitTest { + + private static final TableId KEPT_TABLE = + new TableId("alpha_online", null, "account_histories"); + private static final TableId REMOVED_TABLE = + new TableId("alpha_online", null, "account_interests"); + + /** Converts regular checkpoint table paths to the default Debezium table identifier format. */ + private static final Function<TablePath, TableId> DEFAULT_TABLE_ID_CONVERTER = + tablePath -> + new TableId( + tablePath.getDatabaseName(), + tablePath.getSchemaName(), + tablePath.getTableName()); + + @Test + public void testPruneTablesRemovesDeletedTableState() { + Map<TableId, byte[]> historyTableChanges = new HashMap<>(); + historyTableChanges.put(KEPT_TABLE, new byte[] {1}); + historyTableChanges.put(REMOVED_TABLE, new byte[] {2}); + + IncrementalSplit split = + new IncrementalSplit( + "incremental-split-0", + Arrays.asList(KEPT_TABLE, REMOVED_TABLE), + null, + null, + Arrays.asList( + completedSnapshotSplitInfo("kept-split", KEPT_TABLE), + completedSnapshotSplitInfo("removed-split", REMOVED_TABLE)), + Arrays.asList(catalogTable(KEPT_TABLE), catalogTable(REMOVED_TABLE)), + historyTableChanges); + + IncrementalSplit pruned = + split.pruneTables( + Collections.singletonList(KEPT_TABLE), DEFAULT_TABLE_ID_CONVERTER); + + Assertions.assertEquals(Collections.singletonList(KEPT_TABLE), pruned.getTableIds()); + Assertions.assertEquals(1, pruned.getCompletedSnapshotSplitInfos().size()); + Assertions.assertEquals( + KEPT_TABLE, pruned.getCompletedSnapshotSplitInfos().get(0).getTableId()); + Assertions.assertEquals(1, pruned.getCheckpointTables().size()); + Assertions.assertEquals( + TablePath.of("alpha_online.account_histories"), + pruned.getCheckpointTables().get(0).getTablePath()); + Assertions.assertEquals( + Collections.singleton(KEPT_TABLE), pruned.getHistoryTableChanges().keySet()); + } + + @Test + public void testPruneTablesPreservesLegacyCheckpointDataTypeWhenTableRemoved() { + IncrementalSplit split = legacyCheckpointSplit(); + + IncrementalSplit pruned = + split.pruneTables( + Collections.singletonList(KEPT_TABLE), DEFAULT_TABLE_ID_CONVERTER); + + Assertions.assertEquals(Collections.singletonList(KEPT_TABLE), pruned.getTableIds()); + Assertions.assertSame(BasicType.STRING_TYPE, pruned.getCheckpointDataType()); + } + + @Test + public void testPruneTablesPreservesLegacyCheckpointDataTypeWhenTablesUnchanged() { + IncrementalSplit split = legacyCheckpointSplit(); + + IncrementalSplit pruned = + split.pruneTables( + Arrays.asList(KEPT_TABLE, REMOVED_TABLE), DEFAULT_TABLE_ID_CONVERTER); + + Assertions.assertEquals(Arrays.asList(KEPT_TABLE, REMOVED_TABLE), pruned.getTableIds()); + Assertions.assertSame(BasicType.STRING_TYPE, pruned.getCheckpointDataType()); + } + + /** Verifies checkpoint schemas use the dialect converter instead of a generic table id. */ + @Test + public void testPruneTablesUsesDialectSpecificTableIdConverter() { + TableId db2TableId = new TableId("", "DB2INST1", "CUSTOMERS"); + IncrementalSplit split = + new IncrementalSplit( + "incremental-split-db2", + Collections.singletonList(db2TableId), + null, + null, + Collections.emptyList(), + Collections.singletonList( + catalogTable(TablePath.of("SAMPLE", "DB2INST1", "CUSTOMERS"))), + Collections.emptyMap()); + + IncrementalSplit pruned = + split.pruneTables( + Collections.singletonList(db2TableId), + tablePath -> + new TableId( + "", tablePath.getSchemaName(), tablePath.getTableName())); + + Assertions.assertEquals(Collections.singletonList(db2TableId), pruned.getTableIds()); + Assertions.assertEquals(1, pruned.getCheckpointTables().size()); + Assertions.assertEquals( + TablePath.of("SAMPLE", "DB2INST1", "CUSTOMERS"), + pruned.getCheckpointTables().get(0).getTablePath()); + } + + /** + * Verifies pruneTables tolerates a null tableIds/completedSnapshotSplitInfos the same way it + * already tolerates null checkpointTables/historyTableChanges, instead of throwing an NPE + * during checkpoint recovery. + */ + @Test + public void testPruneTablesToleratesNullTableIdsAndCompletedSnapshotSplitInfos() { + IncrementalSplit split = + new IncrementalSplit( + "incremental-split-null-state", null, null, null, null, null, null); + + IncrementalSplit pruned = + split.pruneTables( + Collections.singletonList(KEPT_TABLE), DEFAULT_TABLE_ID_CONVERTER); + + Assertions.assertTrue(pruned.getTableIds().isEmpty()); + Assertions.assertTrue(pruned.getCompletedSnapshotSplitInfos().isEmpty()); + } + + @SuppressWarnings("deprecation") + private static IncrementalSplit legacyCheckpointSplit() { + IncrementalSplit split = + new IncrementalSplit( + "incremental-split-legacy", + Arrays.asList(KEPT_TABLE, REMOVED_TABLE), + null, + null, + Collections.emptyList()); + return new IncrementalSplit(split, BasicType.STRING_TYPE); + } + + private static CompletedSnapshotSplitInfo completedSnapshotSplitInfo( + String splitId, TableId tableId) { + return new CompletedSnapshotSplitInfo(splitId, tableId, null, null, null, null); + } + + private static CatalogTable catalogTable(TableId tableId) { + return catalogTable(TablePath.of(tableId.catalog(), tableId.table())); + } + + /** Creates a catalog table whose identifier is stored in checkpoint schema state. */ + private static CatalogTable catalogTable(TablePath tablePath) { + return CatalogTable.of( + TableIdentifier.of("test", tablePath), + TableSchema.builder().build(), + Collections.emptyMap(), + Collections.emptyList(), + ""); + } +} diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-db2/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/db2/source/Db2Dialect.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-db2/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/db2/source/Db2Dialect.java index ffaefb3f3b..107feac297 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-db2/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/db2/source/Db2Dialect.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-db2/src/main/java/org/apache/seatunnel/connectors/seatunnel/cdc/db2/source/Db2Dialect.java @@ -20,6 +20,7 @@ package org.apache.seatunnel.connectors.seatunnel.cdc.db2.source; import org.apache.seatunnel.api.table.catalog.CatalogTable; import org.apache.seatunnel.api.table.catalog.ConstraintKey; import org.apache.seatunnel.api.table.catalog.PrimaryKey; +import org.apache.seatunnel.api.table.catalog.TablePath; import org.apache.seatunnel.common.utils.SeaTunnelException; import org.apache.seatunnel.connectors.cdc.base.config.JdbcSourceConfig; import org.apache.seatunnel.connectors.cdc.base.dialect.JdbcDataSourceDialect; @@ -101,6 +102,17 @@ public class Db2Dialect implements JdbcDataSourceDialect { } } + /** + * Converts a SeaTunnel table path to the empty-catalog identifier emitted by Db2 Debezium. + * + * @param tablePath table path from checkpoint schema state + * @return Db2 Debezium table identifier + */ + @Override + public TableId toTableId(TablePath tablePath) { + return new TableId("", tablePath.getSchemaName(), tablePath.getTableName()); + } + @Override public void checkAllTablesEnabledCapture(JdbcConnection jdbcConnection, List<TableId> tableIds) throws SQLException { diff --git a/seatunnel-connectors-v2/connector-cdc/connector-cdc-db2/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/db2/source/Db2IncrementalSourceFactoryTest.java b/seatunnel-connectors-v2/connector-cdc/connector-cdc-db2/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/db2/source/Db2IncrementalSourceFactoryTest.java index 853215c542..a977cb58a1 100644 --- a/seatunnel-connectors-v2/connector-cdc/connector-cdc-db2/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/db2/source/Db2IncrementalSourceFactoryTest.java +++ b/seatunnel-connectors-v2/connector-cdc/connector-cdc-db2/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/db2/source/Db2IncrementalSourceFactoryTest.java @@ -17,12 +17,33 @@ package org.apache.seatunnel.connectors.seatunnel.cdc.db2.source; +import org.apache.seatunnel.api.table.catalog.TablePath; +import org.apache.seatunnel.connectors.seatunnel.cdc.db2.config.Db2SourceConfig; +import org.apache.seatunnel.connectors.seatunnel.cdc.db2.config.Db2SourceConfigFactory; + import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import org.mockito.Mockito; + +import io.debezium.relational.TableId; + +import java.util.Collections; class Db2IncrementalSourceFactoryTest { @Test public void testOptionRule() { Assertions.assertNotNull((new Db2IncrementalSourceFactory()).optionRule()); } + + /** Verifies Db2 checkpoint tables use the empty-catalog identifier returned by discovery. */ + @Test + public void testToTableIdDropsCatalog() { + Db2SourceConfigFactory configFactory = Mockito.mock(Db2SourceConfigFactory.class); + Mockito.when(configFactory.create(0)).thenReturn(Mockito.mock(Db2SourceConfig.class)); + Db2Dialect dialect = new Db2Dialect(configFactory, Collections.emptyList()); + + Assertions.assertEquals( + new TableId("", "DB2INST1", "CUSTOMERS"), + dialect.toTableId(TablePath.of("SAMPLE", "DB2INST1", "CUSTOMERS"))); + } }
