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