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