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