|
|
@ -111,6 +111,7 @@ public class JdbcSourceScanFetcher implements Fetcher<SourceRecord, SourceSplitB
|
|
|
|
boolean reachBinlogEnd = false;
|
|
|
|
boolean reachBinlogEnd = false;
|
|
|
|
final List<SourceRecord> sourceRecords = new ArrayList<>();
|
|
|
|
final List<SourceRecord> sourceRecords = new ArrayList<>();
|
|
|
|
while (!reachBinlogEnd) {
|
|
|
|
while (!reachBinlogEnd) {
|
|
|
|
|
|
|
|
checkReadException();
|
|
|
|
List<DataChangeEvent> batch = queue.poll();
|
|
|
|
List<DataChangeEvent> batch = queue.poll();
|
|
|
|
for (DataChangeEvent event : batch) {
|
|
|
|
for (DataChangeEvent event : batch) {
|
|
|
|
sourceRecords.add(event.getRecord());
|
|
|
|
sourceRecords.add(event.getRecord());
|
|
|
|