This is an automated email from the ASF dual-hosted git repository.

gosonzhang 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 0bf55ecee4 [INLONG-8597][DataProxy] Adjust the format of the metric 
output to the file (#8600)
0bf55ecee4 is described below

commit 0bf55ecee462d5d32b9718d58be66b4b299be171
Author: Goson Zhang <[email protected]>
AuthorDate: Thu Jul 27 18:50:31 2023 +0800

    [INLONG-8597][DataProxy] Adjust the format of the metric output to the file 
(#8600)
---
 .../inlong/dataproxy/consts/ConfigConstants.java   |  1 +
 .../dataproxy/metrics/stats/MonitorIndex.java      |  5 +-
 .../inlong/dataproxy/sink/common/SinkContext.java  | 59 +++++++++++++++++++---
 .../inlong/dataproxy/sink/common/TubeUtils.java    | 13 ++---
 .../sink/mq/MessageQueueZoneSinkContext.java       | 42 ---------------
 .../dataproxy/sink/mq/SimplePackProfile.java       |  4 +-
 .../dataproxy/sink/mq/kafka/KafkaHandler.java      | 22 ++++----
 .../dataproxy/sink/mq/pulsar/PulsarHandler.java    | 17 +++----
 .../inlong/dataproxy/sink/mq/tube/TubeHandler.java | 26 +++++-----
 .../apache/inlong/dataproxy/source/BaseSource.java | 46 ++++++++++++++---
 .../dataproxy/source/ServerMessageHandler.java     | 25 +++------
 .../source/httpMsg/HttpMessageHandler.java         | 25 +++------
 .../dataproxy/source/v0msg/AbsV0MsgCodec.java      | 11 ++--
 .../inlong/dataproxy/source/v0msg/CodecBinMsg.java |  3 +-
 .../dataproxy/source/v0msg/CodecTextMsg.java       |  3 +-
 .../inlong/dataproxy/utils/DateTimeUtils.java      | 14 +++++
 16 files changed, 167 insertions(+), 149 deletions(-)

diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/consts/ConfigConstants.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/consts/ConfigConstants.java
index ea8f51cee1..47f2e9fcb9 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/consts/ConfigConstants.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/consts/ConfigConstants.java
@@ -78,6 +78,7 @@ public class ConfigConstants {
     public static final String REMOTE_IP_KEY = "srcIp";
     public static final String DATAPROXY_IP_KEY = "dpIp";
     public static final String MSG_ENCODE_VER = "msgEnType";
+    public static final String MSG_SEND_TIME = "st";
     public static final String REMOTE_IDC_KEY = "idc";
     public static final String MSG_COUNTER_KEY = "msgcnt";
     public static final String PKG_COUNTER_KEY = "pkgcnt";
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/metrics/stats/MonitorIndex.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/metrics/stats/MonitorIndex.java
index 1aaa7c0ddb..9f21707554 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/metrics/stats/MonitorIndex.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/metrics/stats/MonitorIndex.java
@@ -34,7 +34,8 @@ import java.util.concurrent.atomic.LongAdder;
 public class MonitorIndex extends AbsStatsDaemon {
 
     private static final Logger LOGGER = 
LoggerFactory.getLogger(MonitorIndex.class);
-
+    // Indicator record format version
+    private static final String INDEX_RECORD_VER = "v1";
     private static final AtomicLong RECODE_ID = new AtomicLong(0);
     private final StatsUnit[] statsUnits = new StatsUnit[2];
 
@@ -134,7 +135,7 @@ public class MonitorIndex extends AbsStatsDaemon {
                 if (entry == null || entry.getKey() == null || 
entry.getValue() == null) {
                     continue;
                 }
-                LOGGER.info("{}#{}_{}#{}#{}", this.statsName, printTime,
+                LOGGER.info("{}#{}#{}_{}#{}={}", this.statsName, 
INDEX_RECORD_VER, printTime,
                         RECODE_ID.incrementAndGet(), entry.getKey(), 
entry.getValue().toString());
                 printCnt++;
             }
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/common/SinkContext.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/common/SinkContext.java
index 9ccdb82a05..61e901dddd 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/common/SinkContext.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/common/SinkContext.java
@@ -18,16 +18,23 @@
 package org.apache.inlong.dataproxy.sink.common;
 
 import org.apache.inlong.common.metric.MetricRegister;
+import org.apache.inlong.common.msg.AttributeConstants;
 import org.apache.inlong.dataproxy.config.CommonConfigHolder;
 import org.apache.inlong.dataproxy.config.pojo.CacheClusterConfig;
 import org.apache.inlong.dataproxy.consts.AttrConstants;
+import org.apache.inlong.dataproxy.consts.ConfigConstants;
+import org.apache.inlong.dataproxy.consts.StatConstants;
 import org.apache.inlong.dataproxy.metrics.DataProxyMetricItemSet;
 import org.apache.inlong.dataproxy.metrics.stats.MonitorIndex;
 import org.apache.inlong.dataproxy.metrics.stats.MonitorStats;
 import org.apache.inlong.dataproxy.sink.mq.MessageQueueHandler;
+import org.apache.inlong.dataproxy.sink.mq.PackProfile;
+import org.apache.inlong.dataproxy.sink.mq.SimplePackProfile;
 import org.apache.inlong.dataproxy.sink.mq.pulsar.PulsarHandler;
+import org.apache.inlong.dataproxy.utils.DateTimeUtils;
 
 import org.apache.commons.lang.ClassUtils;
+import org.apache.commons.lang3.math.NumberUtils;
 import org.apache.flume.Channel;
 import org.apache.flume.Context;
 import org.slf4j.Logger;
@@ -123,15 +130,55 @@ public class SinkContext {
         }
     }
 
-    public void fileMetricAddSuccCnt(String key, int cnt, int packCnt, long 
packSize) {
-        if (CommonConfigHolder.getInstance().isEnableFileMetric()) {
-            monitorIndex.addSuccStats(key, cnt, packCnt, packSize);
+    public void fileMetricAddSuccStats(PackProfile profile, String topic, 
String brokerIP) {
+        if (!CommonConfigHolder.getInstance().isEnableFileMetric()
+                || !(profile instanceof SimplePackProfile)) {
+            return;
         }
+        fileMetricIncStats((SimplePackProfile) profile, true,
+                topic, brokerIP, StatConstants.EVENT_SINK_SUCCESS, "");
     }
 
-    public void fileMetricAddFailCnt(String key, int failCnt) {
-        if (CommonConfigHolder.getInstance().isEnableFileMetric()) {
-            monitorIndex.addFailStats(key, failCnt);
+    public void fileMetricAddFailStats(PackProfile profile, String topic, 
String brokerIP, String detailKey) {
+        if (!CommonConfigHolder.getInstance().isEnableFileMetric()
+                || !(profile instanceof SimplePackProfile)) {
+            return;
+        }
+        fileMetricIncStats((SimplePackProfile) profile, false,
+                topic, brokerIP, StatConstants.EVENT_SINK_FAILURE, detailKey);
+    }
+
+    public void fileMetricAddExceptStats(PackProfile profile, String topic, 
String brokerIP, String detailKey) {
+        if (!CommonConfigHolder.getInstance().isEnableFileMetric()
+                || !(profile instanceof SimplePackProfile)) {
+            return;
+        }
+        fileMetricIncStats((SimplePackProfile) profile, false,
+                topic, brokerIP, StatConstants.EVENT_SINK_RECEIVEEXCEPT, 
detailKey);
+    }
+
+    private void fileMetricIncStats(SimplePackProfile profile, boolean isSucc,
+            String topic, String brokerIP, String eventKey, String 
detailInfoKey) {
+        long dtL = 
Long.parseLong(profile.getProperties().get(AttributeConstants.DATA_TIME));
+        long pkgTimeL = 
Long.parseLong(profile.getProperties().get(ConfigConstants.PKG_TIME_KEY));
+        StringBuilder statsKey = new StringBuilder(512)
+                .append(sinkName).append(AttrConstants.SEP_HASHTAG)
+                
.append(profile.getInlongGroupId()).append(AttrConstants.SEP_HASHTAG)
+                
.append(profile.getInlongStreamId()).append(AttrConstants.SEP_HASHTAG)
+                .append(topic).append(AttrConstants.SEP_HASHTAG)
+                
.append(profile.getProperties().get(ConfigConstants.DATAPROXY_IP_KEY)).append(AttrConstants.SEP_HASHTAG)
+                .append(brokerIP).append(AttrConstants.SEP_HASHTAG)
+                
.append(DateTimeUtils.ms2yyyyMMddHHmmTenMins(dtL)).append(AttrConstants.SEP_HASHTAG)
+                .append(DateTimeUtils.ms2yyyyMMddHHmm(pkgTimeL));
+        if (isSucc) {
+            monitorIndex.addSuccStats(statsKey.toString(), NumberUtils.toInt(
+                    
profile.getProperties().get(ConfigConstants.MSG_COUNTER_KEY), 1),
+                    1, profile.getSize());
+            monitorStats.incSumStats(eventKey);
+        } else {
+            monitorIndex.addFailStats(statsKey.toString(), 1);
+            monitorStats.incSumStats(eventKey);
+            monitorStats.incDetailStats(eventKey + "#" + detailInfoKey);
         }
     }
 
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/common/TubeUtils.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/common/TubeUtils.java
index 9bd682429b..fb9e0b26e3 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/common/TubeUtils.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/common/TubeUtils.java
@@ -17,8 +17,6 @@
 
 package org.apache.inlong.dataproxy.sink.common;
 
-import org.apache.inlong.common.enums.DataProxyMsgEncType;
-import org.apache.inlong.common.msg.AttributeConstants;
 import org.apache.inlong.dataproxy.config.pojo.MQClusterConfig;
 import org.apache.inlong.dataproxy.consts.ConfigConstants;
 import org.apache.inlong.dataproxy.utils.Constants;
@@ -62,14 +60,9 @@ public class TubeUtils {
         Map<String, String> headers = event.getHeaders();
         Message message = new Message(topicName, event.getBody());
         String pkgVersion = headers.get(ConfigConstants.MSG_ENCODE_VER);
-        if 
(DataProxyMsgEncType.MSG_ENCODE_TYPE_PB.getStrId().equalsIgnoreCase(pkgVersion))
 {
-            long dataTimeL = 
Long.parseLong(headers.get(ConfigConstants.PKG_TIME_KEY));
-            message.putSystemHeader(headers.get(Constants.INLONG_STREAM_ID),
-                    DateTimeUtils.ms2yyyyMMddHHmm(dataTimeL));
-        } else {
-            message.putSystemHeader(headers.get(AttributeConstants.STREAM_ID),
-                    headers.get(ConfigConstants.PKG_TIME_KEY));
-        }
+        long dataTimeL = 
Long.parseLong(headers.get(ConfigConstants.PKG_TIME_KEY));
+        message.putSystemHeader(headers.get(Constants.INLONG_STREAM_ID),
+                DateTimeUtils.ms2yyyyMMddHHmm(dataTimeL));
         Map<String, String> extraAttrMap = MessageUtils.getXfsAttrs(headers, 
pkgVersion);
         for (Map.Entry<String, String> entry : extraAttrMap.entrySet()) {
             message.setAttrKeyVal(entry.getKey(), entry.getValue());
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSinkContext.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSinkContext.java
index 6989c94ea0..33168baf8c 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSinkContext.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSinkContext.java
@@ -18,10 +18,7 @@
 package org.apache.inlong.dataproxy.sink.mq;
 
 import org.apache.inlong.common.enums.DataProxyErrCode;
-import org.apache.inlong.common.util.NetworkUtils;
 import org.apache.inlong.dataproxy.config.CommonConfigHolder;
-import org.apache.inlong.dataproxy.consts.AttrConstants;
-import org.apache.inlong.dataproxy.consts.ConfigConstants;
 import org.apache.inlong.dataproxy.consts.StatConstants;
 import org.apache.inlong.dataproxy.metrics.DataProxyMetricItem;
 import org.apache.inlong.dataproxy.metrics.audit.AuditUtils;
@@ -30,7 +27,6 @@ import 
org.apache.inlong.sdk.commons.protocol.ProxySdk.INLONG_COMPRESSED_TYPE;
 
 import org.apache.commons.lang.ClassUtils;
 import org.apache.commons.lang3.StringUtils;
-import org.apache.commons.lang3.math.NumberUtils;
 import org.apache.flume.Channel;
 import org.apache.flume.Context;
 import org.apache.flume.conf.Configurable;
@@ -274,42 +270,4 @@ public class MessageQueueZoneSinkContext extends 
SinkContext {
         return null;
     }
 
-    public void fileMetricAddSuccCnt(PackProfile packProfile, String topic, 
String remoteId) {
-        if (!CommonConfigHolder.getInstance().isEnableFileMetric()) {
-            return;
-        }
-        if (packProfile instanceof SimplePackProfile) {
-            SimplePackProfile simpleProfile = (SimplePackProfile) packProfile;
-            StringBuilder statsKey = new StringBuilder(512)
-                    .append(sinkName).append(AttrConstants.SEP_HASHTAG)
-                    
.append(simpleProfile.getInlongGroupId()).append(AttrConstants.SEP_HASHTAG)
-                    
.append(simpleProfile.getInlongStreamId()).append(AttrConstants.SEP_HASHTAG)
-                    .append(topic).append(AttrConstants.SEP_HASHTAG)
-                    
.append(NetworkUtils.getLocalIp()).append(AttrConstants.SEP_HASHTAG)
-                    .append(remoteId).append(AttrConstants.SEP_HASHTAG)
-                    
.append(simpleProfile.getProperties().get(ConfigConstants.PKG_TIME_KEY));
-            monitorIndex.addSuccStats(statsKey.toString(), NumberUtils.toInt(
-                    
simpleProfile.getProperties().get(ConfigConstants.MSG_COUNTER_KEY), 1),
-                    1, simpleProfile.getSize());
-        }
-    }
-
-    public void fileMetricAddFailCnt(PackProfile packProfile, String topic, 
String remoteId) {
-        if (!CommonConfigHolder.getInstance().isEnableFileMetric()) {
-            return;
-        }
-
-        if (packProfile instanceof SimplePackProfile) {
-            SimplePackProfile simpleProfile = (SimplePackProfile) packProfile;
-            StringBuilder statsKey = new StringBuilder(512)
-                    .append(sinkName).append(AttrConstants.SEP_HASHTAG)
-                    
.append(simpleProfile.getInlongGroupId()).append(AttrConstants.SEP_HASHTAG)
-                    
.append(simpleProfile.getInlongStreamId()).append(AttrConstants.SEP_HASHTAG)
-                    .append(topic).append(AttrConstants.SEP_HASHTAG)
-                    
.append(NetworkUtils.getLocalIp()).append(AttrConstants.SEP_HASHTAG)
-                    .append(remoteId).append(AttrConstants.SEP_HASHTAG)
-                    
.append(simpleProfile.getProperties().get(ConfigConstants.PKG_TIME_KEY));
-            monitorIndex.addFailStats(statsKey.toString(), 1);
-        }
-    }
 }
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/SimplePackProfile.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/SimplePackProfile.java
index 242042512c..6d8ecf044a 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/SimplePackProfile.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/SimplePackProfile.java
@@ -148,11 +148,13 @@ public class SimplePackProfile extends PackProfile {
     /**
      * get required properties sent to MQ
      *
+     * @param sendTime  send time
      * @return the properties
      */
-    public Map<String, String> getPropsToMQ() {
+    public Map<String, String> getPropsToMQ(long sendTime) {
         Map<String, String> result = new HashMap<>();
         result.put(AttributeConstants.RCV_TIME, 
event.getHeaders().get(AttributeConstants.RCV_TIME));
+        result.put(ConfigConstants.MSG_SEND_TIME, String.valueOf(sendTime));
         result.put(ConfigConstants.MSG_ENCODE_VER, 
event.getHeaders().get(ConfigConstants.MSG_ENCODE_VER));
         result.put(EventConstants.HEADER_KEY_VERSION, 
event.getHeaders().get(EventConstants.HEADER_KEY_VERSION));
         result.put(ConfigConstants.REMOTE_IP_KEY, 
event.getHeaders().get(ConfigConstants.REMOTE_IP_KEY));
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/kafka/KafkaHandler.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/kafka/KafkaHandler.java
index 3d7303473e..e10ff68dd9 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/kafka/KafkaHandler.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/kafka/KafkaHandler.java
@@ -223,18 +223,16 @@ public class KafkaHandler implements MessageQueueHandler {
      */
     private void sendSimplePackProfile(SimplePackProfile simpleProfile, 
IdTopicConfig idConfig,
             String topic) throws Exception {
-        // headers
-        Map<String, String> headers = simpleProfile.getPropsToMQ();
-        // body
-        byte[] bodyBytes = simpleProfile.getEvent().getBody();
         // metric
-        sinkContext.addSendMetric(simpleProfile, clusterName, topic, 
bodyBytes.length);
+        sinkContext.addSendMetric(simpleProfile, clusterName,
+                topic, simpleProfile.getEvent().getBody().length);
+        // prepare ProducerRecord
+        ProducerRecord<String, byte[]> producerRecord =
+                new ProducerRecord<>(topic, 
simpleProfile.getEvent().getBody());
         // sendAsync
         long sendTime = System.currentTimeMillis();
-
-        // prepare ProducerRecord
-        ProducerRecord<String, byte[]> producerRecord = new 
ProducerRecord<>(topic, bodyBytes);
         // add headers
+        Map<String, String> headers = simpleProfile.getPropsToMQ(sendTime);
         headers.forEach((key, value) -> {
             producerRecord.headers().add(key, value.getBytes());
         });
@@ -245,18 +243,16 @@ public class KafkaHandler implements MessageQueueHandler {
             @Override
             public void onCompletion(RecordMetadata arg0, Exception ex) {
                 if (ex != null) {
-                    
sinkContext.fileMetricIncWithDetailStats(StatConstants.EVENT_SINK_FAILURE, 
topic);
-                    sinkContext.fileMetricAddFailCnt(simpleProfile, topic,
-                            arg0 == null ? "" : 
String.valueOf(arg0.partition()));
+                    sinkContext.fileMetricAddFailStats(simpleProfile, topic,
+                            arg0 == null ? "" : 
String.valueOf(arg0.partition()), topic);
                     sinkContext.processSendFail(simpleProfile, clusterName, 
topic, sendTime,
                             DataProxyErrCode.MQ_RETURN_ERROR, ex.getMessage());
                     if (logCounter.shouldPrint()) {
                         logger.error("Send SimplePackProfile to Kafka 
failure", ex);
                     }
                 } else {
-                    sinkContext.fileMetricAddSuccCnt(simpleProfile, topic,
+                    sinkContext.fileMetricAddSuccStats(simpleProfile, topic,
                             arg0 == null ? "" : 
String.valueOf(arg0.partition()));
-                    
sinkContext.fileMetricIncSumStats(StatConstants.EVENT_SINK_SUCCESS);
                     sinkContext.addSendResultMetric(simpleProfile, 
clusterName, topic, true, sendTime);
                     
sinkContext.getMqZoneSink().releaseAcquiredSizePermit(simpleProfile);
                     simpleProfile.ack();
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java
index 149f3e5e1b..1b76eaf72a 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java
@@ -344,29 +344,24 @@ public class PulsarHandler implements MessageQueueHandler 
{
     private void sendSimplePackProfile(SimplePackProfile simpleProfile, 
IdTopicConfig idConfig,
             Producer<byte[]> producer,
             String producerTopic) throws Exception {
-        // headers
-        Map<String, String> headers = simpleProfile.getPropsToMQ();
-        // body
-        byte[] bodyBytes = simpleProfile.getEvent().getBody();
         // metric
-        sinkContext.addSendMetric(simpleProfile, clusterName, producerTopic, 
bodyBytes.length);
+        sinkContext.addSendMetric(simpleProfile, clusterName,
+                producerTopic, simpleProfile.getEvent().getBody().length);
         // sendAsync
         long sendTime = System.currentTimeMillis();
-        CompletableFuture<MessageId> future = 
producer.newMessage().properties(headers)
-                .value(bodyBytes).sendAsync();
+        CompletableFuture<MessageId> future = producer.newMessage().properties(
+                
simpleProfile.getPropsToMQ(sendTime)).value(simpleProfile.getEvent().getBody()).sendAsync();
         // callback
         future.whenCompleteAsync((msgId, ex) -> {
             if (ex != null) {
-                
sinkContext.fileMetricIncWithDetailStats(StatConstants.EVENT_SINK_FAILURE, 
producerTopic);
-                sinkContext.fileMetricAddFailCnt(simpleProfile, producerTopic, 
msgId.toString());
+                sinkContext.fileMetricAddFailStats(simpleProfile, 
producerTopic, "", producerTopic);
                 sinkContext.processSendFail(simpleProfile, clusterName, 
producerTopic, sendTime,
                         DataProxyErrCode.MQ_RETURN_ERROR, ex.getMessage());
                 if (logCounter.shouldPrint()) {
                     logger.error("Send SimpleProfileV0 to Pulsar failure", ex);
                 }
             } else {
-                sinkContext.fileMetricAddSuccCnt(simpleProfile, producerTopic, 
msgId.toString());
-                
sinkContext.fileMetricIncSumStats(StatConstants.EVENT_SINK_SUCCESS);
+                sinkContext.fileMetricAddSuccStats(simpleProfile, 
producerTopic, "");
                 sinkContext.addSendResultMetric(simpleProfile, clusterName, 
producerTopic, true, sendTime);
                 
sinkContext.getMqZoneSink().releaseAcquiredSizePermit(simpleProfile);
                 simpleProfile.ack();
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/tube/TubeHandler.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/tube/TubeHandler.java
index c5c9dcbcfb..b363f79804 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/tube/TubeHandler.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/tube/TubeHandler.java
@@ -300,31 +300,30 @@ public class TubeHandler implements MessageQueueHandler {
      */
     private void sendSimplePackProfile(SimplePackProfile simpleProfile, 
IdTopicConfig idConfig,
             String topic) throws Exception {
+        // metric
+        sinkContext.addSendMetric(simpleProfile, clusterName,
+                topic, simpleProfile.getEvent().getBody().length);
         // build message
         Message message = new Message(topic, 
simpleProfile.getEvent().getBody());
-        message.putSystemHeader(simpleProfile.getInlongStreamId(),
-                
simpleProfile.getProperties().get(ConfigConstants.PKG_TIME_KEY));
+        long dataTimeL = 
Long.parseLong(simpleProfile.getProperties().get(ConfigConstants.PKG_TIME_KEY));
+        message.putSystemHeader(simpleProfile.getInlongStreamId(), 
DateTimeUtils.ms2yyyyMMddHHmm(dataTimeL));
         // add headers
-        Map<String, String> headers = simpleProfile.getPropsToMQ();
+        long sendTime = System.currentTimeMillis();
+        Map<String, String> headers = simpleProfile.getPropsToMQ(sendTime);
         headers.forEach(message::setAttrKeyVal);
-        // metric
-        sinkContext.addSendMetric(simpleProfile, clusterName, topic, 
simpleProfile.getEvent().getBody().length);
         // callback
-        long sendTime = System.currentTimeMillis();
         MessageSentCallback callback = new MessageSentCallback() {
 
             @Override
             public void onMessageSent(MessageSentResult result) {
                 if (result.isSuccess()) {
-                    sinkContext.fileMetricAddSuccCnt(simpleProfile, topic, 
result.getPartition().getHost());
-                    
sinkContext.fileMetricIncSumStats(StatConstants.EVENT_SINK_SUCCESS);
+                    sinkContext.fileMetricAddSuccStats(simpleProfile, topic, 
result.getPartition().getHost());
                     sinkContext.addSendResultMetric(simpleProfile, 
clusterName, topic, true, sendTime);
                     
sinkContext.getMqZoneSink().releaseAcquiredSizePermit(simpleProfile);
                     simpleProfile.ack();
                 } else {
-                    
sinkContext.fileMetricIncWithDetailStats(StatConstants.EVENT_SINK_FAILURE,
-                            topic + "." + result.getErrCode());
-                    sinkContext.fileMetricAddFailCnt(simpleProfile, topic, 
result.getPartition().getHost());
+                    sinkContext.fileMetricAddFailStats(simpleProfile, topic,
+                            result.getPartition().getHost(), topic + "." + 
result.getErrCode());
                     sinkContext.processSendFail(simpleProfile, clusterName, 
topic, sendTime,
                             DataProxyErrCode.MQ_RETURN_ERROR, 
result.getErrMsg());
                     if (logCounter.shouldPrint()) {
@@ -335,12 +334,11 @@ public class TubeHandler implements MessageQueueHandler {
 
             @Override
             public void onException(Throwable ex) {
-                
sinkContext.fileMetricIncSumStats(StatConstants.EVENT_SINK_RECEIVEEXCEPT);
-                sinkContext.fileMetricAddFailCnt(simpleProfile, topic, "");
+                sinkContext.fileMetricAddExceptStats(simpleProfile, topic, "", 
topic);
                 sinkContext.processSendFail(simpleProfile, clusterName, topic, 
sendTime,
                         DataProxyErrCode.MQ_RETURN_ERROR, ex.getMessage());
                 if (logCounter.shouldPrint()) {
-                    logger.error("Send SimpleProfileV0 to tube exception", ex);
+                    logger.error("Send Message to {} tube exception", topic, 
ex);
                 }
             }
         };
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/BaseSource.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/BaseSource.java
index 3917d22b83..e93c43910b 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/BaseSource.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/BaseSource.java
@@ -24,6 +24,7 @@ import org.apache.inlong.dataproxy.config.CommonConfigHolder;
 import org.apache.inlong.dataproxy.config.ConfigManager;
 import org.apache.inlong.dataproxy.config.holder.ConfigUpdateCallback;
 import org.apache.inlong.dataproxy.consts.AttrConstants;
+import org.apache.inlong.dataproxy.consts.StatConstants;
 import org.apache.inlong.dataproxy.metrics.DataProxyMetricItem;
 import org.apache.inlong.dataproxy.metrics.DataProxyMetricItemSet;
 import org.apache.inlong.dataproxy.metrics.audit.AuditUtils;
@@ -32,6 +33,7 @@ import org.apache.inlong.dataproxy.metrics.stats.MonitorStats;
 import org.apache.inlong.dataproxy.source.httpMsg.HttpMessageHandler;
 import org.apache.inlong.dataproxy.utils.AddressUtils;
 import org.apache.inlong.dataproxy.utils.ConfStringUtils;
+import org.apache.inlong.dataproxy.utils.DateTimeUtils;
 import org.apache.inlong.dataproxy.utils.FailoverChannelProcessorHolder;
 import org.apache.inlong.sdk.commons.admin.AdminServiceRegister;
 
@@ -256,7 +258,7 @@ public abstract class BaseSource
             this.monitorIndex.start();
             this.monitorStats = new MonitorStats(
                     
CommonConfigHolder.getInstance().getFileMetricEventOutName()
-                            + AttrConstants.SEP_HASHTAG + 
this.getProtocolName(),
+                            + AttrConstants.SEP_HASHTAG + this.getName(),
                     
CommonConfigHolder.getInstance().getFileMetricStatInvlSec() * 1000L,
                     
CommonConfigHolder.getInstance().getFileMetricStatCacheCnt());
             this.monitorStats.start();
@@ -413,17 +415,45 @@ public abstract class BaseSource
         }
     }
 
-    public void fileMetricAddSuccCnt(String key, int cnt, int packCnt, long 
packSize) {
-        if (CommonConfigHolder.getInstance().isEnableFileMetric()) {
-            monitorIndex.addSuccStats(key, cnt, packCnt, packSize);
-        }
+    public void fileMetricAddSuccStats(StringBuilder strBuff, String groupId, 
String streamId,
+            String topicName, String clientIP, String msgProcType,
+            long dt, long pkgTime, int cnt, int packCnt, long packSize) {
+        fileMetricIncStats(strBuff, true, groupId, streamId, topicName,
+                clientIP, msgProcType, dt, pkgTime, cnt, packCnt, packSize, 0);
     }
 
-    public void fileMetricAddFailCnt(String key, int failCnt) {
-        if (CommonConfigHolder.getInstance().isEnableFileMetric()) {
-            monitorIndex.addFailStats(key, failCnt);
+    public void fileMetricAddFailStats(StringBuilder strBuff, String groupId, 
String streamId,
+            String topicName, String clientIP, String msgProcType, long dt, 
long pkgTime, int failCnt) {
+        fileMetricIncStats(strBuff, false, groupId, streamId, topicName,
+                clientIP, msgProcType, dt, pkgTime, 0, 0, 0, failCnt);
+    }
+
+    private void fileMetricIncStats(StringBuilder strBuff, boolean isSucc, 
String groupId,
+            String streamId, String topicName, String clientIP, String 
msgProcType,
+            long dt, long pkgTime, int cnt, int packCnt, long packSize, int 
failCnt) {
+        if (!CommonConfigHolder.getInstance().isEnableFileMetric()) {
+            return;
+        }
+        String tenMinsDt = DateTimeUtils.ms2yyyyMMddHHmmTenMins(dt);
+        strBuff.append(getName()).append(AttrConstants.SEP_HASHTAG)
+                .append(groupId).append(AttrConstants.SEP_HASHTAG)
+                .append(streamId).append(AttrConstants.SEP_HASHTAG)
+                .append(topicName).append(AttrConstants.SEP_HASHTAG)
+                .append(msgProcType).append(AttrConstants.SEP_HASHTAG)
+                .append(srcHost).append(AttrConstants.SEP_HASHTAG)
+                .append(clientIP).append(AttrConstants.SEP_HASHTAG)
+                .append(tenMinsDt).append(AttrConstants.SEP_HASHTAG)
+                .append(DateTimeUtils.ms2yyyyMMddHHmm(pkgTime));
+        if (isSucc) {
+            monitorStats.incSumStats(StatConstants.EVENT_MSG_V0_POST_SUCCESS);
+            monitorIndex.addSuccStats(strBuff.toString(), cnt, packCnt, 
packSize);
+        } else {
+            monitorIndex.addFailStats(strBuff.toString(), failCnt);
+            monitorStats.incSumStats(StatConstants.EVENT_MSG_V0_POST_FAILURE);
         }
+        strBuff.delete(0, strBuff.length());
     }
+
     /**
      * addMetric
      *
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/ServerMessageHandler.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/ServerMessageHandler.java
index fc8e15d048..7dd19dc157 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/ServerMessageHandler.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/ServerMessageHandler.java
@@ -23,7 +23,6 @@ import org.apache.inlong.common.msg.AttributeConstants;
 import org.apache.inlong.common.msg.MsgType;
 import org.apache.inlong.dataproxy.config.CommonConfigHolder;
 import org.apache.inlong.dataproxy.config.ConfigManager;
-import org.apache.inlong.dataproxy.consts.AttrConstants;
 import org.apache.inlong.dataproxy.consts.ConfigConstants;
 import org.apache.inlong.dataproxy.consts.StatConstants;
 import org.apache.inlong.dataproxy.source.v0msg.AbsV0MsgCodec;
@@ -31,7 +30,6 @@ import org.apache.inlong.dataproxy.source.v0msg.CodecBinMsg;
 import org.apache.inlong.dataproxy.source.v0msg.CodecTextMsg;
 import org.apache.inlong.dataproxy.source.v1msg.InlongTcpSourceCallback;
 import org.apache.inlong.dataproxy.utils.AddressUtils;
-import org.apache.inlong.dataproxy.utils.DateTimeUtils;
 import org.apache.inlong.sdk.commons.protocol.EventUtils;
 import org.apache.inlong.sdk.commons.protocol.ProxyEvent;
 import org.apache.inlong.sdk.commons.protocol.ProxyPackEvent;
@@ -271,30 +269,21 @@ public class ServerMessageHandler extends 
ChannelInboundHandlerAdapter {
         }
         // build InLong event.
         Event event = msgCodec.encEventPackage(source, channel);
-        // build metric data item
-        long longDataTime = msgCodec.getDataTimeMs() / 1000 / 60 / 10;
-        longDataTime = longDataTime * 1000 * 60 * 10;
-        String statsKey = 
strBuff.append(source.getProtocolName()).append(AttrConstants.SEP_HASHTAG)
-                
.append(msgCodec.getGroupId()).append(AttrConstants.SEP_HASHTAG)
-                
.append(msgCodec.getStreamId()).append(AttrConstants.SEP_HASHTAG)
-                
.append(msgCodec.getStrRemoteIP()).append(AttrConstants.SEP_HASHTAG)
-                .append(source.getSrcHost()).append(AttrConstants.SEP_HASHTAG)
-                
.append(msgCodec.getMsgProcType()).append(AttrConstants.SEP_HASHTAG)
-                
.append(DateTimeUtils.ms2yyyyMMddHHmm(longDataTime)).append(AttrConstants.SEP_HASHTAG)
-                
.append(DateTimeUtils.ms2yyyyMMddHHmm(msgCodec.getMsgRcvTime())).toString();
-        strBuff.delete(0, strBuff.length());
         try {
             source.getChannelProcessor().processEvent(event);
-            
source.fileMetricIncSumStats(StatConstants.EVENT_MSG_V0_POST_SUCCESS);
-            source.fileMetricAddSuccCnt(statsKey, msgCodec.getMsgCount(), 1, 
event.getBody().length);
+            source.fileMetricAddSuccStats(strBuff, msgCodec.getGroupId(), 
msgCodec.getStreamId(),
+                    msgCodec.getTopicName(), msgCodec.getStrRemoteIP(), 
msgCodec.getMsgProcType(),
+                    msgCodec.getDataTimeMs(), msgCodec.getMsgPkgTime(), 
msgCodec.getMsgCount(), 1,
+                    event.getBody().length);
             source.addMetric(true, event.getBody().length, event);
             if (msgCodec.isNeedResp() && !msgCodec.isOrderOrProxy()) {
                 msgCodec.setSuccessInfo();
                 responseV0Msg(channel, msgCodec, strBuff);
             }
         } catch (Throwable ex) {
-            
source.fileMetricIncSumStats(StatConstants.EVENT_MSG_V0_POST_FAILURE);
-            source.fileMetricAddFailCnt(statsKey, 1);
+            source.fileMetricAddFailStats(strBuff, msgCodec.getGroupId(), 
msgCodec.getStreamId(),
+                    msgCodec.getTopicName(), msgCodec.getStrRemoteIP(), 
msgCodec.getMsgProcType(),
+                    msgCodec.getDataTimeMs(), msgCodec.getMsgPkgTime(), 1);
             source.addMetric(false, event.getBody().length, event);
             if (msgCodec.isNeedResp() && !msgCodec.isOrderOrProxy()) {
                 
msgCodec.setFailureInfo(DataProxyErrCode.PUT_EVENT_TO_CHANNEL_FAILURE,
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/httpMsg/HttpMessageHandler.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/httpMsg/HttpMessageHandler.java
index 9704e72607..ec759c8c30 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/httpMsg/HttpMessageHandler.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/httpMsg/HttpMessageHandler.java
@@ -23,13 +23,11 @@ import org.apache.inlong.common.monitor.LogCounter;
 import org.apache.inlong.common.msg.AttributeConstants;
 import org.apache.inlong.common.msg.InLongMsg;
 import org.apache.inlong.dataproxy.config.ConfigManager;
-import org.apache.inlong.dataproxy.consts.AttrConstants;
 import org.apache.inlong.dataproxy.consts.ConfigConstants;
 import org.apache.inlong.dataproxy.consts.HttpAttrConst;
 import org.apache.inlong.dataproxy.consts.StatConstants;
 import org.apache.inlong.dataproxy.source.BaseSource;
 import org.apache.inlong.dataproxy.utils.AddressUtils;
-import org.apache.inlong.dataproxy.utils.DateTimeUtils;
 import org.apache.inlong.sdk.commons.protocol.EventConstants;
 
 import io.netty.buffer.ByteBuf;
@@ -343,35 +341,24 @@ public class HttpMessageHandler extends 
SimpleChannelInboundHandler<FullHttpRequ
         eventHeaders.put(ConfigConstants.TOPIC_KEY, topicName);
         eventHeaders.put(AttributeConstants.DATA_TIME, 
String.valueOf(dataTime));
         eventHeaders.put(ConfigConstants.REMOTE_IP_KEY, clientIp);
+        eventHeaders.put(ConfigConstants.DATAPROXY_IP_KEY, 
source.getSrcHost());
         eventHeaders.put(ConfigConstants.MSG_COUNTER_KEY, strMsgCount);
         eventHeaders.put(ConfigConstants.MSG_ENCODE_VER,
                 DataProxyMsgEncType.MSG_ENCODE_TYPE_INLONGMSG.getStrId());
         eventHeaders.put(EventConstants.HEADER_KEY_VERSION,
                 DataProxyMsgEncType.MSG_ENCODE_TYPE_INLONGMSG.getStrId());
         eventHeaders.put(AttributeConstants.RCV_TIME, 
String.valueOf(msgRcvTime));
-        eventHeaders.put(ConfigConstants.PKG_TIME_KEY, 
DateTimeUtils.ms2yyyyMMddHHmm(pkgTime));
+        eventHeaders.put(ConfigConstants.PKG_TIME_KEY, 
String.valueOf(pkgTime));
         Event event = EventBuilder.withBody(inlongMsgData, eventHeaders);
-        // build metric data item
-        dataTime = dataTime / 1000 / 60 / 10;
-        dataTime = dataTime * 1000 * 60 * 10;
-        String statsKey = 
strBuff.append(source.getProtocolName()).append(AttrConstants.SEP_HASHTAG)
-                .append(groupId).append(AttrConstants.SEP_HASHTAG)
-                .append(streamId).append(AttrConstants.SEP_HASHTAG)
-                .append(clientIp).append(AttrConstants.SEP_HASHTAG)
-                .append(source.getSrcHost()).append(AttrConstants.SEP_HASHTAG)
-                .append("b2b").append(AttrConstants.SEP_HASHTAG)
-                
.append(DateTimeUtils.ms2yyyyMMddHHmm(dataTime)).append(AttrConstants.SEP_HASHTAG)
-                .append(DateTimeUtils.ms2yyyyMMddHHmm(msgRcvTime)).toString();
-        strBuff.delete(0, strBuff.length());
         try {
             source.getChannelProcessor().processEvent(event);
-            
source.fileMetricIncSumStats(StatConstants.EVENT_MSG_V0_POST_SUCCESS);
-            source.fileMetricAddSuccCnt(statsKey, intMsgCnt, 1, 
event.getBody().length);
+            source.fileMetricAddSuccStats(strBuff, groupId, streamId, 
topicName, clientIp,
+                    "b2b", dataTime, pkgTime, intMsgCnt, 1, 
event.getBody().length);
             source.addMetric(true, event.getBody().length, event);
             sendSuccessResponse(ctx, isCloseCon, callback);
         } catch (Throwable ex) {
-            
source.fileMetricIncSumStats(StatConstants.EVENT_MSG_V0_POST_FAILURE);
-            source.fileMetricAddFailCnt(statsKey, 1);
+            source.fileMetricAddFailStats(strBuff, groupId, streamId, 
topicName, clientIp,
+                    "b2b", dataTime, pkgTime, 1);
             source.addMetric(false, event.getBody().length, event);
             sendErrorMsg(ctx, DataProxyErrCode.PUT_EVENT_TO_CHANNEL_FAILURE,
                     strBuff.append("Put event to channel failure: 
").append(ex.getMessage()).toString(), callback);
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/AbsV0MsgCodec.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/AbsV0MsgCodec.java
index 11a180a1bd..78a24792be 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/AbsV0MsgCodec.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/AbsV0MsgCodec.java
@@ -23,7 +23,6 @@ import org.apache.inlong.common.msg.AttributeConstants;
 import org.apache.inlong.dataproxy.consts.ConfigConstants;
 import org.apache.inlong.dataproxy.consts.StatConstants;
 import org.apache.inlong.dataproxy.source.BaseSource;
-import org.apache.inlong.dataproxy.utils.DateTimeUtils;
 import org.apache.inlong.sdk.commons.protocol.EventConstants;
 
 import com.google.common.base.Joiner;
@@ -66,6 +65,7 @@ public abstract class AbsV0MsgCodec {
     protected boolean isOrderOrProxy = false;
     protected String msgProcType = "b2b";
     protected boolean needResp = true;
+    protected long msgPkgTime;
 
     public AbsV0MsgCodec(int totalDataLen, int msgTypeValue,
             long msgRcvTime, String strRemoteIP) {
@@ -117,6 +117,10 @@ public abstract class AbsV0MsgCodec {
         return this.dataTimeMs;
     }
 
+    public long getMsgPkgTime() {
+        return msgPkgTime;
+    }
+
     public String getGroupId() {
         return this.groupId;
     }
@@ -219,7 +223,7 @@ public abstract class AbsV0MsgCodec {
         return true;
     }
 
-    protected Map<String, String> buildEventHeaders(long pkgTime) {
+    protected Map<String, String> buildEventHeaders(BaseSource source) {
         // build headers
         Map<String, String> headers = new HashMap<>();
         headers.put(AttributeConstants.GROUP_ID, groupId);
@@ -227,6 +231,7 @@ public abstract class AbsV0MsgCodec {
         headers.put(ConfigConstants.TOPIC_KEY, topicName);
         headers.put(AttributeConstants.DATA_TIME, String.valueOf(dataTimeMs));
         headers.put(ConfigConstants.REMOTE_IP_KEY, strRemoteIP);
+        headers.put(ConfigConstants.DATAPROXY_IP_KEY, source.getSrcHost());
         headers.put(ConfigConstants.MSG_COUNTER_KEY, String.valueOf(msgCount));
         headers.put(ConfigConstants.MSG_ENCODE_VER,
                 DataProxyMsgEncType.MSG_ENCODE_TYPE_INLONGMSG.getStrId());
@@ -234,7 +239,7 @@ public abstract class AbsV0MsgCodec {
                 DataProxyMsgEncType.MSG_ENCODE_TYPE_INLONGMSG.getStrId());
         headers.put(AttributeConstants.RCV_TIME, String.valueOf(msgRcvTime));
         headers.put(AttributeConstants.UNIQ_ID, String.valueOf(uniq));
-        headers.put(ConfigConstants.PKG_TIME_KEY, 
DateTimeUtils.ms2yyyyMMddHHmm(pkgTime));
+        headers.put(ConfigConstants.PKG_TIME_KEY, String.valueOf(msgPkgTime));
         // add extra key-value information
         if (!needResp) {
             headers.put(AttributeConstants.MESSAGE_IS_ACK, "false");
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/CodecBinMsg.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/CodecBinMsg.java
index 71b80eb17e..a0e4adc6b1 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/CodecBinMsg.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/CodecBinMsg.java
@@ -236,7 +236,8 @@ public class CodecBinMsg extends AbsV0MsgCodec {
         InLongMsg inLongMsg = InLongMsg.newInLongMsg(source.isCompressed(), 4);
         inLongMsg.addMsg(dataBuf.array());
         byte[] inlongMsgData = inLongMsg.buildArray();
-        Event event = EventBuilder.withBody(inlongMsgData, 
buildEventHeaders(inLongMsg.getCreatetime()));
+        msgPkgTime = inLongMsg.getCreatetime();
+        Event event = EventBuilder.withBody(inlongMsgData, 
buildEventHeaders(source));
         if (isOrderOrProxy) {
             event = new SinkRspEvent(event, MsgType.MSG_BIN_MULTI_BODY, 
channel);
         }
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/CodecTextMsg.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/CodecTextMsg.java
index 5abadf20c9..c3bf9281ac 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/CodecTextMsg.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/v0msg/CodecTextMsg.java
@@ -267,7 +267,8 @@ public class CodecTextMsg extends AbsV0MsgCodec {
             inLongMsg.addMsg(mapJoiner.join(attrMap), bodyData);
         }
         byte[] inlongMsgData = inLongMsg.buildArray();
-        Event event = EventBuilder.withBody(inlongMsgData, 
buildEventHeaders(inLongMsg.getCreatetime()));
+        msgPkgTime = inLongMsg.getCreatetime();
+        Event event = EventBuilder.withBody(inlongMsgData, 
buildEventHeaders(source));
         inLongMsg.reset();
         return event;
     }
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/utils/DateTimeUtils.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/utils/DateTimeUtils.java
index 1a94b2f64c..1eea8d2049 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/utils/DateTimeUtils.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/utils/DateTimeUtils.java
@@ -38,4 +38,18 @@ public class DateTimeUtils {
                 LocalDateTime.ofInstant(Instant.ofEpochMilli(timestamp), 
defZoneId);
         return DATE_FORMATTER.format(localDateTime);
     }
+
+    /**
+     * convert ms value to ten minute level 'yyyyMMddHHmm' string
+     *
+     * @param timestamp The millisecond value of the specified time
+     * @return the time string in ten-minute format yyyyMMddHHmmss
+     */
+    public static String ms2yyyyMMddHHmmTenMins(long timestamp) {
+        long longDataTime = timestamp / 1000 / 60 / 10;
+        LocalDateTime localDateTime = LocalDateTime.ofInstant(
+                Instant.ofEpochMilli(longDataTime * 1000 * 60 * 10), 
defZoneId);
+        return DATE_FORMATTER.format(localDateTime);
+    }
+
 }


Reply via email to