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());