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 5c6aaeead7 [INLONG-10311][Sort] Implement TubeMQ Source report audit 
information exactly once And fix consuming TubeMQ data twice (#10440)
5c6aaeead7 is described below

commit 5c6aaeead7c61a40537eae2894d123ea83be02f0
Author: XiaoYou201 <[email protected]>
AuthorDate: Wed Jun 19 15:43:42 2024 +0800

    [INLONG-10311][Sort] Implement TubeMQ Source report audit information 
exactly once And fix consuming TubeMQ data twice (#10440)
---
 .../inlong/sort/base/metric/MetricsCollector.java  |  11 +-
 ...rceMetricData.java => SourceExactlyMetric.java} | 152 ++++++---------------
 .../inlong/sort/base/metric/SourceMetricData.java  |   7 +-
 .../sort/base/metric/SourceMetricsReporter.java    |  24 ++++
 .../inlong/sort/tubemq/FlinkTubeMQConsumer.java    |  53 +++++--
 .../table/DynamicTubeMQDeserializationSchema.java  |   6 +-
 .../DynamicTubeMQTableDeserializationSchema.java   |  27 +++-
 7 files changed, 139 insertions(+), 141 deletions(-)

diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/MetricsCollector.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/MetricsCollector.java
index 0bf1a6010a..0e26c6dce9 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/MetricsCollector.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/MetricsCollector.java
@@ -31,11 +31,10 @@ public class MetricsCollector<T> implements 
TimestampedCollector<T> {
 
     private long timestampMillis;
 
-    SourceMetricData metricData;
-
+    private final SourceMetricsReporter sourceMetricsReporter;
     public MetricsCollector(Collector<T> collector,
-            SourceMetricData sourceMetricData) {
-        this.metricData = sourceMetricData;
+            SourceMetricsReporter sourceMetricsReporter) {
+        this.sourceMetricsReporter = sourceMetricsReporter;
         this.collector = collector;
     }
 
@@ -45,8 +44,8 @@ public class MetricsCollector<T> implements 
TimestampedCollector<T> {
 
     @Override
     public void collect(T record) {
-        if (metricData != null) {
-            metricData.outputMetricsWithEstimate(record, timestampMillis);
+        if (sourceMetricsReporter != null) {
+            sourceMetricsReporter.outputMetricsWithEstimate(record, 
timestampMillis);
         }
         collector.collect(record);
     }
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricData.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceExactlyMetric.java
similarity index 66%
copy from 
inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricData.java
copy to 
inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceExactlyMetric.java
index 0d2035e71c..19f9f1eda9 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricData.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceExactlyMetric.java
@@ -17,20 +17,20 @@
 
 package org.apache.inlong.sort.base.metric;
 
-import org.apache.inlong.audit.AuditOperator;
+import org.apache.inlong.audit.AuditReporterImpl;
 
 import org.apache.flink.metrics.Counter;
 import org.apache.flink.metrics.Gauge;
 import org.apache.flink.metrics.Meter;
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.metrics.SimpleCounter;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
 
 import java.io.Serializable;
 import java.util.List;
 import java.util.Map;
 
+import static org.apache.inlong.audit.consts.ConfigConstants.DEFAULT_AUDIT_TAG;
+import static 
org.apache.inlong.common.constant.Constants.DEFAULT_AUDIT_VERSION;
 import static 
org.apache.inlong.sort.base.Constants.CURRENT_EMIT_EVENT_TIME_LAG;
 import static 
org.apache.inlong.sort.base.Constants.CURRENT_FETCH_EVENT_TIME_LAG;
 import static org.apache.inlong.sort.base.Constants.NUM_BYTES_IN;
@@ -41,13 +41,9 @@ import static 
org.apache.inlong.sort.base.Constants.NUM_RECORDS_IN_FOR_METER;
 import static org.apache.inlong.sort.base.Constants.NUM_RECORDS_IN_PER_SECOND;
 import static 
org.apache.inlong.sort.base.util.CalculateObjectSizeUtils.getDataSize;
 
-/**
- * A collection class for handling metrics
- */
-public class SourceMetricData implements MetricData, Serializable {
+public class SourceExactlyMetric implements MetricData, Serializable, 
SourceMetricsReporter {
 
     private static final long serialVersionUID = 1L;
-    private static final Logger LOG = 
LoggerFactory.getLogger(SourceMetricData.class);
     private MetricGroup metricGroup;
     private final Map<String, String> labels;
     private Counter numRecordsIn;
@@ -56,8 +52,10 @@ public class SourceMetricData implements MetricData, 
Serializable {
     private Counter numBytesInForMeter;
     private Meter numRecordsInPerSecond;
     private Meter numBytesInPerSecond;
-    private AuditOperator auditOperator;
+    private AuditReporterImpl auditReporter;
     private List<Integer> auditKeys;
+    private Long currentCheckpointId = 0L;
+    private Long lastCheckpointId = 0L;
 
     /**
      * currentFetchEventTimeLag = FetchTime - messageTimestamp, where the 
FetchTime is the time the
@@ -82,7 +80,7 @@ public class SourceMetricData implements MetricData, 
Serializable {
      */
     private volatile long emitDelay = 0L;
 
-    public SourceMetricData(MetricOption option, MetricGroup metricGroup) {
+    public SourceExactlyMetric(MetricOption option, MetricGroup metricGroup) {
         this.metricGroup = metricGroup;
         this.labels = option.getLabels();
 
@@ -104,90 +102,63 @@ public class SourceMetricData implements MetricData, 
Serializable {
         }
 
         if (option.getIpPorts().isPresent()) {
-            AuditOperator.getInstance().setAuditProxy(option.getIpPortSet());
-            this.auditOperator = AuditOperator.getInstance();
+            this.auditReporter = new AuditReporterImpl();
+            auditReporter.setAutoFlush(false);
+            auditReporter.setAuditProxy(option.getIpPortSet());
             this.auditKeys = option.getInlongAuditKeys();
         }
     }
 
-    public SourceMetricData(MetricOption option) {
+    public SourceExactlyMetric(MetricOption option) {
         this.labels = option.getLabels();
-
         if (option.getIpPorts().isPresent()) {
-            AuditOperator.getInstance().setAuditProxy(option.getIpPortSet());
-            this.auditOperator = AuditOperator.getInstance();
+            this.auditReporter = new AuditReporterImpl();
+            auditReporter.setAutoFlush(false);
+            auditReporter.setAuditProxy(option.getIpPortSet());
             this.auditKeys = option.getInlongAuditKeys();
         }
     }
 
     /**
-     * Default counter is {@link SimpleCounter}
-     * groupId and streamId and nodeId are label value, user can use it filter 
metric data when use metric reporter
-     * prometheus
+     * Users can custom counter that extends from {@link SimpleCounter}
+     * groupId and streamId and nodeId are label values,
+     * user can use it to filter metric data when using metric reporter 
Prometheus
+     * The following method is similar
      */
     public void registerMetricsForNumRecordsInForMeter() {
         registerMetricsForNumRecordsInForMeter(new SimpleCounter());
     }
 
     /**
-     * User can use custom counter that extends from {@link Counter}
-     * groupId and streamId and nodeId are label value, user can use it filter 
metric data when use metric reporter
-     * prometheus
+     * Users can custom counter that extends from {@link Counter}
+     * groupId and streamId and nodeId are label values,
+     * user can use it to filter metric data when using metric reporter 
Prometheus
+     * The following method is similar
      */
     public void registerMetricsForNumRecordsInForMeter(Counter counter) {
         numRecordsInForMeter = registerCounter(NUM_RECORDS_IN_FOR_METER, 
counter);
     }
 
-    /**
-     * Default counter is {@link SimpleCounter}
-     * groupId and streamId and nodeId are label value, user can use it filter 
metric data when use metric reporter
-     * prometheus
-     */
     public void registerMetricsForNumBytesInForMeter() {
         registerMetricsForNumBytesInForMeter(new SimpleCounter());
     }
 
-    /**
-     * User can use custom counter that extends from {@link Counter}
-     * groupId and streamId and nodeId are label value, user can use it filter 
metric data when use metric reporter
-     * prometheus
-     */
     public void registerMetricsForNumBytesInForMeter(Counter counter) {
         numBytesInForMeter = registerCounter(NUM_BYTES_IN_FOR_METER, counter);
     }
 
-    /**
-     * Default counter is {@link SimpleCounter}
-     * groupId and streamId and nodeId are label value, user can use it filter 
metric data when use metric reporter
-     * prometheus
-     */
     public void registerMetricsForNumRecordsIn() {
         registerMetricsForNumRecordsIn(new SimpleCounter());
     }
 
-    /**
-     * User can use custom counter that extends from {@link Counter}
-     * groupId and streamId and nodeId are label value, user can use it filter 
metric data when use metric reporter
-     * prometheus
-     */
     public void registerMetricsForNumRecordsIn(Counter counter) {
         numRecordsIn = registerCounter(NUM_RECORDS_IN, counter);
     }
 
-    /**
-     * Default counter is {@link SimpleCounter}
-     * groupId and streamId and nodeId are label value, user can use it filter 
metric data when use metric reporter
-     * prometheus
-     */
     public void registerMetricsForNumBytesIn() {
         registerMetricsForNumBytesIn(new SimpleCounter());
     }
 
-    /**
-     * User can use custom counter that extends from {@link Counter}
-     * groupId and streamId and nodeId are label value, user can use it filter 
metric data when use metric reporter
-     * prometheus
-     */
     public void registerMetricsForNumBytesIn(Counter counter) {
         numBytesIn = registerCounter(NUM_BYTES_IN, counter);
     }
@@ -250,73 +221,29 @@ public class SourceMetricData implements MetricData, 
Serializable {
         return labels;
     }
 
-    public void outputMetricsWithEstimate(Object data) {
-        outputMetrics(1, getDataSize(data));
-    }
-
-    public void outputMetricsWithEstimate(Object data, long fetchDelay, long 
emitDelay) {
-        outputMetrics(1, getDataSize(data));
-        this.fetchDelay = fetchDelay;
-        this.emitDelay = emitDelay;
-    }
-
+    @Override
     public void outputMetricsWithEstimate(Object data, long dataTime) {
         outputMetrics(1, getDataSize(data), dataTime);
     }
 
-    public void outputMetrics(long rowCountSize, long rowDataSize) {
-        outputDefaultMetrics(rowCountSize, rowDataSize);
-
-        if (auditOperator != null) {
-            for (Integer key : auditKeys) {
-                auditOperator.add(
-                        key,
-                        getGroupId(),
-                        getStreamId(),
-                        System.currentTimeMillis(),
-                        rowCountSize,
-                        rowDataSize);
-            }
-
-        }
-    }
-
-    public void outputMetrics(long rowCountSize, long rowDataSize, long 
fetchDelay, long emitDelay) {
-        outputDefaultMetrics(rowCountSize, rowDataSize, fetchDelay, emitDelay);
-
-        if (auditOperator != null) {
-            for (Integer key : auditKeys) {
-                auditOperator.add(
-                        key,
-                        getGroupId(),
-                        getStreamId(),
-                        System.currentTimeMillis(),
-                        rowCountSize,
-                        rowDataSize);
-            }
-
-        }
-    }
-
     public void outputMetrics(long rowCountSize, long rowDataSize, long 
dataTime) {
         outputDefaultMetrics(rowCountSize, rowDataSize);
-        if (auditOperator != null) {
+        if (auditReporter != null) {
             for (Integer key : auditKeys) {
-                auditOperator.add(
+                auditReporter.add(
+                        this.currentCheckpointId,
                         key,
+                        DEFAULT_AUDIT_TAG,
                         getGroupId(),
                         getStreamId(),
-                        getCurrentOrProvidedTime(dataTime),
+                        dataTime,
                         rowCountSize,
-                        rowDataSize);
+                        rowDataSize,
+                        DEFAULT_AUDIT_VERSION);
             }
         }
     }
 
-    private long getCurrentOrProvidedTime(long dataTime) {
-        return dataTime == 0 ? System.currentTimeMillis() : dataTime;
-    }
-
     private void outputDefaultMetrics(long rowCountSize, long rowDataSize) {
         if (numRecordsIn != null) {
             this.numRecordsIn.inc(rowCountSize);
@@ -339,16 +266,17 @@ public class SourceMetricData implements MetricData, 
Serializable {
      * flush audit data
      * usually call this method in close method or when checkpointing
      */
-    public void flushAuditData() {
-        if (auditOperator != null) {
-            auditOperator.flush();
+    public void flushAudit() {
+        if (auditReporter != null) {
+            auditReporter.flush(lastCheckpointId);
         }
     }
 
-    private void outputDefaultMetrics(long rowCountSize, long rowDataSize, 
long fetchDelay, long emitDelay) {
-        outputDefaultMetrics(rowCountSize, rowDataSize);
-        this.fetchDelay = fetchDelay;
-        this.emitDelay = emitDelay;
+    public void updateLastCheckpointId(Long checkpointId) {
+        lastCheckpointId = checkpointId;
+    }
+    public void updateCurrentCheckpointId(Long checkpointId) {
+        currentCheckpointId = checkpointId;
     }
 
     @Override
@@ -364,7 +292,7 @@ public class SourceMetricData implements MetricData, 
Serializable {
                 + ", numBytesInPerSecond=" + numBytesInPerSecond.getRate()
                 + ", currentFetchEventTimeLag=" + 
currentFetchEventTimeLag.getValue()
                 + ", currentEmitEventTimeLag=" + 
currentEmitEventTimeLag.getValue()
-                + ", auditOperator=" + auditOperator
+                + ", auditReporter=" + auditReporter
                 + '}';
     }
 }
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricData.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricData.java
index 0d2035e71c..0d84039a14 100644
--- 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricData.java
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricData.java
@@ -24,8 +24,6 @@ import org.apache.flink.metrics.Gauge;
 import org.apache.flink.metrics.Meter;
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.metrics.SimpleCounter;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
 
 import java.io.Serializable;
 import java.util.List;
@@ -44,10 +42,10 @@ import static 
org.apache.inlong.sort.base.util.CalculateObjectSizeUtils.getDataS
 /**
  * A collection class for handling metrics
  */
-public class SourceMetricData implements MetricData, Serializable {
+@Deprecated
+public class SourceMetricData implements MetricData, Serializable, 
SourceMetricsReporter {
 
     private static final long serialVersionUID = 1L;
-    private static final Logger LOG = 
LoggerFactory.getLogger(SourceMetricData.class);
     private MetricGroup metricGroup;
     private final Map<String, String> labels;
     private Counter numRecordsIn;
@@ -260,6 +258,7 @@ public class SourceMetricData implements MetricData, 
Serializable {
         this.emitDelay = emitDelay;
     }
 
+    @Override
     public void outputMetricsWithEstimate(Object data, long dataTime) {
         outputMetrics(1, getDataSize(data), dataTime);
     }
diff --git 
a/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricsReporter.java
 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricsReporter.java
new file mode 100644
index 0000000000..376a538d49
--- /dev/null
+++ 
b/inlong-sort/sort-flink/base/src/main/java/org/apache/inlong/sort/base/metric/SourceMetricsReporter.java
@@ -0,0 +1,24 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sort.base.metric;
+
+public interface SourceMetricsReporter {
+
+    void outputMetricsWithEstimate(Object data, long dataTime);
+
+}
diff --git 
a/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/FlinkTubeMQConsumer.java
 
b/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/FlinkTubeMQConsumer.java
index 1f261cfef5..47f17eb95d 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/FlinkTubeMQConsumer.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/FlinkTubeMQConsumer.java
@@ -20,6 +20,7 @@ package org.apache.inlong.sort.tubemq;
 import org.apache.inlong.sort.tubemq.table.DynamicTubeMQDeserializationSchema;
 import org.apache.inlong.sort.tubemq.table.TubeMQOptions;
 import org.apache.inlong.tubemq.client.config.ConsumerConfig;
+import org.apache.inlong.tubemq.client.consumer.ConsumeOffsetInfo;
 import org.apache.inlong.tubemq.client.consumer.ConsumePosition;
 import org.apache.inlong.tubemq.client.consumer.ConsumerResult;
 import org.apache.inlong.tubemq.client.consumer.PullMessageConsumer;
@@ -28,6 +29,7 @@ import org.apache.inlong.tubemq.corebase.Message;
 import org.apache.inlong.tubemq.corebase.TErrCodeConstants;
 
 import org.apache.flink.api.common.functions.util.ListCollector;
+import org.apache.flink.api.common.state.CheckpointListener;
 import org.apache.flink.api.common.state.ListState;
 import org.apache.flink.api.common.state.ListStateDescriptor;
 import org.apache.flink.api.common.state.OperatorStateStore;
@@ -64,7 +66,8 @@ import static org.apache.flink.util.TimeUtils.parseDuration;
  */
 public class FlinkTubeMQConsumer<T> extends RichParallelSourceFunction<T>
         implements
-            CheckpointedFunction {
+            CheckpointedFunction,
+            CheckpointListener {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(FlinkTubeMQConsumer.class);
     private static final String TUBE_OFFSET_STATE = "tube-offset-state";
@@ -257,18 +260,25 @@ public class FlinkTubeMQConsumer<T> extends 
RichParallelSourceFunction<T>
             lastConsumeInstant = Instant.now();
 
             List<T> records = new ArrayList<>();
-            lastConsumeInstant = getRecords(lastConsumeInstant, messageList, 
records);
 
             synchronized (ctx.getCheckpointLock()) {
-
+                lastConsumeInstant = getRecords(lastConsumeInstant, 
messageList, records);
                 for (T record : records) {
                     ctx.collect(record);
                 }
-                currentOffsets.put(
-                        consumeResult.getPartitionKey(),
-                        consumeResult.getCurrOffset());
-            }
 
+                ConsumerResult confirmResult = confirmOffset(consumeResult);
+                if (confirmResult != null && confirmResult.isSuccess()) {
+                    currentOffsets.put(confirmResult.getPartitionKey(), 
confirmResult.getCurrOffset());
+                } else {
+                    LOG.warn("Confirm offset failed, fallback to use offset in 
consume result.");
+                    currentOffsets.put(consumeResult.getPartitionKey(), 
consumeResult.getCurrOffset());
+                }
+            }
+        }
+    }
+    private ConsumerResult confirmOffset(ConsumerResult consumeResult) {
+        try {
             ConsumerResult confirmResult = messagePullConsumer
                     .confirmConsume(consumeResult.getConfirmContext(), true);
             if (!confirmResult.isSuccess()) {
@@ -283,9 +293,12 @@ public class FlinkTubeMQConsumer<T> extends 
RichParallelSourceFunction<T>
                             confirmResult.getErrMsg());
                 }
             }
+            return confirmResult;
+        } catch (Exception e) {
+            LOG.error("Confirm tube offset exception.", e);
         }
+        return null;
     }
-
     private Instant getRecords(Instant lastConsumeInstant, List<Message> 
messageList, List<T> records)
             throws Exception {
         if (messageList != null) {
@@ -311,14 +324,21 @@ public class FlinkTubeMQConsumer<T> extends 
RichParallelSourceFunction<T>
 
     @Override
     public void snapshotState(FunctionSnapshotContext context) throws 
Exception {
-
+        
deserializationSchema.setCurrentCheckpointId(context.getCheckpointId());
+        // Make sure snapshot all tube partitions' offset to avoid 
inconsistency.
+        Map<String, ConsumeOffsetInfo> curConsumedPartitions = 
messagePullConsumer.getCurConsumedPartitions();
+        for (Map.Entry<String, ConsumeOffsetInfo> consumedPartition : 
curConsumedPartitions.entrySet()) {
+            ConsumeOffsetInfo consumeOffsetInfo = consumedPartition.getValue();
+            if 
(!currentOffsets.containsKey(consumeOffsetInfo.getPartitionKey())) {
+                LOG.info("snapshot assigned but not consume partition: {}", 
consumedPartition.getValue());
+                currentOffsets.put(consumeOffsetInfo.getPartitionKey(), 
consumedPartition.getValue().getCurrOffset());
+            }
+        }
         offsetsState.clear();
         for (Map.Entry<String, Long> entry : currentOffsets.entrySet()) {
             offsetsState.add(new Tuple2<>(entry.getKey(), entry.getValue()));
         }
 
-        deserializationSchema.flushAudit();
-
         LOG.info("Successfully save the offsets in checkpoint {}: {}.",
                 context.getCheckpointId(), currentOffsets);
     }
@@ -353,4 +373,15 @@ public class FlinkTubeMQConsumer<T> extends 
RichParallelSourceFunction<T>
 
         LOG.info("Closed the tubemq source.");
     }
+
+    @Override
+    public void notifyCheckpointComplete(long checkpointId) throws Exception {
+        deserializationSchema.flushAudit();
+        deserializationSchema.updateLastCheckpointId(checkpointId);
+    }
+
+    @Override
+    public void notifyCheckpointAborted(long checkpointId) throws Exception {
+        CheckpointListener.super.notifyCheckpointAborted(checkpointId);
+    }
 }
diff --git 
a/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/table/DynamicTubeMQDeserializationSchema.java
 
b/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/table/DynamicTubeMQDeserializationSchema.java
index c6ec9ea9cb..80532a2dc1 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/table/DynamicTubeMQDeserializationSchema.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/table/DynamicTubeMQDeserializationSchema.java
@@ -29,7 +29,7 @@ import java.io.Serializable;
 public interface DynamicTubeMQDeserializationSchema<T> extends Serializable, 
ResultTypeQueryable<T> {
 
     @PublicEvolving
-    default void open() throws Exception {
+    default void open() {
     }
 
     /**
@@ -59,5 +59,9 @@ public interface DynamicTubeMQDeserializationSchema<T> 
extends Serializable, Res
         }
     }
 
+    void setCurrentCheckpointId(long checkpointId);
+
+    void updateLastCheckpointId(Long checkpointId);
+
     void flushAudit();
 }
diff --git 
a/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/table/DynamicTubeMQTableDeserializationSchema.java
 
b/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/table/DynamicTubeMQTableDeserializationSchema.java
index 3f2a57d7c7..5d4a3bd2c6 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/table/DynamicTubeMQTableDeserializationSchema.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.15/sort-connectors/tubemq/src/main/java/org/apache/inlong/sort/tubemq/table/DynamicTubeMQTableDeserializationSchema.java
@@ -19,7 +19,7 @@ package org.apache.inlong.sort.tubemq.table;
 
 import org.apache.inlong.sort.base.metric.MetricOption;
 import org.apache.inlong.sort.base.metric.MetricsCollector;
-import org.apache.inlong.sort.base.metric.SourceMetricData;
+import org.apache.inlong.sort.base.metric.SourceExactlyMetric;
 import org.apache.inlong.tubemq.corebase.Message;
 
 import com.google.common.base.Objects;
@@ -61,9 +61,9 @@ public class DynamicTubeMQTableDeserializationSchema 
implements DynamicTubeMQDes
 
     private final boolean innerFormat;
 
-    private SourceMetricData sourceMetricData;
+    private SourceExactlyMetric sourceExactlyMetric;
 
-    private MetricOption metricOption;
+    private final MetricOption metricOption;
 
     public DynamicTubeMQTableDeserializationSchema(
             DeserializationSchema<RowData> schema,
@@ -83,7 +83,7 @@ public class DynamicTubeMQTableDeserializationSchema 
implements DynamicTubeMQDes
     @Override
     public void open() {
         if (metricOption != null) {
-            sourceMetricData = new SourceMetricData(metricOption);
+            sourceExactlyMetric = new SourceExactlyMetric(metricOption);
         }
     }
 
@@ -97,7 +97,7 @@ public class DynamicTubeMQTableDeserializationSchema 
implements DynamicTubeMQDes
         List<RowData> rows = new ArrayList<>();
 
         MetricsCollector<RowData> metricsCollector =
-                new MetricsCollector<>(new ListCollector<>(rows), 
sourceMetricData);
+                new MetricsCollector<>(new ListCollector<>(rows), 
sourceExactlyMetric);
 
         // reset time stamp if the deserialize schema has not inner format
         if (!innerFormat) {
@@ -111,8 +111,21 @@ public class DynamicTubeMQTableDeserializationSchema 
implements DynamicTubeMQDes
 
     @Override
     public void flushAudit() {
-        if (sourceMetricData != null) {
-            sourceMetricData.flushAuditData();
+        if (sourceExactlyMetric != null) {
+            sourceExactlyMetric.flushAudit();
+        }
+    }
+    @Override
+    public void setCurrentCheckpointId(long checkpointId) {
+        if (sourceExactlyMetric != null) {
+            sourceExactlyMetric.updateCurrentCheckpointId(checkpointId);
+        }
+    }
+
+    @Override
+    public void updateLastCheckpointId(Long checkpointId) {
+        if (sourceExactlyMetric != null) {
+            sourceExactlyMetric.updateLastCheckpointId(checkpointId);
         }
     }
 

Reply via email to