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

dockerzhang 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 7dae4a193 [INLONG-8005][Sort] Fix duplicate split request when add new 
table in mysql connector (#8059)
7dae4a193 is described below

commit 7dae4a1938305656793a7467bff07d121c73c8d7
Author: Schnapps <[email protected]>
AuthorDate: Fri May 19 22:22:25 2023 +0800

    [INLONG-8005][Sort] Fix duplicate split request when add new table in mysql 
connector (#8059)
---
 .../inlong/sort/cdc/mysql/source/reader/MySqlSourceReader.java | 10 ++++++++--
 1 file changed, 8 insertions(+), 2 deletions(-)

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/reader/MySqlSourceReader.java
 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/reader/MySqlSourceReader.java
index 10ae8b4f5..03bdaa3f8 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/reader/MySqlSourceReader.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/source/reader/MySqlSourceReader.java
@@ -168,22 +168,28 @@ public class MySqlSourceReader<T>
 
     @Override
     protected void onSplitFinished(Map<String, MySqlSplitState> 
finishedSplitIds) {
+        boolean requestNextSplit = true;
         for (MySqlSplitState mySqlSplitState : finishedSplitIds.values()) {
             MySqlSplit mySqlSplit = mySqlSplitState.toMySqlSplit();
             if (mySqlSplit.isBinlogSplit()) {
                 LOG.info(
                         "binlog split reader suspended due to newly added 
table, offset {}",
                         
mySqlSplitState.asBinlogSplitState().getStartingOffset());
-
                 mySqlSourceReaderContext.resetStopBinlogSplitReader();
                 suspendedBinlogSplit = 
toSuspendedBinlogSplit(mySqlSplit.asBinlogSplit());
                 context.sendSourceEventToCoordinator(new 
SuspendBinlogReaderAckEvent());
+                // do not request next split when the reader is suspended, the 
suspended reader will
+                // automatically request the next split after it has been 
wakeup
+                requestNextSplit = false;
             } else {
                 finishedUnackedSplits.put(mySqlSplit.splitId(), 
mySqlSplit.asSnapshotSplit());
             }
         }
         reportFinishedSnapshotSplitsIfNeed();
-        context.sendSplitRequest();
+        LOG.info("request for new split finished unacked splits:{}", 
finishedUnackedSplits);
+        if (requestNextSplit) {
+            context.sendSplitRequest();
+        }
     }
 
     @Override

Reply via email to