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 017e4006d7 [INLONG-8318][DataProxy] Change notification 
synchronization through condition variables and locks (#8320)
017e4006d7 is described below

commit 017e4006d72ab26372826078e8f516c8758228f8
Author: Goson Zhang <[email protected]>
AuthorDate: Mon Jun 26 15:49:50 2023 +0800

    [INLONG-8318][DataProxy] Change notification synchronization through 
condition variables and locks (#8320)
---
 .../inlong/dataproxy/consts/StatConstants.java     |  3 ++
 .../sink/mq/MessageQueueZoneProducer.java          | 62 ++++++++++++----------
 .../dataproxy/sink/mq/MessageQueueZoneSink.java    | 27 +++++++---
 .../dataproxy/source2/v0msg/CodecBinMsg.java       |  2 +-
 4 files changed, 57 insertions(+), 37 deletions(-)

diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/consts/StatConstants.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/consts/StatConstants.java
index 9ffe88ffff..f7cc229fc9 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/consts/StatConstants.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/consts/StatConstants.java
@@ -85,6 +85,9 @@ public class StatConstants {
     public static final java.lang.String EVENT_SINK_DEFAULT_TOPIC_MISSING = 
"default.topic.empty";
     public static final java.lang.String EVENT_SINK_DEFAULT_TOPIC_USED = 
"default.topic.used";
     public static final java.lang.String EVENT_SINK_PRODUCER_NULL = 
"sink.producer.null";
+    public static final java.lang.String EVENT_SINK_CLUSTER_EMPTY = 
"sink.cluster.empty";
+    public static final java.lang.String EVENT_SINK_CLUSTER_UNMATCHED = 
"sink.cluster.unmatched";
+    public static final java.lang.String EVENT_SINK_CPRODUCER_NULL = 
"sink.cluster.producer.null";
     public static final java.lang.String EVENT_SINK_SEND_EXCEPTION = 
"sink.send.exception";
 
     public static final java.lang.String EVENT_SINK_FAILRETRY = "sink.retry";
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneProducer.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneProducer.java
index be20c21ccf..358165a3af 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneProducer.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneProducer.java
@@ -19,6 +19,7 @@ package org.apache.inlong.dataproxy.sink.mq;
 
 import org.apache.inlong.dataproxy.config.ConfigManager;
 import org.apache.inlong.dataproxy.config.pojo.CacheClusterConfig;
+import org.apache.inlong.dataproxy.consts.StatConstants;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -30,8 +31,6 @@ import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.atomic.AtomicInteger;
-import java.util.concurrent.locks.ReadWriteLock;
-import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 /**
  * 
@@ -46,7 +45,6 @@ public class MessageQueueZoneProducer {
     private final CacheClusterSelector cacheClusterSelector;
 
     private final AtomicInteger clusterIndex = new AtomicInteger(0);
-    private final ReadWriteLock readWriteLock = new ReentrantReadWriteLock();
     private List<String> currentClusterNames = new ArrayList<>();
     private final ConcurrentHashMap<String, Long> usingTimeMap = new 
ConcurrentHashMap<>();
     private final ConcurrentHashMap<String, MessageQueueClusterProducer> 
usingClusterMap = new ConcurrentHashMap<>();
@@ -150,24 +148,29 @@ public class MessageQueueZoneProducer {
      */
     public boolean send(PackProfile profile) {
         String clusterName;
+        List<String> tmpClusters;
         MessageQueueClusterProducer clusterProducer;
-        readWriteLock.readLock().lock();
-        try {
-            do {
-                clusterName = currentClusterNames.get(
-                        Math.abs(clusterIndex.getAndIncrement()) % 
currentClusterNames.size());
-                if (clusterName == null) {
-                    continue;
-                }
-                clusterProducer = usingClusterMap.get(clusterName);
-                if (clusterProducer == null) {
-                    continue;
-                }
-                return clusterProducer.send(profile);
-            } while (true);
-        } finally {
-            readWriteLock.readLock().unlock();
-        }
+        do {
+            tmpClusters = currentClusterNames;
+            if (tmpClusters == null || tmpClusters.isEmpty()) {
+                
context.fileMetricIncSumStats(StatConstants.EVENT_SINK_CLUSTER_EMPTY);
+                sleepSomeTime(100);
+                continue;
+            }
+            clusterName = 
tmpClusters.get(Math.abs(clusterIndex.getAndIncrement()) % tmpClusters.size());
+            if (clusterName == null) {
+                
context.fileMetricIncSumStats(StatConstants.EVENT_SINK_CLUSTER_UNMATCHED);
+                sleepSomeTime(100);
+                continue;
+            }
+            clusterProducer = usingClusterMap.get(clusterName);
+            if (clusterProducer == null) {
+                
context.fileMetricIncWithDetailStats(StatConstants.EVENT_SINK_CPRODUCER_NULL, 
clusterName);
+                sleepSomeTime(100);
+                continue;
+            }
+            return clusterProducer.send(profile);
+        } while (true);
     }
 
     private void checkAndReloadClusterInfo() {
@@ -225,14 +228,9 @@ public class MessageQueueZoneProducer {
                 }
             }
             // replace cluster names
-            readWriteLock.writeLock().lock();
-            try {
-                if (!lastClusterNames.equals(currentClusterNames)) {
-                    changed = true;
-                    currentClusterNames = lastClusterNames;
-                }
-            } finally {
-                readWriteLock.writeLock().unlock();
+            if (!lastClusterNames.equals(currentClusterNames)) {
+                currentClusterNames = lastClusterNames;
+                changed = true;
             }
             // filter removed records
             Set<String> needRmvs = new HashSet<>();
@@ -293,4 +291,12 @@ public class MessageQueueZoneProducer {
             clusterProducer.publishTopic(curTopicSet);
         }
     }
+
+    private void sleepSomeTime(long millis) {
+        try {
+            Thread.sleep(millis);
+        } catch (Throwable e) {
+            //
+        }
+    }
 }
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSink.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSink.java
index 6edb9f5bc0..76e50d697e 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSink.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/sink/mq/MessageQueueZoneSink.java
@@ -46,6 +46,8 @@ import java.util.concurrent.Executors;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.locks.Condition;
+import java.util.concurrent.locks.ReentrantLock;
 
 /**
  * MessageQueueZoneSink
@@ -72,7 +74,8 @@ public class MessageQueueZoneSink extends AbstractSink 
implements Configurable,
 
     private MessageQueueZoneProducer zoneProducer;
     // configure change notify
-    private final Object syncLock = new Object();
+    private final ReentrantLock reentrantLock = new ReentrantLock();
+    private final Condition condition = reentrantLock.newCondition();
     private final AtomicLong lastNotifyTime = new AtomicLong(0);
     // changeListerThread
     private Thread configListener;
@@ -295,8 +298,13 @@ public class MessageQueueZoneSink extends AbstractSink 
implements Configurable,
         if (zoneProducer == null) {
             return;
         }
-        lastNotifyTime.set(System.currentTimeMillis());
-        syncLock.notifyAll();
+        reentrantLock.lock();
+        try {
+            lastNotifyTime.set(System.currentTimeMillis());
+            condition.signal();
+        } finally {
+            reentrantLock.unlock();
+        }
     }
 
     /**
@@ -311,14 +319,16 @@ public class MessageQueueZoneSink extends AbstractSink 
implements Configurable,
         @Override
         public void run() {
             long lastCheckTime;
+            logger.info("{} config-change processor start!", getName());
             while (!isShutdown) {
+                reentrantLock.lock();
                 try {
-                    syncLock.wait();
-                } catch (InterruptedException e) {
-                    logger.error("{} config-change processor meet interrupt, 
exit!", getName());
+                    condition.await();
+                } catch (InterruptedException e1) {
+                    logger.info("{} config-change processor meet interrupt, 
break!", getName());
                     break;
-                } catch (Throwable e2) {
-                    //
+                } finally {
+                    reentrantLock.unlock();
                 }
                 if (zoneProducer == null) {
                     continue;
@@ -328,6 +338,7 @@ public class MessageQueueZoneSink extends AbstractSink 
implements Configurable,
                     zoneProducer.reloadMetaConfig();
                 } while (lastCheckTime != lastNotifyTime.get());
             }
+            logger.info("{} config-change processor exit!", getName());
         }
     }
 }
diff --git 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecBinMsg.java
 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecBinMsg.java
index 9f31e7f97e..cd2353f01e 100644
--- 
a/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecBinMsg.java
+++ 
b/inlong-dataproxy/dataproxy-source/src/main/java/org/apache/inlong/dataproxy/source2/v0msg/CodecBinMsg.java
@@ -229,7 +229,7 @@ public class CodecBinMsg extends AbsV0MsgCodec {
                 - BIN_MSG_ATTRLEN_SIZE - BIN_MSG_MAGIC_SIZE - 
origAttr.length(), (short) origAttr.length());
         if (origAttr.length() > 0) {
             System.arraycopy(origAttr.getBytes(StandardCharsets.UTF_8), 0, 
dataBuf.array(),
-                    totalPkgLength - BIN_MSG_MAGIC_SIZE - origAttr.length(), 
bodyData.length);
+                    totalPkgLength - BIN_MSG_MAGIC_SIZE - origAttr.length(), 
origAttr.length());
         }
         dataBuf.putShort(totalPkgLength - BIN_MSG_MAGIC_SIZE, (short) 
BIN_MSG_MAGIC);
         // build InLong message

Reply via email to