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