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 2b64674d6a [INLONG-8332][DataProxy] Return original content for
MSG_ORIGINAL_RETURN type messages (#8333)
2b64674d6a is described below
commit 2b64674d6a61e4438ed864911380994324c2d3dc
Author: Goson Zhang <[email protected]>
AuthorDate: Tue Jun 27 17:00:54 2023 +0800
[INLONG-8332][DataProxy] Return original content for MSG_ORIGINAL_RETURN
type messages (#8333)
---
.../sink/mq/MessageQueueZoneSinkContext.java | 12 +++----
.../dataproxy/source2/InLongMessageHandler.java | 39 +++++++++++++++++++++-
.../dataproxy/source2/v0msg/AbsV0MsgCodec.java | 5 +++
.../dataproxy/source2/v0msg/CodecTextMsg.java | 4 +++
4 files changed, 53 insertions(+), 7 deletions(-)
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 becfade7bd..6989c94ea0 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
@@ -302,12 +302,12 @@ public class MessageQueueZoneSinkContext extends
SinkContext {
if (packProfile instanceof SimplePackProfile) {
SimplePackProfile simpleProfile = (SimplePackProfile) packProfile;
StringBuilder statsKey = new StringBuilder(512)
- .append(sinkName).append(AttrConstants.SEPARATOR)
-
.append(simpleProfile.getInlongGroupId()).append(AttrConstants.SEPARATOR)
-
.append(simpleProfile.getInlongStreamId()).append(AttrConstants.SEPARATOR)
- .append(topic).append(AttrConstants.SEPARATOR)
-
.append(NetworkUtils.getLocalIp()).append(AttrConstants.SEPARATOR)
- .append(remoteId).append(AttrConstants.SEPARATOR)
+ .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/source2/InLongMessageHandler.java
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/InLongMessageHandler.java
index ad34d1e3b9..0a2d27c1c8 100644
---
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/InLongMessageHandler.java
+++
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/InLongMessageHandler.java
@@ -479,7 +479,7 @@ public class InLongMessageHandler extends
ChannelInboundHandlerAdapter {
if (MsgType.MSG_BIN_MULTI_BODY.equals(msgType)) {
retData = buildBinMsgRspPackage(strBuff.toString(),
msgObj.getUniq());
} else {
- retData = buildTxtMsgRspPackage(msgType, strBuff.toString());
+ retData = buildTxtMsgRspPackage(msgType, strBuff.toString(),
msgObj);
}
strBuff.delete(0, strBuff.length());
flushV0MsgPackage(source, channel, retData, msgObj.getAttr());
@@ -598,6 +598,43 @@ public class InLongMessageHandler extends
ChannelInboundHandlerAdapter {
return buffer;
}
+ /**
+ * Build default-msg response message ByteBuf
+ *
+ * @param msgType the message type
+ * @param attrs the return attribute
+ * @param msgObj the request message object
+ * @return ByteBuf
+ */
+ private ByteBuf buildTxtMsgRspPackage(MsgType msgType, String attrs,
AbsV0MsgCodec msgObj) {
+ int attrsLen = 0;
+ int bodyLen = 0;
+ byte[] backBody = null;
+ if (attrs != null) {
+ attrsLen = attrs.length();
+ }
+ if (MsgType.MSG_ORIGINAL_RETURN.equals(msgType)) {
+ backBody = msgObj.getOrigBody();
+ if (backBody != null) {
+ bodyLen = backBody.length;
+ }
+ }
+ // backTotalLen = mstType + bodyLen + body + attrsLen + attrs
+ int backTotalLen = 1 + 4 + bodyLen + 4 + attrsLen;
+ ByteBuf buffer = ByteBufAllocator.DEFAULT.buffer(4 + backTotalLen);
+ buffer.writeInt(backTotalLen);
+ buffer.writeByte(msgType.getValue());
+ buffer.writeInt(bodyLen);
+ if (bodyLen > 0) {
+ buffer.writeBytes(backBody);
+ }
+ buffer.writeInt(attrsLen);
+ if (attrsLen > 0) {
+ buffer.writeBytes(attrs.getBytes(StandardCharsets.UTF_8));
+ }
+ return buffer;
+ }
+
/**
* Build heartbeat response message ByteBuf
*
diff --git
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java
index 0764f7e599..fd367a3522 100644
---
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java
+++
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java
@@ -56,6 +56,7 @@ public abstract class AbsV0MsgCodec {
protected int msgCount;
protected String origAttr = "";
protected byte[] bodyData;
+ protected byte[] origBody = null;
protected long dataTimeMs;
protected String groupId;
protected String streamId = "";
@@ -144,6 +145,10 @@ public abstract class AbsV0MsgCodec {
return strRemoteIP;
}
+ public byte[] getOrigBody() {
+ return origBody;
+ }
+
public long getMsgRcvTime() {
return msgRcvTime;
}
diff --git
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecTextMsg.java
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecTextMsg.java
index 07a8f1aac5..c7189dbe80 100644
---
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecTextMsg.java
+++
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecTextMsg.java
@@ -71,6 +71,10 @@ public class CodecTextMsg extends AbsV0MsgCodec {
// extract body bytes
this.bodyData = new byte[bodyLen];
cb.getBytes(msgHeadPos + TXT_MSG_BODY_OFFSET, this.bodyData, 0,
bodyLen);
+ if (MsgType.MSG_ORIGINAL_RETURN.equals(MsgType.valueOf(msgType))) {
+ this.origBody = new byte[bodyLen];
+ System.arraycopy(this.bodyData, 0, this.origBody, 0, bodyLen);
+ }
// get attribute length
int attrLen = cb.getInt(msgHeadPos + TXT_MSG_BODY_OFFSET + bodyLen);
if (attrLen < 0) {