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"));
+    }
 }

Reply via email to