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 83ded2149d [INLONG-10313][DataProxy] Replace audit ID macro with audit 
API (#10315)
83ded2149d is described below

commit 83ded2149d032c6e4cd2085e9785bc52bc49faea
Author: Goson Zhang <[email protected]>
AuthorDate: Thu May 30 15:11:20 2024 +0800

    [INLONG-10313][DataProxy] Replace audit ID macro with audit API (#10315)
---
 .../inlong/dataproxy/metrics/audit/AuditUtils.java | 33 ++++++++++++++++++----
 .../sink/mq/MessageQueueZoneSinkContext.java       |  5 ++--
 .../apache/inlong/dataproxy/source/BaseSource.java |  2 +-
 3 files changed, 31 insertions(+), 9 deletions(-)

diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/metrics/audit/AuditUtils.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/metrics/audit/AuditUtils.java
index 899aca38d4..ec81ff36ff 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/metrics/audit/AuditUtils.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/metrics/audit/AuditUtils.java
@@ -17,6 +17,7 @@
 
 package org.apache.inlong.dataproxy.metrics.audit;
 
+import org.apache.inlong.audit.AuditIdEnum;
 import org.apache.inlong.audit.AuditOperator;
 import org.apache.inlong.audit.util.AuditConfig;
 import org.apache.inlong.common.enums.MessageWrapType;
@@ -39,8 +40,8 @@ import static 
org.apache.inlong.audit.consts.ConfigConstants.DEFAULT_AUDIT_TAG;
  */
 public class AuditUtils {
 
-    public static final int AUDIT_ID_DATAPROXY_READ_SUCCESS = 5;
-    public static final int AUDIT_ID_DATAPROXY_SEND_SUCCESS = 6;
+    private static int auditIdReadSuccess = 5;
+    private static int auditIdSendSuccess = 6;
 
     /**
      * Init audit
@@ -55,16 +56,38 @@ public class AuditUtils {
                     CommonConfigHolder.getInstance().getAuditFilePath(),
                     CommonConfigHolder.getInstance().getAuditMaxCacheRows());
             AuditOperator.getInstance().setAuditConfig(auditConfig);
+            auditIdReadSuccess =
+                    
AuditOperator.getInstance().buildSuccessfulAuditId(AuditIdEnum.DATA_PROXY_INPUT);
+            auditIdSendSuccess =
+                    
AuditOperator.getInstance().buildSuccessfulAuditId(AuditIdEnum.DATA_PROXY_OUTPUT);
         }
     }
 
     /**
-     * Add audit data
+     * Add input audit data
+     *
+     * @param event    event to be counted
+     */
+    public static void addInputSuccess(Event event) {
+        if (event == null || 
!CommonConfigHolder.getInstance().isEnableAudit()) {
+            return;
+        }
+        addAuditData(event, auditIdReadSuccess);
+    }
+
+    /**
+     * Add output audit data
+     *
+     * @param event    event to be counted
      */
-    public static void add(int auditID, Event event) {
-        if (!CommonConfigHolder.getInstance().isEnableAudit() || event == 
null) {
+    public static void addOutputSuccess(Event event) {
+        if (event == null || 
!CommonConfigHolder.getInstance().isEnableAudit()) {
             return;
         }
+        addAuditData(event, auditIdSendSuccess);
+    }
+
+    private static void addAuditData(Event event, int auditID) {
         Map<String, String> headers = event.getHeaders();
         String pkgVersion = headers.get(ConfigConstants.MSG_ENCODE_VER);
         if 
(MessageWrapType.INLONG_MSG_V1.getStrId().equalsIgnoreCase(pkgVersion)) {
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 9bae668283..3681148456 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
@@ -134,8 +134,7 @@ public class MessageQueueZoneSinkContext extends 
SinkContext {
             long sendTime) {
         if (currentRecord instanceof SimplePackProfile) {
             if (result) {
-                AuditUtils.add(AuditUtils.AUDIT_ID_DATAPROXY_SEND_SUCCESS,
-                        ((SimplePackProfile) currentRecord).getEvent());
+                AuditUtils.addOutputSuccess(((SimplePackProfile) 
currentRecord).getEvent());
             }
             return;
         }
@@ -166,7 +165,7 @@ public class MessageQueueZoneSinkContext extends 
SinkContext {
                     metricItem.nodeDuration.addAndGet(nodeDuration);
                     metricItem.wholeDuration.addAndGet(wholeDuration);
                 }
-                AuditUtils.add(AuditUtils.AUDIT_ID_DATAPROXY_SEND_SUCCESS, 
event);
+                AuditUtils.addOutputSuccess(event);
             } else {
                 metricItem.sendFailCount.addAndGet(1);
                 metricItem.sendFailSize.addAndGet(event.getBody().length);
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/BaseSource.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/BaseSource.java
index 71f5c538d4..c3061fa7b4 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/BaseSource.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source/BaseSource.java
@@ -485,7 +485,7 @@ public abstract class BaseSource
         if (result) {
             metricItem.readSuccessCount.incrementAndGet();
             metricItem.readSuccessSize.addAndGet(size);
-            AuditUtils.add(AuditUtils.AUDIT_ID_DATAPROXY_READ_SUCCESS, event);
+            AuditUtils.addInputSuccess(event);
         } else {
             metricItem.readFailCount.incrementAndGet();
             metricItem.readFailSize.addAndGet(size);

Reply via email to