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/DataStoreUpdateListener.java b/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStoreUpdateListener.java new file mode 100644 index 0000000000..de994191f7 --- /dev/null +++ b/dubbo-common/src/main/java/org/apache/dubbo/common/store/DataStoreUpdateListener.java @@ -0,0 +1,21 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.dubbo.common.store; + +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); // > private final ConcurrentMap> data = new ConcurrentHashMap<>(); + private final ConcurrentHashSet listeners = new ConcurrentHashSet<>(); @Override public Map get(String componentName) { @@ -52,6 +59,7 @@ public class SimpleDataStore implements DataStore { Map 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 componentNameCaptor = ArgumentCaptor.forClass(String.class); + ArgumentCaptor keyCaptor = ArgumentCaptor.forClass(String.class); + ArgumentCaptor 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 executors = dataStore.get(EXECUTOR_SERVICE_COMPONENT_KEY); for (Map.Entry 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 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 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 executors = (Map) 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) 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")); + } }