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