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 3b20180e8f [INLONG-8368][DataProxy] Sink does not have audit data 
(#8375)
3b20180e8f is described below

commit 3b20180e8f611dece88a2f15e4140a89f745f892
Author: Goson Zhang <[email protected]>
AuthorDate: Fri Jun 30 10:07:14 2023 +0800

    [INLONG-8368][DataProxy] Sink does not have audit data (#8375)
---
 .../java/org/apache/inlong/dataproxy/sink/mq/SimplePackProfile.java   | 1 -
 .../java/org/apache/inlong/dataproxy/source/ServerMessageHandler.java | 2 +-
 .../java/org/apache/inlong/dataproxy/source/v0msg/CodecBinMsg.java    | 4 ++--
 .../java/org/apache/inlong/dataproxy/source/v0msg/CodecTextMsg.java   | 4 ++--
 4 files changed, 5 insertions(+), 6 deletions(-)

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 c131a4172a..242042512c 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
@@ -156,7 +156,6 @@ public class SimplePackProfile extends PackProfile {
         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));
-        result.put(ConfigConstants.PKG_TIME_KEY, 
event.getHeaders().get(ConfigConstants.PKG_TIME_KEY));
         result.put(ConfigConstants.DATAPROXY_IP_KEY, 
NetworkUtils.getLocalIp());
         return result;
     }
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 929407a778..41690d4f69 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
@@ -287,7 +287,7 @@ public class ServerMessageHandler extends 
ChannelInboundHandlerAdapter {
         try {
             source.getChannelProcessor().processEvent(event);
             
source.fileMetricIncSumStats(StatConstants.EVENT_MSG_V0_POST_SUCCESS);
-            source.fileMetricAddSuccCnt(statsKey, msgCodec.getMsgCount(), 1, 
msgCodec.getBodyLength());
+            source.fileMetricAddSuccCnt(statsKey, msgCodec.getMsgCount(), 1, 
event.getBody().length);
             source.addMetric(true, event.getBody().length, event);
             if (msgCodec.isNeedResp() && !msgCodec.isOrderOrProxy()) {
                 msgCodec.setSuccessInfo();
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 260f38bb8f..71b80eb17e 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
@@ -235,8 +235,8 @@ public class CodecBinMsg extends AbsV0MsgCodec {
         // build InLong message
         InLongMsg inLongMsg = InLongMsg.newInLongMsg(source.isCompressed(), 4);
         inLongMsg.addMsg(dataBuf.array());
-        long pkgTime = inLongMsg.getCreatetime();
-        Event event = EventBuilder.withBody(inLongMsg.buildArray(), 
buildEventHeaders(pkgTime));
+        byte[] inlongMsgData = inLongMsg.buildArray();
+        Event event = EventBuilder.withBody(inlongMsgData, 
buildEventHeaders(inLongMsg.getCreatetime()));
         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 772529f38c..5abadf20c9 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
@@ -266,8 +266,8 @@ public class CodecTextMsg extends AbsV0MsgCodec {
             }
             inLongMsg.addMsg(mapJoiner.join(attrMap), bodyData);
         }
-        long pkgTime = inLongMsg.getCreatetime();
-        Event event = EventBuilder.withBody(inLongMsg.buildArray(), 
buildEventHeaders(pkgTime));
+        byte[] inlongMsgData = inLongMsg.buildArray();
+        Event event = EventBuilder.withBody(inlongMsgData, 
buildEventHeaders(inLongMsg.getCreatetime()));
         inLongMsg.reset();
         return event;
     }

Reply via email to