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