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