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

Reply via email to