This is an automated email from the ASF dual-hosted git repository.

zirui pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 792001139 [INLONG-8038][Sort] Optimize MySQL CDC chunk splitting to 
prevent large chunk (#8044)
792001139 is described below

commit 79200113935c88745b5ea22d41440bec7ac4aac7
Author: emhui <[email protected]>
AuthorDate: Tue May 30 14:27:46 2023 +0800

    [INLONG-8038][Sort] Optimize MySQL CDC chunk splitting to prevent large 
chunk (#8044)
---
 .../cdc/mysql/source/assigners/MySqlSnapshotSplitAssigner.java | 10 +++++++++-
 1 file changed, 9 insertions(+), 1 deletion(-)

diff --git 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/assigners/MySqlSnapshotSplitAssigner.java
 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/assigners/MySqlSnapshotSplitAssigner.java
index dc17534c2..0d7b14c42 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/assigners/MySqlSnapshotSplitAssigner.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/assigners/MySqlSnapshotSplitAssigner.java
@@ -230,7 +230,7 @@ public class MySqlSnapshotSplitAssigner implements 
MySqlSplitAssigner {
                                                 
.map(MySqlSnapshotSplit::toSchemaLessSnapshotSplit)
                                                 .collect(Collectors.toList());
                                 synchronized (lock) {
-                                    
remainingSplits.addAll(schemaLessSnapshotSplits);
+                                    
addNewlyAddedSplits(schemaLessSnapshotSplits);
                                     remainingTables.remove(nextTable);
                                     
addAlreadyProcessedTablesIfNotExists(nextTable);
                                     lock.notify();
@@ -399,6 +399,14 @@ public class MySqlSnapshotSplitAssigner implements 
MySqlSplitAssigner {
         }
     }
 
+    private void addNewlyAddedSplits(List<MySqlSchemalessSnapshotSplit> 
schemaLessSnapshotSplits) {
+        int size = schemaLessSnapshotSplits.size();
+        // move the last snapshot split to the front of the remaining splits 
to prevent OOM
+        // caused by the excessive data of the last snapshot split.
+        remainingSplits.add(0, schemaLessSnapshotSplits.get(size - 1));
+        remainingSplits.addAll(schemaLessSnapshotSplits.subList(0, size - 1));
+    }
+
     private void addAlreadyProcessedTablesIfNotExists(TableId tableId) {
         if (!alreadyProcessedTables.contains(tableId)) {
             alreadyProcessedTables.add(tableId);

Reply via email to