diff --git a/flink-connector-tidb-cdc/src/main/java/com/ververica/cdc/connectors/tidb/TiKVRichParallelSourceFunction.java b/flink-connector-tidb-cdc/src/main/java/com/ververica/cdc/connectors/tidb/TiKVRichParallelSourceFunction.java index bf1fae2fa..697d736bb 100644 --- a/flink-connector-tidb-cdc/src/main/java/com/ververica/cdc/connectors/tidb/TiKVRichParallelSourceFunction.java +++ b/flink-connector-tidb-cdc/src/main/java/com/ververica/cdc/connectors/tidb/TiKVRichParallelSourceFunction.java @@ -243,7 +243,7 @@ public class TiKVRichParallelSourceFunction extends RichParallelSourceFunctio } handleRow(row); } - resolvedTs = cdcClient.getMinResolvedTs(); + resolvedTs = cdcClient.getMaxResolvedTs(); if (commits.size() > 0) { flushRows(resolvedTs); }