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 6dd635496 [INLONG-8060][Sort] Let mysql source reader throws 
runtimeException when connect times out or OOM (#8061)
6dd635496 is described below

commit 6dd635496c59321afa971ea10f3ee448e3e09259
Author: Schnapps <[email protected]>
AuthorDate: Fri May 19 22:21:54 2023 +0800

    [INLONG-8060][Sort] Let mysql source reader throws runtimeException when 
connect times out or OOM (#8061)
---
 .../inlong/sort/cdc/mysql/debezium/reader/SnapshotSplitReader.java       | 1 +
 1 file changed, 1 insertion(+)

diff --git 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/debezium/reader/SnapshotSplitReader.java
 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/debezium/reader/SnapshotSplitReader.java
index 1d62a929c..87effb4d7 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/debezium/reader/SnapshotSplitReader.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/mysql-cdc/src/main/java/org/apache/inlong/sort/cdc/mysql/debezium/reader/SnapshotSplitReader.java
@@ -222,6 +222,7 @@ public class SnapshotSplitReader implements 
DebeziumReader<SourceRecord, MySqlSp
             boolean reachBinlogEnd = false;
             final List<SourceRecord> sourceRecords = new ArrayList<>();
             while (!reachBinlogEnd) {
+                checkReadException();
                 List<DataChangeEvent> batch = queue.poll();
                 for (DataChangeEvent event : batch) {
                     sourceRecords.add(event.getRecord());

Reply via email to