This is an automated email from the ASF dual-hosted git repository.

lizhimins pushed a commit to branch rocketmq-studio
in repository https://gitbox.apache.org/repos/asf/rocketmq-dashboard.git


The following commit(s) were added to refs/heads/rocketmq-studio by this push:
     new 08587b9d perf(dashboard): count topics and groups per master broker 
(#1683)
08587b9d is described below

commit 08587b9dad17216df035b6647bb5520812d07450
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 11 21:31:25 2026 +0800

    perf(dashboard): count topics and groups per master broker (#1683)
    
    * fix: bound dashboard topic counting by broker
    
    * test: cover dashboard topic deduplication
    
    * fix: deduplicate dashboard consumer groups
    
    * fix: surface unavailable dashboard group counts
---
 .../provider/apache/RocketMQDashboardProvider.java |  88 +++++++------
 .../apache/RocketMQDashboardProviderTest.java      | 137 +++++++++++++++++++++
 2 files changed, 179 insertions(+), 46 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProvider.java
index 5db3dd08..92c6c9d3 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProvider.java
@@ -26,10 +26,8 @@ import java.util.Set;
 import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
 import org.apache.rocketmq.remoting.protocol.body.KVTable;
 import org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
-import org.apache.rocketmq.remoting.protocol.body.TopicList;
+import org.apache.rocketmq.remoting.protocol.body.TopicConfigSerializeWrapper;
 import org.apache.rocketmq.remoting.protocol.route.BrokerData;
-import org.apache.rocketmq.remoting.protocol.route.QueueData;
-import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
 import 
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
 import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
@@ -149,48 +147,36 @@ public class RocketMQDashboardProvider implements 
DashboardProvider {
                 }
             }
 
-            // Count topics globally and per-cluster
-            Map<String, Integer> topicsByCluster = new HashMap<>();
-            try {
-                TopicList topicList = admin.fetchAllTopicList();
-                Set<String> topics = topicList == null || 
topicList.getTopicList() == null
-                        ? Set.of() : topicList.getTopicList();
-                totalTopics = (int) topics.stream()
-                        .filter(t -> !isSystemTopic(t))
-                        .count();
-                // For each non-system topic, determine its cluster via route 
data
-                for (String topicName : topics) {
-                    if (isSystemTopic(topicName)) {
-                        continue;
-                    }
-                    try {
-                        TopicRouteData route = 
admin.examineTopicRouteInfo(topicName);
-                        if (route != null && route.getQueueDatas() != null) {
-                            for (QueueData qd : route.getQueueDatas()) {
-                                BrokerData bd = 
brokerAddrTable.get(qd.getBrokerName());
-                                if (bd != null && bd.getBrokerAddrs() != null) 
{
-                                    String addr = bd.getBrokerAddrs().get(0L);
-                                    String cluster = 
brokerAddrToCluster.get(addr);
-                                    if (cluster != null) {
-                                        topicsByCluster.merge(cluster, 1, 
Integer::sum);
-                                    }
-                                    break; // one queue is enough to determine 
the cluster
-                                }
-                            }
-                        }
-                    } catch (Exception ignored) {
-                        // Skip topics whose route data is unavailable
-                    }
-                }
-            } catch (Exception e) {
-                log.warn("Failed to fetch topic list: {}", e.getMessage());
-            }
-
-            // Count subscription groups globally and per-cluster, and collect 
TPS from each master broker
+            // Count topics and subscription groups from each master broker. 
Topic metadata is
+            // available in the broker config response, so this remains 
bounded by broker count
+            // instead of resolving a NameServer route for every Topic.
+            Set<String> allTopics = new HashSet<>();
+            Map<String, Set<String>> topicsByCluster = new HashMap<>();
+            Set<String> topicCountsUnavailableClusters = new HashSet<>();
             Set<String> allGroups = new HashSet<>();
-            Map<String, Integer> groupsByCluster = new HashMap<>();
+            Map<String, Set<String>> groupsByCluster = new HashMap<>();
+            Set<String> groupCountsUnavailableClusters = new HashSet<>();
             for (String brokerAddr : masterAddrs) {
                 String clusterName = brokerAddrToCluster.get(brokerAddr);
+                try {
+                    TopicConfigSerializeWrapper topicConfig = 
admin.getAllTopicConfig(brokerAddr, 5000);
+                    if (topicConfig == null || 
topicConfig.getTopicConfigTable() == null) {
+                        markCountUnavailable(topicCountsUnavailableClusters, 
clusterName);
+                    } else {
+                        Set<String> clusterTopics = 
topicsByCluster.computeIfAbsent(
+                                clusterName, ignored -> new HashSet<>());
+                        topicConfig.getTopicConfigTable().keySet().stream()
+                                .filter(topic -> !isSystemTopic(topic))
+                                .forEach(topic -> {
+                                    allTopics.add(topic);
+                                    clusterTopics.add(topic);
+                                });
+                    }
+                } catch (Exception e) {
+                    markCountUnavailable(topicCountsUnavailableClusters, 
clusterName);
+                    log.warn("Failed to get topic config from broker {}: {}", 
brokerAddr, e.getMessage());
+                }
+
                 try {
                     SubscriptionGroupWrapper subscriptionGroupWrapper = 
admin.getAllSubscriptionGroup(brokerAddr, 5000);
                     if (subscriptionGroupWrapper != null && 
subscriptionGroupWrapper.getSubscriptionGroupTable() != null) {
@@ -200,12 +186,14 @@ public class RocketMQDashboardProvider implements 
DashboardProvider {
                             if (!isSystemGroup(groupName)) {
                                 allGroups.add(groupName);
                                 if (clusterName != null) {
-                                    groupsByCluster.merge(clusterName, 1, 
Integer::sum);
+                                    
groupsByCluster.computeIfAbsent(clusterName, ignored -> new HashSet<>())
+                                            .add(groupName);
                                 }
                             }
                         }
                     }
                 } catch (Exception e) {
+                    markCountUnavailable(groupCountsUnavailableClusters, 
clusterName);
                     log.warn("Failed to get subscription groups from broker 
{}: {}", brokerAddr, e.getMessage());
                 }
 
@@ -230,6 +218,7 @@ public class RocketMQDashboardProvider implements 
DashboardProvider {
                     log.warn("Failed to get runtime info from broker {}: {}", 
brokerAddr, e.getMessage());
                 }
             }
+            totalTopics = allTopics.size();
             totalGroups = allGroups.size();
 
             // Build per-cluster overview
@@ -240,7 +229,8 @@ public class RocketMQDashboardProvider implements 
DashboardProvider {
                 long clusterTpsIn = 0;
                 long clusterTpsOut = 0;
                 String version = "unknown";
-                boolean runtimeMetricsUnavailable = false;
+                boolean runtimeMetricsUnavailable = 
topicCountsUnavailableClusters.contains(clusterName)
+                        || 
groupCountsUnavailableClusters.contains(clusterName);
 
                 for (String brokerName : brokerNames) {
                     BrokerData brokerData = brokerAddrTable.get(brokerName);
@@ -273,8 +263,8 @@ public class RocketMQDashboardProvider implements 
DashboardProvider {
                         .status(runtimeMetricsUnavailable ? 
ClusterStatus.warning : ClusterStatus.healthy)
                         .brokers(clusterBrokers)
                         .proxies(0)
-                        .topics(topicsByCluster.getOrDefault(clusterName, 0))
-                        .groups(groupsByCluster.getOrDefault(clusterName, 0))
+                        .topics(topicsByCluster.getOrDefault(clusterName, 
Set.of()).size())
+                        .groups(groupsByCluster.getOrDefault(clusterName, 
Set.of()).size())
                         .tpsIn(clusterTpsIn)
                         .tpsOut(clusterTpsOut)
                         .version(version)
@@ -369,4 +359,10 @@ public class RocketMQDashboardProvider implements 
DashboardProvider {
     private boolean isSystemGroup(String group) {
         return SystemGroupFilter.isSystem(group);
     }
+
+    private void markCountUnavailable(Set<String> unavailableClusters, String 
clusterName) {
+        if (clusterName != null) {
+            unavailableClusters.add(clusterName);
+        }
+    }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProviderTest.java
index 1849931f..2e40d50a 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDashboardProviderTest.java
@@ -20,11 +20,16 @@ import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 
+import org.apache.rocketmq.common.TopicConfig;
 import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
 import org.apache.rocketmq.remoting.protocol.body.KVTable;
+import org.apache.rocketmq.remoting.protocol.body.SubscriptionGroupWrapper;
 import org.apache.rocketmq.remoting.protocol.body.TopicList;
+import org.apache.rocketmq.remoting.protocol.body.TopicConfigSerializeWrapper;
 import org.apache.rocketmq.remoting.protocol.route.BrokerData;
+import 
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
 import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
@@ -42,6 +47,7 @@ import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
@@ -53,6 +59,7 @@ class RocketMQDashboardProviderTest {
         DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
         when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
         when(adminExt.fetchAllTopicList()).thenReturn(topicList());
+        when(adminExt.getAllTopicConfig("10.0.0.11:10911", 
5000)).thenReturn(topicConfig("order-topic"));
         
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtimeStats());
 
         RocketMQDashboardProvider provider = newProvider(adminExt);
@@ -66,6 +73,97 @@ class RocketMQDashboardProviderTest {
         verify(adminExt, times(1)).fetchBrokerRuntimeStats("10.0.0.11:10911");
     }
 
+    @Test
+    void 
dashboardShouldCountTopicsFromBrokerConfigWithoutPerTopicRouteRequests() throws 
Exception {
+        DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+        when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
+        when(adminExt.getAllTopicConfig("10.0.0.11:10911", 5000))
+                .thenReturn(topicConfig("order-topic", "payments", 
"SCHEDULE_TOPIC_XXXX"));
+        
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtimeStats());
+
+        DashboardDataVO dashboard = newProvider(adminExt).getDashboardData();
+
+        assertThat(dashboard.getStats().getTotalTopics()).isEqualTo(2);
+        assertThat(dashboard.getClusters()).singleElement().satisfies(cluster 
-> {
+            assertThat(cluster.getTopics()).isEqualTo(2);
+            assertThat(cluster.getStatus()).isEqualTo(ClusterStatus.healthy);
+        });
+        verify(adminExt).getAllTopicConfig("10.0.0.11:10911", 5000);
+        verify(adminExt, never()).examineTopicRouteInfo(anyString());
+    }
+
+    @Test
+    void dashboardShouldDeduplicateTopicsReportedByMultipleClusterBrokers() 
throws Exception {
+        DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+        
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithTwoMasters());
+        when(adminExt.getAllTopicConfig("10.0.0.11:10911", 5000))
+                .thenReturn(topicConfig("orders", "payments"));
+        when(adminExt.getAllTopicConfig("10.0.0.12:10911", 5000))
+                .thenReturn(topicConfig("orders", "inventory"));
+        
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtimeStats());
+        
when(adminExt.fetchBrokerRuntimeStats("10.0.0.12:10911")).thenReturn(runtimeStats());
+
+        DashboardDataVO dashboard = newProvider(adminExt).getDashboardData();
+
+        assertThat(dashboard.getStats().getTotalTopics()).isEqualTo(3);
+        assertThat(dashboard.getClusters()).singleElement()
+                .extracting(cluster -> cluster.getTopics())
+                .isEqualTo(3);
+        verify(adminExt, never()).examineTopicRouteInfo(anyString());
+    }
+
+    @Test
+    void dashboardShouldDeduplicateGroupsReportedByMultipleClusterBrokers() 
throws Exception {
+        DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+        
when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfoWithTwoMasters());
+        when(adminExt.getAllSubscriptionGroup("10.0.0.11:10911", 5000))
+                .thenReturn(subscriptionGroups("cg-orders", "cg-payments"));
+        when(adminExt.getAllSubscriptionGroup("10.0.0.12:10911", 5000))
+                .thenReturn(subscriptionGroups("cg-orders", "cg-inventory"));
+        
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtimeStats());
+        
when(adminExt.fetchBrokerRuntimeStats("10.0.0.12:10911")).thenReturn(runtimeStats());
+
+        DashboardDataVO dashboard = newProvider(adminExt).getDashboardData();
+
+        assertThat(dashboard.getStats().getTotalConsumerGroups()).isEqualTo(3);
+        assertThat(dashboard.getClusters()).singleElement()
+                .extracting(cluster -> cluster.getGroups())
+                .isEqualTo(3);
+    }
+
+    @Test
+    void dashboardShouldMarkClusterWarningWhenTopicCountsAreUnavailable() 
throws Exception {
+        DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+        when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
+        when(adminExt.getAllTopicConfig("10.0.0.11:10911", 5000))
+                .thenThrow(new RuntimeException("broker unavailable"));
+        
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtimeStats());
+
+        DashboardDataVO dashboard = newProvider(adminExt).getDashboardData();
+
+        assertThat(dashboard.getClusters()).singleElement()
+                .extracting(cluster -> cluster.getStatus())
+                .isEqualTo(ClusterStatus.warning);
+        assertThat(dashboard.getStats().getHealthyClusters()).isZero();
+    }
+
+    @Test
+    void dashboardShouldMarkClusterWarningWhenGroupCountsAreUnavailable() 
throws Exception {
+        DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+        when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
+        when(adminExt.getAllTopicConfig("10.0.0.11:10911", 
5000)).thenReturn(topicConfig("orders"));
+        when(adminExt.getAllSubscriptionGroup("10.0.0.11:10911", 5000))
+                .thenThrow(new RuntimeException("broker unavailable"));
+        
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtimeStats());
+
+        DashboardDataVO dashboard = newProvider(adminExt).getDashboardData();
+
+        assertThat(dashboard.getClusters()).singleElement()
+                .extracting(cluster -> cluster.getStatus())
+                .isEqualTo(ClusterStatus.warning);
+        assertThat(dashboard.getStats().getHealthyClusters()).isZero();
+    }
+
     @Test
     void dashboardShouldSurviveNullTopologyTables() throws Exception {
         DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
@@ -258,12 +356,51 @@ class RocketMQDashboardProviderTest {
         return info;
     }
 
+    private ClusterInfo clusterInfoWithTwoMasters() {
+        ClusterInfo info = new ClusterInfo();
+        HashMap<String, BrokerData> brokerAddrTable = new HashMap<>();
+        brokerAddrTable.put("broker-a", brokerData("broker-a", 
"10.0.0.11:10911"));
+        brokerAddrTable.put("broker-b", brokerData("broker-b", 
"10.0.0.12:10911"));
+        info.setBrokerAddrTable(brokerAddrTable);
+
+        HashMap<String, Set<String>> clusterAddrTable = new HashMap<>();
+        clusterAddrTable.put("DefaultCluster", Set.of("broker-a", "broker-b"));
+        info.setClusterAddrTable(clusterAddrTable);
+        return info;
+    }
+
+    private BrokerData brokerData(String name, String address) {
+        HashMap<Long, String> addresses = new HashMap<>();
+        addresses.put(0L, address);
+        return new BrokerData("DefaultCluster", name, addresses);
+    }
+
     private TopicList topicList() {
         TopicList topicList = new TopicList();
         topicList.setTopicList(Set.of("order-topic"));
         return topicList;
     }
 
+    private TopicConfigSerializeWrapper topicConfig(String... names) {
+        TopicConfigSerializeWrapper wrapper = new 
TopicConfigSerializeWrapper();
+        ConcurrentHashMap<String, TopicConfig> configs = new 
ConcurrentHashMap<>();
+        for (String name : names) {
+            configs.put(name, new TopicConfig(name));
+        }
+        wrapper.setTopicConfigTable(configs);
+        return wrapper;
+    }
+
+    private SubscriptionGroupWrapper subscriptionGroups(String... names) {
+        SubscriptionGroupWrapper wrapper = new SubscriptionGroupWrapper();
+        ConcurrentHashMap<String, SubscriptionGroupConfig> groups = new 
ConcurrentHashMap<>();
+        for (String name : names) {
+            groups.put(name, new SubscriptionGroupConfig());
+        }
+        wrapper.setSubscriptionGroupTable(groups);
+        return wrapper;
+    }
+
     private KVTable runtimeStats() {
         return runtimeStats("2.0", "5.0");
     }

Reply via email to