This is an automated email from the ASF dual-hosted git repository.
albumenj pushed a commit to branch 3.2
in repository https://gitbox.apache.org/repos/asf/dubbo.git
The following commit(s) were added to refs/heads/3.2 by this push:
new 7d5fa50783 Fix metrics unable to retrieve lazy inited executor status
(#14348)
7d5fa50783 is described below
commit 7d5fa50783c5321447368230ca75138425697459
Author: Albumen Kevin <[email protected]>
AuthorDate: Mon Jun 24 09:26:45 2024 +0800
Fix metrics unable to retrieve lazy inited executor status (#14348)
* Fix metrics unable to retrieve lazy inited executor status
* Enhance
---
.../org/apache/dubbo/common/store/DataStore.java | 2 ++
...DataStore.java => DataStoreUpdateListener.java} | 20 ++------------
.../common/store/support/SimpleDataStore.java | 30 +++++++++++++++++++++
.../java/org/apache/dubbo/config/Constants.java | 4 +++
.../common/store/support/SimpleDataStoreTest.java | 30 +++++++++++++++++++++
.../collector/sample/ThreadPoolMetricsSampler.java | 30 ++++++++++++++++-----
.../sample/ThreadPoolMetricsSamplerTest.java | 31 ++++++++++++++++++++++
7 files changed, 123 insertions(+), 24 deletions(-)
diff --git
a/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStore.java
b/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStore.java
index 4b8fc31ab6..d0c20b9e8b 100644
--- a/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStore.java
+++ b/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStore.java
@@ -34,4 +34,6 @@ public interface DataStore {
void put(String componentName, String key, Object value);
void remove(String componentName, String key);
+
+ default void addListener(DataStoreUpdateListener dataStoreUpdateListener)
{}
}
diff --git
a/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStore.java
b/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStoreUpdateListener.java
similarity index 62%
copy from
dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStore.java
copy to
dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStoreUpdateListener.java
index 4b8fc31ab6..de994191f7 100644
--- a/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStore.java
+++
b/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStoreUpdateListener.java
@@ -16,22 +16,6 @@
*/
package org.apache.dubbo.common.store;
-import org.apache.dubbo.common.extension.ExtensionScope;
-import org.apache.dubbo.common.extension.SPI;
-
-import java.util.Map;
-
-@SPI(value = "simple", scope = ExtensionScope.APPLICATION)
-public interface DataStore {
-
- /**
- * return a snapshot value of componentName
- */
- Map<String, Object> get(String componentName);
-
- Object get(String componentName, String key);
-
- void put(String componentName, String key, Object value);
-
- void remove(String componentName, String key);
+public interface DataStoreUpdateListener {
+ void onUpdate(String componentName, String key, Object value);
}
diff --git
a/dubbo-common/src/main/java/org/apache/dubbo/common/store/support/SimpleDataStore.java
b/dubbo-common/src/main/java/org/apache/dubbo/common/store/support/SimpleDataStore.java
index cabb9e903d..6bda3cfada 100644
---
a/dubbo-common/src/main/java/org/apache/dubbo/common/store/support/SimpleDataStore.java
+++
b/dubbo-common/src/main/java/org/apache/dubbo/common/store/support/SimpleDataStore.java
@@ -16,8 +16,13 @@
*/
package org.apache.dubbo.common.store.support;
+import org.apache.dubbo.common.constants.LoggerCodeConstants;
+import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
+import org.apache.dubbo.common.logger.LoggerFactory;
import org.apache.dubbo.common.store.DataStore;
+import org.apache.dubbo.common.store.DataStoreUpdateListener;
import org.apache.dubbo.common.utils.ConcurrentHashMapUtils;
+import org.apache.dubbo.common.utils.ConcurrentHashSet;
import java.util.HashMap;
import java.util.Map;
@@ -25,9 +30,11 @@ import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
public class SimpleDataStore implements DataStore {
+ private static final ErrorTypeAwareLogger logger =
LoggerFactory.getErrorTypeAwareLogger(SimpleDataStore.class);
// <component name or id, <data-name, data-value>>
private final ConcurrentMap<String, ConcurrentMap<String, Object>> data =
new ConcurrentHashMap<>();
+ private final ConcurrentHashSet<DataStoreUpdateListener> listeners = new
ConcurrentHashSet<>();
@Override
public Map<String, Object> get(String componentName) {
@@ -52,6 +59,7 @@ public class SimpleDataStore implements DataStore {
Map<String, Object> componentData =
ConcurrentHashMapUtils.computeIfAbsent(data, componentName, k
-> new ConcurrentHashMap<>());
componentData.put(key, value);
+ notifyListeners(componentName, key, value);
}
@Override
@@ -60,5 +68,27 @@ public class SimpleDataStore implements DataStore {
return;
}
data.get(componentName).remove(key);
+ notifyListeners(componentName, key, null);
+ }
+
+ @Override
+ public void addListener(DataStoreUpdateListener dataStoreUpdateListener) {
+ listeners.add(dataStoreUpdateListener);
+ }
+
+ private void notifyListeners(String componentName, String key, Object
value) {
+ for (DataStoreUpdateListener listener : listeners) {
+ try {
+ listener.onUpdate(componentName, key, value);
+ } catch (Throwable t) {
+ logger.warn(
+ LoggerCodeConstants.INTERNAL_ERROR,
+ "",
+ "",
+ "Failed to notify data store update listener. " +
"ComponentName: " + componentName + " Key: "
+ + key,
+ t);
+ }
+ }
}
}
diff --git a/dubbo-common/src/main/java/org/apache/dubbo/config/Constants.java
b/dubbo-common/src/main/java/org/apache/dubbo/config/Constants.java
index da6ab9f5d7..a281cb2d7b 100644
--- a/dubbo-common/src/main/java/org/apache/dubbo/config/Constants.java
+++ b/dubbo-common/src/main/java/org/apache/dubbo/config/Constants.java
@@ -149,7 +149,11 @@ public interface Constants {
String SERVER_THREAD_POOL_NAME = "DubboServerHandler";
+ String SERVER_THREAD_POOL_PREFIX = SERVER_THREAD_POOL_NAME + "-";
+
String CLIENT_THREAD_POOL_NAME = "DubboClientHandler";
+ String CLIENT_THREAD_POOL_PREFIX = CLIENT_THREAD_POOL_NAME + "-";
+
String REST_PROTOCOL = "rest";
}
diff --git
a/dubbo-common/src/test/java/org/apache/dubbo/common/store/support/SimpleDataStoreTest.java
b/dubbo-common/src/test/java/org/apache/dubbo/common/store/support/SimpleDataStoreTest.java
index fa9470426c..4186df4985 100644
---
a/dubbo-common/src/test/java/org/apache/dubbo/common/store/support/SimpleDataStoreTest.java
+++
b/dubbo-common/src/test/java/org/apache/dubbo/common/store/support/SimpleDataStoreTest.java
@@ -16,9 +16,13 @@
*/
package org.apache.dubbo.common.store.support;
+import org.apache.dubbo.common.store.DataStoreUpdateListener;
+
import java.util.Map;
import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mockito;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
@@ -57,4 +61,30 @@ class SimpleDataStoreTest {
dataStore.remove("component", "key");
assertNotEquals(map, dataStore.get("component"));
}
+
+ @Test
+ void testNotify() {
+ DataStoreUpdateListener listener =
Mockito.mock(DataStoreUpdateListener.class);
+ dataStore.addListener(listener);
+
+ ArgumentCaptor<String> componentNameCaptor =
ArgumentCaptor.forClass(String.class);
+ ArgumentCaptor<String> keyCaptor =
ArgumentCaptor.forClass(String.class);
+ ArgumentCaptor<Object> valueCaptor =
ArgumentCaptor.forClass(Object.class);
+
+ dataStore.put("name", "key", "1");
+ Mockito.verify(listener).onUpdate(componentNameCaptor.capture(),
keyCaptor.capture(), valueCaptor.capture());
+ assertEquals("name", componentNameCaptor.getValue());
+ assertEquals("key", keyCaptor.getValue());
+ assertEquals("1", valueCaptor.getValue());
+
+ dataStore.remove("name", "key");
+ Mockito.verify(listener, Mockito.times(2))
+ .onUpdate(componentNameCaptor.capture(), keyCaptor.capture(),
valueCaptor.capture());
+ assertEquals("name", componentNameCaptor.getValue());
+ assertEquals("key", keyCaptor.getValue());
+ assertNull(valueCaptor.getValue());
+
+ dataStore.remove("name2", "key");
+ Mockito.verify(listener, Mockito.times(0)).onUpdate("name2", "key",
null);
+ }
}
diff --git
a/dubbo-metrics/dubbo-metrics-default/src/main/java/org/apache/dubbo/metrics/collector/sample/ThreadPoolMetricsSampler.java
b/dubbo-metrics/dubbo-metrics-default/src/main/java/org/apache/dubbo/metrics/collector/sample/ThreadPoolMetricsSampler.java
index d39e2fa2db..d7d26d448a 100644
---
a/dubbo-metrics/dubbo-metrics-default/src/main/java/org/apache/dubbo/metrics/collector/sample/ThreadPoolMetricsSampler.java
+++
b/dubbo-metrics/dubbo-metrics-default/src/main/java/org/apache/dubbo/metrics/collector/sample/ThreadPoolMetricsSampler.java
@@ -19,6 +19,7 @@ package org.apache.dubbo.metrics.collector.sample;
import org.apache.dubbo.common.logger.ErrorTypeAwareLogger;
import org.apache.dubbo.common.logger.LoggerFactory;
import org.apache.dubbo.common.store.DataStore;
+import org.apache.dubbo.common.store.DataStoreUpdateListener;
import org.apache.dubbo.common.threadpool.manager.FrameworkExecutorRepository;
import org.apache.dubbo.common.threadpool.support.AbortPolicyWithReport;
import org.apache.dubbo.common.utils.ConcurrentHashMapUtils;
@@ -42,11 +43,12 @@ import java.util.concurrent.atomic.AtomicBoolean;
import static
org.apache.dubbo.common.constants.CommonConstants.CONSUMER_SHARED_EXECUTOR_SERVICE_COMPONENT_KEY;
import static
org.apache.dubbo.common.constants.CommonConstants.EXECUTOR_SERVICE_COMPONENT_KEY;
import static
org.apache.dubbo.common.constants.LoggerCodeConstants.COMMON_METRICS_COLLECTOR_EXCEPTION;
-import static org.apache.dubbo.config.Constants.CLIENT_THREAD_POOL_NAME;
+import static org.apache.dubbo.config.Constants.CLIENT_THREAD_POOL_PREFIX;
import static org.apache.dubbo.config.Constants.SERVER_THREAD_POOL_NAME;
+import static org.apache.dubbo.config.Constants.SERVER_THREAD_POOL_PREFIX;
import static org.apache.dubbo.metrics.model.MetricsCategory.THREAD_POOL;
-public class ThreadPoolMetricsSampler implements MetricsSampler {
+public class ThreadPoolMetricsSampler implements MetricsSampler,
DataStoreUpdateListener {
private final ErrorTypeAwareLogger logger =
LoggerFactory.getErrorTypeAwareLogger(ThreadPoolMetricsSampler.class);
@@ -61,14 +63,28 @@ public class ThreadPoolMetricsSampler implements
MetricsSampler {
this.collector = collector;
}
+ @Override
+ public void onUpdate(String componentName, String key, Object value) {
+ if (EXECUTOR_SERVICE_COMPONENT_KEY.equals(componentName)) {
+ if (value instanceof ThreadPoolExecutor) {
+ addExecutors(SERVER_THREAD_POOL_PREFIX + key,
(ThreadPoolExecutor) value);
+ }
+ } else if
(CONSUMER_SHARED_EXECUTOR_SERVICE_COMPONENT_KEY.equals(componentName)) {
+ if (value instanceof ThreadPoolExecutor) {
+ addExecutors(CLIENT_THREAD_POOL_PREFIX + key,
(ThreadPoolExecutor) value);
+ }
+ }
+ }
+
public void addExecutors(String name, ExecutorService executorService) {
Optional.ofNullable(executorService)
.filter(Objects::nonNull)
.filter(e -> e instanceof ThreadPoolExecutor)
.map(e -> (ThreadPoolExecutor) e)
.ifPresent(threadPoolExecutor -> {
- sampleThreadPoolExecutor.put(name, threadPoolExecutor);
- samplesChanged.set(true);
+ if (sampleThreadPoolExecutor.put(name, threadPoolExecutor)
== null) {
+ samplesChanged.set(true);
+ }
});
}
@@ -152,18 +168,20 @@ public class ThreadPoolMetricsSampler implements
MetricsSampler {
}
if (dataStore != null) {
+ dataStore.addListener(this);
+
Map<String, Object> executors =
dataStore.get(EXECUTOR_SERVICE_COMPONENT_KEY);
for (Map.Entry<String, Object> entry : executors.entrySet()) {
ExecutorService executor = (ExecutorService) entry.getValue();
if (executor instanceof ThreadPoolExecutor) {
- this.addExecutors(SERVER_THREAD_POOL_NAME + "-" +
entry.getKey(), executor);
+ this.addExecutors(SERVER_THREAD_POOL_PREFIX +
entry.getKey(), executor);
}
}
executors =
dataStore.get(CONSUMER_SHARED_EXECUTOR_SERVICE_COMPONENT_KEY);
for (Map.Entry<String, Object> entry : executors.entrySet()) {
ExecutorService executor = (ExecutorService) entry.getValue();
if (executor instanceof ThreadPoolExecutor) {
- this.addExecutors(CLIENT_THREAD_POOL_NAME + "-" +
entry.getKey(), executor);
+ this.addExecutors(CLIENT_THREAD_POOL_PREFIX +
entry.getKey(), executor);
}
}
diff --git
a/dubbo-metrics/dubbo-metrics-default/src/test/java/org/apache/dubbo/metrics/collector/sample/ThreadPoolMetricsSamplerTest.java
b/dubbo-metrics/dubbo-metrics-default/src/test/java/org/apache/dubbo/metrics/collector/sample/ThreadPoolMetricsSamplerTest.java
index 6b20d6c69d..0d6bcec95d 100644
---
a/dubbo-metrics/dubbo-metrics-default/src/test/java/org/apache/dubbo/metrics/collector/sample/ThreadPoolMetricsSamplerTest.java
+++
b/dubbo-metrics/dubbo-metrics-default/src/test/java/org/apache/dubbo/metrics/collector/sample/ThreadPoolMetricsSamplerTest.java
@@ -19,6 +19,7 @@ package org.apache.dubbo.metrics.collector.sample;
import org.apache.dubbo.common.beans.factory.ScopeBeanFactory;
import org.apache.dubbo.common.extension.ExtensionLoader;
import org.apache.dubbo.common.store.DataStore;
+import org.apache.dubbo.common.store.DataStoreUpdateListener;
import org.apache.dubbo.common.threadpool.manager.FrameworkExecutorRepository;
import org.apache.dubbo.metrics.collector.DefaultMetricsCollector;
import org.apache.dubbo.metrics.model.ThreadPoolMetric;
@@ -37,11 +38,13 @@ import java.util.concurrent.ThreadPoolExecutor;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
import static
org.apache.dubbo.common.constants.CommonConstants.CONSUMER_SHARED_EXECUTOR_SERVICE_COMPONENT_KEY;
import static
org.apache.dubbo.common.constants.CommonConstants.EXECUTOR_SERVICE_COMPONENT_KEY;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@SuppressWarnings("all")
@@ -178,4 +181,32 @@ public class ThreadPoolMetricsSamplerTest {
serverExecutor.shutdown();
clientExecutor.shutdown();
}
+
+ @Test
+ void testDataSourceNotify() throws Exception {
+ ArgumentCaptor<DataStoreUpdateListener> captor =
ArgumentCaptor.forClass(DataStoreUpdateListener.class);
+
when(scopeBeanFactory.getBean(FrameworkExecutorRepository.class)).thenReturn(frameworkExecutorRepository);
+ when(frameworkExecutorRepository.getSharedExecutor()).thenReturn(null);
+ sampler2.registryDefaultSampleThreadPoolExecutor();
+
+ Field f =
ThreadPoolMetricsSampler.class.getDeclaredField("sampleThreadPoolExecutor");
+ f.setAccessible(true);
+ Map<String, ThreadPoolExecutor> executors = (Map<String,
ThreadPoolExecutor>) f.get(sampler2);
+
+ Assertions.assertEquals(0, executors.size());
+
+ verify(dataStore).addListener(captor.capture());
+ Assertions.assertEquals(sampler2, captor.getValue());
+
+ ExecutorService executorService = Executors.newFixedThreadPool(5);
+ sampler2.onUpdate(EXECUTOR_SERVICE_COMPONENT_KEY, "20880",
executorService);
+
+ executors = (Map<String, ThreadPoolExecutor>) f.get(sampler2);
+ Assertions.assertEquals(1, executors.size());
+
Assertions.assertTrue(executors.containsKey("DubboServerHandler-20880"));
+
+ sampler2.onUpdate(CONSUMER_SHARED_EXECUTOR_SERVICE_COMPONENT_KEY,
"client", executorService);
+ Assertions.assertEquals(2, executors.size());
+
Assertions.assertTrue(executors.containsKey("DubboClientHandler-client"));
+ }
}