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;
 

Reply via email to