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 721486143f [INLONG-8323][DataProxy] Add Topic detailed information
output when Producer is null (#8324)
721486143f is described below
commit 721486143f98aa7c613bee73eca32c89c240427f
Author: Goson Zhang <[email protected]>
AuthorDate: Mon Jun 26 19:59:08 2023 +0800
[INLONG-8323][DataProxy] Add Topic detailed information output when
Producer is null (#8324)
---
.../org/apache/inlong/dataproxy/sink/mq/SimplePackProfile.java | 9 +++++++++
.../org/apache/inlong/dataproxy/sink/mq/kafka/KafkaHandler.java | 2 +-
.../apache/inlong/dataproxy/sink/mq/pulsar/PulsarHandler.java | 2 +-
.../org/apache/inlong/dataproxy/sink/mq/tube/TubeHandler.java | 2 +-
.../apache/inlong/dataproxy/source2/InLongMessageHandler.java | 6 +++---
.../org/apache/inlong/dataproxy/source2/v0msg/AbsV0MsgCodec.java | 3 +++
6 files changed, 18 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 5949d38e4e..fa292893e0 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,6 +156,7 @@ 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;
}
@@ -191,6 +192,14 @@ public class SimplePackProfile extends PackProfile {
.append(AttributeConstants.KEY_VALUE_SEPARATOR).append(errMsg);
}
}
+ String origAttr =
event.getHeaders().getOrDefault(ConfigConstants.DECODER_ATTRS, "");
+ if (StringUtils.isNotEmpty(origAttr)) {
+ if (strBuff.length() > 0) {
+
strBuff.append(AttributeConstants.SEPARATOR).append(origAttr);
+ } else {
+ strBuff.append(origAttr);
+ }
+ }
// build and send response message
ByteBuf retData;
if (MsgType.MSG_BIN_MULTI_BODY.equals(msgType)) {
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 126ca385dc..3d7303473e 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
@@ -146,7 +146,7 @@ public class KafkaHandler implements MessageQueueHandler {
}
// create producer failed
if (producer == null) {
-
sinkContext.fileMetricIncSumStats(StatConstants.EVENT_SINK_PRODUCER_NULL);
+
sinkContext.fileMetricIncWithDetailStats(StatConstants.EVENT_SINK_PRODUCER_NULL,
topic);
sinkContext.processSendFail(profile, clusterName, topic, 0,
DataProxyErrCode.PRODUCER_IS_NULL, "");
return false;
}
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 66f8067014..149f3e5e1b 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
@@ -253,7 +253,7 @@ public class PulsarHandler implements MessageQueueHandler {
}
// create producer failed
if (producer == null) {
-
sinkContext.fileMetricIncSumStats(StatConstants.EVENT_SINK_PRODUCER_NULL);
+
sinkContext.fileMetricIncWithDetailStats(StatConstants.EVENT_SINK_PRODUCER_NULL,
producerTopic);
sinkContext.processSendFail(profile, clusterName,
producerTopic, 0,
DataProxyErrCode.PRODUCER_IS_NULL, "");
return false;
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 dcd28d13c3..c5c9dcbcfb 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
@@ -211,7 +211,7 @@ public class TubeHandler implements MessageQueueHandler {
}
// create producer failed
if (producer == null) {
-
sinkContext.fileMetricIncSumStats(StatConstants.EVENT_SINK_PRODUCER_NULL);
+
sinkContext.fileMetricIncWithDetailStats(StatConstants.EVENT_SINK_PRODUCER_NULL,
topic);
sinkContext.processSendFail(profile, clusterName, topic, 0,
DataProxyErrCode.PRODUCER_IS_NULL, "");
return false;
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 1ff018786b..ad34d1e3b9 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
@@ -469,9 +469,9 @@ public class InLongMessageHandler extends
ChannelInboundHandlerAdapter {
strBuff.append(AttributeConstants.SEPARATOR).append(AttributeConstants.MESSAGE_PROCESS_ERRMSG)
.append(AttributeConstants.KEY_VALUE_SEPARATOR).append(msgObj.getErrMsg());
}
- if (StringUtils.isNotEmpty(msgObj.getAttr())) {
-
strBuff.append(AttributeConstants.SEPARATOR).append(msgObj.getAttr());
- }
+ }
+ if (StringUtils.isNotEmpty(msgObj.getAttr())) {
+
strBuff.append(AttributeConstants.SEPARATOR).append(msgObj.getAttr());
}
// build and send response message
ByteBuf retData;
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 83d82da94e..0764f7e599 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
@@ -242,6 +242,9 @@ public abstract class AbsV0MsgCodec {
if (StringUtils.isNotEmpty(proxySend)) {
headers.put(AttributeConstants.MESSAGE_PROXY_SEND, proxySend);
}
+ if (isOrderOrProxy) {
+ headers.put(ConfigConstants.DECODER_ATTRS, this.origAttr);
+ }
String partitionKey =
attrMap.get(AttributeConstants.MESSAGE_PARTITION_KEY);
if (StringUtils.isNotEmpty(partitionKey)) {
headers.put(AttributeConstants.MESSAGE_PARTITION_KEY,
partitionKey);