-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtmp.diff
More file actions
25 lines (24 loc) · 1.69 KB
/
Copy pathtmp.diff
File metadata and controls
25 lines (24 loc) · 1.69 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
diff --git a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/assigners/MySqlSnapshotSplitAssigner.java b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/assigners/MySqlSnapshotSplitAssigner.java
index 1acbeac9..dbaeea6f 100644
--- a/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/assigners/MySqlSnapshotSplitAssigner.java
+++ b/flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/main/java/org/apache/flink/cdc/connectors/mysql/source/assigners/MySqlSnapshotSplitAssigner.java
@@ -301,7 +301,11 @@ public class MySqlSnapshotSplitAssigner implements MySqlSplitAssigner {
boolean hasRecordSchema = false;
// split the given table into chunks (snapshot splits)
do {
- synchronized (lock) {
+ try {
+ Thread.sleep(100);
+ } catch (InterruptedException e) {
+ // ignore
+ }
List<MySqlSnapshotSplit> splits;
try {
splits = chunkSplitter.splitChunks(partition, nextTable);
@@ -310,6 +314,7 @@ public class MySqlSnapshotSplitAssigner implements MySqlSplitAssigner {
"Error when splitting chunks for " + nextTable, e);
}
+ synchronized (lock) {
if (!hasRecordSchema && !splits.isEmpty()) {
hasRecordSchema = true;
final Map<TableId, TableChanges.TableChange> tableSchema = new HashMap<>();