From e004ef309dab8bbae78c101bd77ac2ba936a0a19 Mon Sep 17 00:00:00 2001 From: dev-donghwan Date: Wed, 22 Jul 2026 19:07:43 +0900 Subject: [PATCH] [FLINK-40216][cdc-base] Preserve assignment order of splits across restore to keep meta groups consistent The restore constructor of SnapshotSplitAssigner re-sorted the assigned splits lexicographically by split id, while at runtime splits are kept in assignment order. Since the enumerator partitions getFinishedSplitInfos() into meta groups, the groups served before and after a restore differed, so a reader resuming an incomplete meta synchronization received duplicated split infos while never receiving others, leading to silently dropped change events in the stream phase. Port the FLINK-38218 fix from the MySQL connector to flink-cdc-base: - deserialize assignedSplits into a LinkedHashMap and keep the checkpointed assignment order instead of re-sorting on restore - remove the lexicographic sort in HybridSplitAssigner#createStreamSplit - replace the order-dependent deduplication in IncrementalSourceReader with the order-agnostic prefix-discard approach - reject duplicated split infos in StreamSplit --- .../source/assigner/HybridSplitAssigner.java | 6 +- .../assigner/SnapshotSplitAssigner.java | 31 +-- .../state/PendingSplitsStateSerializer.java | 17 +- .../state/SnapshotPendingSplitsState.java | 6 +- .../base/source/meta/split/StreamSplit.java | 39 +++- .../reader/IncrementalSourceReader.java | 47 ++--- .../assigner/MetaGroupOrderingTest.java | 182 ++++++++++++++++++ .../PendingSplitsStateSerializerTest.java | 53 ++++- .../source/meta/split/StreamSplitTest.java | 62 ++++++ 9 files changed, 374 insertions(+), 69 deletions(-) create mode 100644 flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/assigner/MetaGroupOrderingTest.java create mode 100644 flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplitTest.java diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/HybridSplitAssigner.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/HybridSplitAssigner.java index 7d2ce9f614f..54734602771 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/HybridSplitAssigner.java +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/HybridSplitAssigner.java @@ -38,12 +38,10 @@ import java.io.IOException; import java.util.ArrayList; import java.util.Collection; -import java.util.Comparator; import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Optional; -import java.util.stream.Collectors; import static org.apache.flink.cdc.connectors.base.source.assigner.AssignerStatus.isInitialAssigningFinished; import static org.apache.flink.cdc.connectors.base.source.assigner.AssignerStatus.isNewlyAddedAssigningFinished; @@ -257,9 +255,7 @@ public void close() throws IOException { public StreamSplit createStreamSplit() { final List assignedSnapshotSplit = - snapshotSplitAssigner.getAssignedSplits().values().stream() - .sorted(Comparator.comparing(SourceSplitBase::splitId)) - .collect(Collectors.toList()); + new ArrayList<>(snapshotSplitAssigner.getAssignedSplits().values()); Map splitFinishedOffsets = snapshotSplitAssigner.getSplitFinishedOffsets(); final List finishedSnapshotSplitInfos = new ArrayList<>(); diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/SnapshotSplitAssigner.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/SnapshotSplitAssigner.java index 69de3d3779e..4a35cc781a6 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/SnapshotSplitAssigner.java +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/SnapshotSplitAssigner.java @@ -29,6 +29,7 @@ import org.apache.flink.cdc.connectors.base.source.meta.split.SnapshotSplit; import org.apache.flink.cdc.connectors.base.source.meta.split.SourceSplitBase; import org.apache.flink.cdc.connectors.base.source.metrics.SourceEnumeratorMetrics; +import org.apache.flink.cdc.connectors.base.source.reader.IncrementalSourceReader; import org.apache.flink.util.FlinkRuntimeException; import org.apache.flink.util.Preconditions; @@ -72,7 +73,21 @@ public class SnapshotSplitAssigner implements SplitAssig private final List alreadyProcessedTables; private final List remainingSplits; - private final Map assignedSplits; + + /** + * The splits that have been assigned to a reader. Once a split is finished, it remains in this + * map. An entry added to {@link #splitFinishedOffsets} indicates that the split has been + * finished. If reading the split fails, it is removed from this map. + * + *

{@link IncrementalSourceReader} relies on the order of elements within the map: + * + *

    + *
  1. It must correspond to the order of assignment of the splits to readers. + *
  2. The order must be retained across job restarts. + *
+ */ + private final LinkedHashMap assignedSplits; + private final Map tableSchemas; private final Map splitFinishedOffsets; @@ -152,7 +167,7 @@ private SnapshotSplitAssigner( int currentParallelism, List alreadyProcessedTables, List remainingSplits, - Map assignedSplits, + LinkedHashMap assignedSplits, Map tableSchemas, Map splitFinishedOffsets, AssignerStatus assignerStatus, @@ -167,17 +182,7 @@ private SnapshotSplitAssigner( this.currentParallelism = currentParallelism; this.alreadyProcessedTables = alreadyProcessedTables; this.remainingSplits = remainingSplits; - // When job restore from savepoint, sort the existing tables and newly added tables - // to let enumerator only send newly added tables' StreamSplitMetaEvent - this.assignedSplits = - assignedSplits.entrySet().stream() - .sorted(Map.Entry.comparingByKey()) - .collect( - Collectors.toMap( - Map.Entry::getKey, - Map.Entry::getValue, - (o, o2) -> o, - LinkedHashMap::new)); + this.assignedSplits = assignedSplits; this.tableSchemas = tableSchemas; this.splitFinishedOffsets = splitFinishedOffsets; this.assignerStatus = assignerStatus; diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializer.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializer.java index c823c74f2b5..d8e720b1dce 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializer.java +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializer.java @@ -35,6 +35,7 @@ import java.util.ArrayList; import java.util.Collection; import java.util.HashMap; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; @@ -214,12 +215,12 @@ private SnapshotPendingSplitsState deserializeLegacySnapshotPendingSplitsState( int splitVersion, DataInputDeserializer in) throws IOException { List alreadyProcessedTables = readTableIds(2, in); List remainingSplits = readSnapshotSplits(splitVersion, in); - Map assignedSnapshotSplits = + LinkedHashMap assignedSnapshotSplits = readAssignedSnapshotSplits(splitVersion, in); final List remainingSchemalessSplits = new ArrayList<>(); - final Map assignedSchemalessSnapshotSplits = - new HashMap<>(); + final LinkedHashMap assignedSchemalessSnapshotSplits = + new LinkedHashMap<>(); final Map tableSchemas = new HashMap<>(); remainingSplits.forEach( split -> { @@ -268,7 +269,7 @@ private SnapshotPendingSplitsState deserializeSnapshotPendingSplitsState( int version, int splitVersion, DataInputDeserializer in) throws IOException { List alreadyProcessedTables = readTableIds(version, in); List remainingSplits = readSnapshotSplits(splitVersion, in); - Map assignedSnapshotSplits = + LinkedHashMap assignedSnapshotSplits = readAssignedSnapshotSplits(splitVersion, in); Map finishedOffsets = readFinishedOffsets(splitVersion, in); AssignerStatus assignerStatus; @@ -285,8 +286,8 @@ private SnapshotPendingSplitsState deserializeSnapshotPendingSplitsState( List remainingTableIds = readTableIds(version, in); boolean isTableIdCaseSensitive = in.readBoolean(); final List remainingSchemalessSplits = new ArrayList<>(); - final Map assignedSchemalessSnapshotSplits = - new HashMap<>(); + final LinkedHashMap assignedSchemalessSnapshotSplits = + new LinkedHashMap<>(); final Map tableSchemas = new HashMap<>(); remainingSplits.forEach( split -> { @@ -415,9 +416,9 @@ private void writeAssignedSnapshotSplits( } } - private Map readAssignedSnapshotSplits( + private LinkedHashMap readAssignedSnapshotSplits( int splitVersion, DataInputDeserializer in) throws IOException { - Map assignedSplits = new HashMap<>(); + LinkedHashMap assignedSplits = new LinkedHashMap<>(); final int size = in.readInt(); for (int i = 0; i < size; i++) { String splitId = in.readUTF(); diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/SnapshotPendingSplitsState.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/SnapshotPendingSplitsState.java index ca5e578bddc..14830e75e3c 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/SnapshotPendingSplitsState.java +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/assigner/state/SnapshotPendingSplitsState.java @@ -52,7 +52,7 @@ public class SnapshotPendingSplitsState extends PendingSplitsState { * The snapshot splits that the {@link IncrementalSourceEnumerator} has assigned to {@link * IncrementalSourceSplitReader}s. */ - private final Map assignedSplits; + private final LinkedHashMap assignedSplits; /** * The offsets of finished (snapshot) splits that the {@link IncrementalSourceEnumerator} has @@ -83,7 +83,7 @@ public class SnapshotPendingSplitsState extends PendingSplitsState { public SnapshotPendingSplitsState( List alreadyProcessedTables, List remainingSplits, - Map assignedSplits, + LinkedHashMap assignedSplits, Map tableSchemas, Map splitFinishedOffsets, AssignerStatus assignerStatus, @@ -119,7 +119,7 @@ public List getRemainingSplits() { return remainingSplits; } - public Map getAssignedSplits() { + public LinkedHashMap getAssignedSplits() { return assignedSplits; } diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplit.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplit.java index cdeb247f3ad..99b45c2ea8f 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplit.java +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplit.java @@ -28,6 +28,7 @@ import java.util.ArrayList; import java.util.HashMap; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Objects; @@ -42,7 +43,10 @@ public class StreamSplit extends SourceSplitBase { private final Offset startingOffset; private final Offset endingOffset; + + /** Split IDs of all elements must be unique. */ private final List finishedSnapshotSplitInfos; + private final Map tableSchemas; private final int totalFinishedSplitSize; @@ -66,6 +70,9 @@ public StreamSplit( boolean isSuspended, boolean isSnapshotCompleted) { super(splitId); + + ensureNoDuplicates(finishedSnapshotSplitInfos); + this.startingOffset = startingOffset; this.endingOffset = endingOffset; this.finishedSnapshotSplitInfos = finishedSnapshotSplitInfos; @@ -82,14 +89,30 @@ public StreamSplit( List finishedSnapshotSplitInfos, Map tableSchemas, int totalFinishedSplitSize) { - super(splitId); - this.startingOffset = startingOffset; - this.endingOffset = endingOffset; - this.finishedSnapshotSplitInfos = finishedSnapshotSplitInfos; - this.tableSchemas = tableSchemas; - this.totalFinishedSplitSize = totalFinishedSplitSize; - this.isSuspended = false; - this.isSnapshotCompleted = false; + this( + splitId, + startingOffset, + endingOffset, + finishedSnapshotSplitInfos, + tableSchemas, + totalFinishedSplitSize, + false, + false); + } + + private static void ensureNoDuplicates( + List finishedSnapshotSplitInfos) { + Set seenSplitIds = new HashSet<>(); + for (FinishedSnapshotSplitInfo splitInfo : finishedSnapshotSplitInfos) { + if (seenSplitIds.contains(splitInfo.getSplitId())) { + throw new IllegalArgumentException( + String.format( + "Found duplicate split ID %s in finished snapshot split infos", + splitInfo.getSplitId())); + } + + seenSplitIds.add(splitInfo.getSplitId()); + } } public Offset getStartingOffset() { diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/reader/IncrementalSourceReader.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/reader/IncrementalSourceReader.java index 15666dc7663..c83de9b3eab 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/reader/IncrementalSourceReader.java +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/main/java/org/apache/flink/cdc/connectors/base/source/reader/IncrementalSourceReader.java @@ -59,10 +59,8 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; -import java.util.HashSet; import java.util.List; import java.util.Map; -import java.util.Set; import java.util.function.Supplier; import java.util.stream.Collectors; @@ -361,28 +359,6 @@ private StreamSplit discoverTableSchemasForStreamSplit( } } - private Set getExistedSplitsOfLastGroup( - List finishedSnapshotSplits, int metaGroupSize) { - int splitsNumOfLastGroup = - finishedSnapshotSplits.size() % sourceConfig.getSplitMetaGroupSize(); - if (splitsNumOfLastGroup != 0) { - int lastGroupStart = - ((int) (finishedSnapshotSplits.size() / sourceConfig.getSplitMetaGroupSize())) - * metaGroupSize; - // Keep same order with HybridSplitAssigner.createStreamSplit() to avoid - // 'invalid request meta group id' error - List sortedFinishedSnapshotSplits = - finishedSnapshotSplits.stream() - .map(FinishedSnapshotSplitInfo::getSplitId) - .sorted() - .collect(Collectors.toList()); - return new HashSet<>( - sortedFinishedSnapshotSplits.subList( - lastGroupStart, lastGroupStart + splitsNumOfLastGroup)); - } - return new HashSet<>(); - } - @Override public void handleSourceEvents(SourceEvent sourceEvent) { if (sourceEvent instanceof FinishedSnapshotSplitsAckEvent) { @@ -437,15 +413,24 @@ private void fillMetaDataForStreamSplit(StreamSplitMetaEvent metadataEvent) { streamSplit = toNormalStreamSplit(streamSplit, receivedTotalFinishedSplitSize); uncompletedStreamSplits.put(streamSplit.splitId(), streamSplit); } else if (receivedMetaGroupId == expectedMetaGroupId) { - Set existedSplitsOfLastGroup = - getExistedSplitsOfLastGroup( - streamSplit.getFinishedSnapshotSplitInfos(), - sourceConfig.getSplitMetaGroupSize()); - + int expectedNumberOfAlreadyRetrievedElements = + streamSplit.getFinishedSnapshotSplitInfos().size() + % sourceConfig.getSplitMetaGroupSize(); + List metaGroup = metadataEvent.getMetaGroup(); + if (expectedNumberOfAlreadyRetrievedElements > 0) { + LOG.info( + "Source reader {} is discarding the first {} out of {} elements of meta group {}.", + subtaskId, + expectedNumberOfAlreadyRetrievedElements, + metaGroup.size(), + receivedMetaGroupId); + metaGroup = + metaGroup.subList( + expectedNumberOfAlreadyRetrievedElements, metaGroup.size()); + } List newAddedMetadataGroup = - metadataEvent.getMetaGroup().stream() + metaGroup.stream() .map(sourceSplitSerializer::deserialize) - .filter(r -> !existedSplitsOfLastGroup.contains(r.getSplitId())) .collect(Collectors.toList()); uncompletedStreamSplits.put( streamSplit.splitId(), diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/assigner/MetaGroupOrderingTest.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/assigner/MetaGroupOrderingTest.java new file mode 100644 index 00000000000..93ad35bc5b6 --- /dev/null +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/assigner/MetaGroupOrderingTest.java @@ -0,0 +1,182 @@ +/* + * 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.flink.cdc.connectors.base.source.assigner; + +import org.apache.flink.cdc.connectors.base.config.SourceConfig; +import org.apache.flink.cdc.connectors.base.dialect.DataSourceDialect; +import org.apache.flink.cdc.connectors.base.source.assigner.state.ChunkSplitterState; +import org.apache.flink.cdc.connectors.base.source.assigner.state.SnapshotPendingSplitsState; +import org.apache.flink.cdc.connectors.base.source.meta.offset.Offset; +import org.apache.flink.cdc.connectors.base.source.meta.offset.OffsetFactory; +import org.apache.flink.cdc.connectors.base.source.meta.split.FinishedSnapshotSplitInfo; +import org.apache.flink.cdc.connectors.base.source.meta.split.SchemalessSnapshotSplit; +import org.apache.flink.cdc.connectors.base.source.reader.IncrementalSourceReader; +import org.apache.flink.table.types.logical.BigIntType; +import org.apache.flink.table.types.logical.RowType; + +import org.apache.flink.shaded.guava31.com.google.common.collect.Lists; + +import io.debezium.relational.TableId; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; + +/** + * Tests that the meta groups of finished snapshot splits stay consistent across a checkpoint + * restore (FLINK-40216). + * + *

The meta groups served by {@code IncrementalSourceEnumerator#sendStreamMetaRequestEvent} are + * partitioned from {@link SnapshotSplitAssigner#getFinishedSplitInfos()}, so the assignment order + * of the splits must be retained across job restarts; otherwise a reader that synchronized only + * part of the meta groups before a restart receives duplicated split infos while never receiving + * others. + */ +class MetaGroupOrderingTest { + + private static final TableId TABLE_ID = new TableId("db", null, "t"); + private static final RowType SPLIT_KEY_TYPE = + new RowType(Collections.singletonList(new RowType.RowField("id", new BigIntType()))); + + @Test + void testFinishedSplitInfosPreserveAssignmentOrderAfterRestore() { + List assignmentOrder = new ArrayList<>(); + SnapshotSplitAssigner restoredAssigner = restoreAssigner(assignmentOrder, 12); + + List restoredOrder = + restoredAssigner.getFinishedSplitInfos().stream() + .map(FinishedSnapshotSplitInfo::getSplitId) + .collect(Collectors.toList()); + + assertThat(restoredOrder).isEqualTo(assignmentOrder); + } + + @Test + void testReaderResumingAfterRestoreReceivesEachSplitExactlyOnce() { + int metaGroupSize = 2; + List assignmentOrder = new ArrayList<>(); + SnapshotSplitAssigner restoredAssigner = restoreAssigner(assignmentOrder, 12); + + // Before the restart the enumerator serves groups partitioned from the assignment order; + // the reader synchronized groups 0..2 into its stream split state: + // [t:0, t:1], [t:2, t:3], [t:4, t:5] + List> groupsBeforeRestart = Lists.partition(assignmentOrder, metaGroupSize); + List readerState = new ArrayList<>(); + for (int groupId = 0; groupId < 3; groupId++) { + readerState.addAll(groupsBeforeRestart.get(groupId)); + } + + // The job restarts: the enumerator rebuilds its meta groups from the restored assigner + List> groupsAfterRestart = + Lists.partition(restoredAssigner.getFinishedSplitInfos(), metaGroupSize); + + // The reader resumes requesting group ids based on how many infos it already holds + int nextGroupId = + IncrementalSourceReader.getNextMetaGroupId(readerState.size(), metaGroupSize); + while (readerState.size() < assignmentOrder.size() + && nextGroupId < groupsAfterRestart.size()) { + for (FinishedSnapshotSplitInfo info : groupsAfterRestart.get(nextGroupId)) { + readerState.add(info.getSplitId()); + } + nextGroupId = + IncrementalSourceReader.getNextMetaGroupId(readerState.size(), metaGroupSize); + } + + assertThat(readerState).containsExactlyElementsOf(assignmentOrder); + } + + private static SnapshotSplitAssigner restoreAssigner( + List assignmentOrder, int splitCount) { + LinkedHashMap assignedSplits = new LinkedHashMap<>(); + Map splitFinishedOffsets = new HashMap<>(); + // 10+ splits so that the lexicographic order of split ids (t:0, t:1, t:10, t:11, t:2, + // ...) differs from the assignment order (t:0, t:1, t:2, ..., t:11) + for (int i = 0; i < splitCount; i++) { + String splitId = TABLE_ID + ":" + i; + assignmentOrder.add(splitId); + assignedSplits.put( + splitId, + new SchemalessSnapshotSplit( + TABLE_ID, + splitId, + SPLIT_KEY_TYPE, + new Object[] {i}, + new Object[] {i + 1}, + null)); + splitFinishedOffsets.put(splitId, null); + } + + return new SnapshotSplitAssigner<>( + mock(SourceConfig.class), + 4, + new SnapshotPendingSplitsState( + Collections.singletonList(TABLE_ID), + Collections.emptyList(), + assignedSplits, + Collections.emptyMap(), + splitFinishedOffsets, + AssignerStatus.INITIAL_ASSIGNING_FINISHED, + Collections.emptyList(), + false, + true, + Collections.emptyMap(), + ChunkSplitterState.NO_SPLITTING_TABLE_STATE), + mock(DataSourceDialect.class), + new MockOffsetFactory()); + } + + private static class MockOffsetFactory extends OffsetFactory { + @Override + public Offset newOffset(Map offset) { + return null; + } + + @Override + public Offset newOffset(String filename, Long position) { + return null; + } + + @Override + public Offset newOffset(Long position) { + return null; + } + + @Override + public Offset createTimestampOffset(long timestampMillis) { + return null; + } + + @Override + public Offset createInitialOffset() { + return null; + } + + @Override + public Offset createNoStoppingOffset() { + return null; + } + } +} diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializerTest.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializerTest.java index 9ebcc36938f..bb9cb00abd7 100644 --- a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializerTest.java +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/assigner/state/PendingSplitsStateSerializerTest.java @@ -20,6 +20,7 @@ import org.apache.flink.cdc.connectors.base.source.assigner.AssignerStatus; import org.apache.flink.cdc.connectors.base.source.meta.offset.Offset; import org.apache.flink.cdc.connectors.base.source.meta.offset.OffsetFactory; +import org.apache.flink.cdc.connectors.base.source.meta.split.SchemalessSnapshotSplit; import org.apache.flink.cdc.connectors.base.source.meta.split.SnapshotSplit; import org.apache.flink.cdc.connectors.base.source.meta.split.SourceSplitSerializer; import org.apache.flink.cdc.connectors.base.source.meta.split.StreamSplit; @@ -37,6 +38,8 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; +import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import static org.assertj.core.api.Assertions.assertThat; @@ -92,7 +95,7 @@ void testSerializeSnapshotPendingSplitsState() throws Exception { new SnapshotPendingSplitsState( Collections.emptyList(), Collections.emptyList(), - Collections.emptyMap(), + new LinkedHashMap<>(), constructTableSchema(), Collections.emptyMap(), AssignerStatus.INITIAL_ASSIGNING, @@ -107,6 +110,54 @@ void testSerializeSnapshotPendingSplitsState() throws Exception { .isEqualTo(state); } + @Test + void testAssignedSplitsKeepAssignmentOrderAfterSerde() throws Exception { + TableId tableId = constructTableId(); + RowType splitKeyType = + new RowType( + Collections.singletonList(new RowType.RowField("id", new BigIntType()))); + // 10+ splits so that the lexicographic order of split ids differs from the assignment + // order, which determines the meta group partitioning of finished snapshot splits + List assignmentOrder = new ArrayList<>(); + LinkedHashMap assignedSplits = new LinkedHashMap<>(); + for (int i = 0; i < 12; i++) { + String splitId = tableId + ":" + i; + assignmentOrder.add(splitId); + assignedSplits.put( + splitId, + new SchemalessSnapshotSplit( + tableId, + splitId, + splitKeyType, + new Object[] {i}, + new Object[] {i + 1}, + null)); + } + PendingSplitsStateSerializer serializer = + new PendingSplitsStateSerializer(constructSourceSplitSerializer()); + SnapshotPendingSplitsState state = + new SnapshotPendingSplitsState( + Collections.singletonList(tableId), + Collections.emptyList(), + assignedSplits, + constructTableSchema(), + Collections.emptyMap(), + AssignerStatus.INITIAL_ASSIGNING_FINISHED, + Collections.emptyList(), + false, + true, + Collections.emptyMap(), + ChunkSplitterState.NO_SPLITTING_TABLE_STATE); + + SnapshotPendingSplitsState restoredState = + (SnapshotPendingSplitsState) + serializer.deserialize( + serializer.getVersion(), serializer.serialize(state)); + + assertThat(new ArrayList<>(restoredState.getAssignedSplits().keySet())) + .isEqualTo(assignmentOrder); + } + private SourceSplitSerializer constructSourceSplitSerializer() { return new SourceSplitSerializer() { @Override diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplitTest.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplitTest.java new file mode 100644 index 00000000000..116a6b90c68 --- /dev/null +++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-cdc-base/src/test/java/org/apache/flink/cdc/connectors/base/source/meta/split/StreamSplitTest.java @@ -0,0 +1,62 @@ +/* + * 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.flink.cdc.connectors.base.source.meta.split; + +import org.apache.flink.cdc.connectors.base.source.meta.offset.OffsetFactory; + +import io.debezium.relational.TableId; +import org.junit.jupiter.api.Test; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.mock; + +/** Tests for {@link StreamSplit}. */ +class StreamSplitTest { + + @Test + void testDuplicateFinishedSnapshotSplitInfosAreRejected() { + FinishedSnapshotSplitInfo info = + new FinishedSnapshotSplitInfo( + new TableId("catalog", "schema", "table"), + "split", + null, + null, + null, + mock(OffsetFactory.class)); + List infos = new ArrayList<>(); + infos.add(info); + infos.add(info); + + assertThatThrownBy( + () -> + new StreamSplit( + "stream-split", + null, + null, + infos, + Collections.emptyMap(), + 0, + false, + false)) + .isExactlyInstanceOf(IllegalArgumentException.class); + } +}