Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -257,9 +255,7 @@ public void close() throws IOException {

public StreamSplit createStreamSplit() {
final List<SchemalessSnapshotSplit> assignedSnapshotSplit =
snapshotSplitAssigner.getAssignedSplits().values().stream()
.sorted(Comparator.comparing(SourceSplitBase::splitId))
.collect(Collectors.toList());
new ArrayList<>(snapshotSplitAssigner.getAssignedSplits().values());

Map<String, Offset> splitFinishedOffsets = snapshotSplitAssigner.getSplitFinishedOffsets();
final List<FinishedSnapshotSplitInfo> finishedSnapshotSplitInfos = new ArrayList<>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -72,7 +73,21 @@ public class SnapshotSplitAssigner<C extends SourceConfig> implements SplitAssig

private final List<TableId> alreadyProcessedTables;
private final List<SchemalessSnapshotSplit> remainingSplits;
private final Map<String, SchemalessSnapshotSplit> 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.
*
* <p>{@link IncrementalSourceReader} relies on the order of elements within the map:
*
* <ol>
* <li>It must correspond to the order of assignment of the splits to readers.
* <li>The order must be retained across job restarts.
* </ol>
*/
private final LinkedHashMap<String, SchemalessSnapshotSplit> assignedSplits;

private final Map<TableId, TableChanges.TableChange> tableSchemas;
private final Map<String, Offset> splitFinishedOffsets;

Expand Down Expand Up @@ -152,7 +167,7 @@ private SnapshotSplitAssigner(
int currentParallelism,
List<TableId> alreadyProcessedTables,
List<SchemalessSnapshotSplit> remainingSplits,
Map<String, SchemalessSnapshotSplit> assignedSplits,
LinkedHashMap<String, SchemalessSnapshotSplit> assignedSplits,
Map<TableId, TableChanges.TableChange> tableSchemas,
Map<String, Offset> splitFinishedOffsets,
AssignerStatus assignerStatus,
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -214,12 +215,12 @@ private SnapshotPendingSplitsState deserializeLegacySnapshotPendingSplitsState(
int splitVersion, DataInputDeserializer in) throws IOException {
List<TableId> alreadyProcessedTables = readTableIds(2, in);
List<SnapshotSplit> remainingSplits = readSnapshotSplits(splitVersion, in);
Map<String, SnapshotSplit> assignedSnapshotSplits =
LinkedHashMap<String, SnapshotSplit> assignedSnapshotSplits =
readAssignedSnapshotSplits(splitVersion, in);

final List<SchemalessSnapshotSplit> remainingSchemalessSplits = new ArrayList<>();
final Map<String, SchemalessSnapshotSplit> assignedSchemalessSnapshotSplits =
new HashMap<>();
final LinkedHashMap<String, SchemalessSnapshotSplit> assignedSchemalessSnapshotSplits =
new LinkedHashMap<>();
final Map<TableId, TableChanges.TableChange> tableSchemas = new HashMap<>();
remainingSplits.forEach(
split -> {
Expand Down Expand Up @@ -268,7 +269,7 @@ private SnapshotPendingSplitsState deserializeSnapshotPendingSplitsState(
int version, int splitVersion, DataInputDeserializer in) throws IOException {
List<TableId> alreadyProcessedTables = readTableIds(version, in);
List<SnapshotSplit> remainingSplits = readSnapshotSplits(splitVersion, in);
Map<String, SnapshotSplit> assignedSnapshotSplits =
LinkedHashMap<String, SnapshotSplit> assignedSnapshotSplits =
readAssignedSnapshotSplits(splitVersion, in);
Map<String, Offset> finishedOffsets = readFinishedOffsets(splitVersion, in);
AssignerStatus assignerStatus;
Expand All @@ -285,8 +286,8 @@ private SnapshotPendingSplitsState deserializeSnapshotPendingSplitsState(
List<TableId> remainingTableIds = readTableIds(version, in);
boolean isTableIdCaseSensitive = in.readBoolean();
final List<SchemalessSnapshotSplit> remainingSchemalessSplits = new ArrayList<>();
final Map<String, SchemalessSnapshotSplit> assignedSchemalessSnapshotSplits =
new HashMap<>();
final LinkedHashMap<String, SchemalessSnapshotSplit> assignedSchemalessSnapshotSplits =
new LinkedHashMap<>();
final Map<TableId, TableChanges.TableChange> tableSchemas = new HashMap<>();
remainingSplits.forEach(
split -> {
Expand Down Expand Up @@ -415,9 +416,9 @@ private void writeAssignedSnapshotSplits(
}
}

private Map<String, SnapshotSplit> readAssignedSnapshotSplits(
private LinkedHashMap<String, SnapshotSplit> readAssignedSnapshotSplits(
int splitVersion, DataInputDeserializer in) throws IOException {
Map<String, SnapshotSplit> assignedSplits = new HashMap<>();
LinkedHashMap<String, SnapshotSplit> assignedSplits = new LinkedHashMap<>();
final int size = in.readInt();
for (int i = 0; i < size; i++) {
String splitId = in.readUTF();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, SchemalessSnapshotSplit> assignedSplits;
private final LinkedHashMap<String, SchemalessSnapshotSplit> assignedSplits;

/**
* The offsets of finished (snapshot) splits that the {@link IncrementalSourceEnumerator} has
Expand Down Expand Up @@ -83,7 +83,7 @@ public class SnapshotPendingSplitsState extends PendingSplitsState {
public SnapshotPendingSplitsState(
List<TableId> alreadyProcessedTables,
List<SchemalessSnapshotSplit> remainingSplits,
Map<String, SchemalessSnapshotSplit> assignedSplits,
LinkedHashMap<String, SchemalessSnapshotSplit> assignedSplits,
Map<TableId, TableChanges.TableChange> tableSchemas,
Map<String, Offset> splitFinishedOffsets,
AssignerStatus assignerStatus,
Expand Down Expand Up @@ -119,7 +119,7 @@ public List<SchemalessSnapshotSplit> getRemainingSplits() {
return remainingSplits;
}

public Map<String, SchemalessSnapshotSplit> getAssignedSplits() {
public LinkedHashMap<String, SchemalessSnapshotSplit> getAssignedSplits() {
return assignedSplits;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<FinishedSnapshotSplitInfo> finishedSnapshotSplitInfos;

private final Map<TableId, TableChange> tableSchemas;
private final int totalFinishedSplitSize;

Expand All @@ -66,6 +70,9 @@ public StreamSplit(
boolean isSuspended,
boolean isSnapshotCompleted) {
super(splitId);

ensureNoDuplicates(finishedSnapshotSplitInfos);

this.startingOffset = startingOffset;
this.endingOffset = endingOffset;
this.finishedSnapshotSplitInfos = finishedSnapshotSplitInfos;
Expand All @@ -82,14 +89,30 @@ public StreamSplit(
List<FinishedSnapshotSplitInfo> finishedSnapshotSplitInfos,
Map<TableId, TableChange> 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<FinishedSnapshotSplitInfo> finishedSnapshotSplitInfos) {
Set<String> 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() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -361,28 +359,6 @@ private StreamSplit discoverTableSchemasForStreamSplit(
}
}

private Set<String> getExistedSplitsOfLastGroup(
List<FinishedSnapshotSplitInfo> 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<String> 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) {
Expand Down Expand Up @@ -437,15 +413,24 @@ private void fillMetaDataForStreamSplit(StreamSplitMetaEvent metadataEvent) {
streamSplit = toNormalStreamSplit(streamSplit, receivedTotalFinishedSplitSize);
uncompletedStreamSplits.put(streamSplit.splitId(), streamSplit);
} else if (receivedMetaGroupId == expectedMetaGroupId) {
Set<String> existedSplitsOfLastGroup =
getExistedSplitsOfLastGroup(
streamSplit.getFinishedSnapshotSplitInfos(),
sourceConfig.getSplitMetaGroupSize());

int expectedNumberOfAlreadyRetrievedElements =
streamSplit.getFinishedSnapshotSplitInfos().size()
% sourceConfig.getSplitMetaGroupSize();
List<byte[]> 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<FinishedSnapshotSplitInfo> newAddedMetadataGroup =
metadataEvent.getMetaGroup().stream()
metaGroup.stream()
.map(sourceSplitSerializer::deserialize)
.filter(r -> !existedSplitsOfLastGroup.contains(r.getSplitId()))
.collect(Collectors.toList());
uncompletedStreamSplits.put(
streamSplit.splitId(),
Expand Down
Loading