This is an automated email from the ASF dual-hosted git repository.
loserwang1024 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git
The following commit(s) were added to refs/heads/master by this push:
new 61e0d2c71 [FLINK-40202][cdc-base] Init enumerator metrics before
opening SnapshotSplitAssigner to avoid ConcurrentModificationException (#4480)
61e0d2c71 is described below
commit 61e0d2c71846af5f49aa14211ace9d10ff0f9077
Author: naivedogger <[email protected]>
AuthorDate: Wed Jul 22 14:41:54 2026 +0800
[FLINK-40202][cdc-base] Init enumerator metrics before opening
SnapshotSplitAssigner to avoid ConcurrentModificationException (#4480)
---
.../cdc/connectors/base/source/assigner/HybridSplitAssigner.java | 6 ++++--
.../cdc/connectors/base/source/assigner/SnapshotSplitAssigner.java | 5 ++++-
2 files changed, 8 insertions(+), 3 deletions(-)
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 a4d21ada6..7d2ce9f61 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
@@ -136,9 +136,11 @@ public class HybridSplitAssigner<C extends SourceConfig>
implements SplitAssigne
enumeratorMetrics.exitStreamReading();
}
- snapshotSplitAssigner.open();
- // init enumerator metrics
+ // Init enumerator metrics before opening the snapshot assigner.
Opening the assigner
+ // starts an asynchronous splitting thread that mutates
remainingSplits, so metrics must be
+ // initialized first to avoid a ConcurrentModificationException while
iterating the splits.
snapshotSplitAssigner.initEnumeratorMetrics(enumeratorMetrics);
+ snapshotSplitAssigner.open();
}
@Override
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 0b1a4100b..69de3d377 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
@@ -295,7 +295,10 @@ public class SnapshotSplitAssigner<C extends SourceConfig>
implements SplitAssig
}
}
- /** This should be invoked after this class's open method. */
+ /**
+ * This must be invoked before this class's {@link #open()} method,
because {@code open()}
+ * starts an asynchronous splitting thread that concurrently mutates
{@code remainingSplits}.
+ */
public void initEnumeratorMetrics(SourceEnumeratorMetrics
enumeratorMetrics) {
this.enumeratorMetrics = enumeratorMetrics;