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 4d5964639 feat(cluster): show per-broker daily message counters (#3997)
4d5964639 is described below
commit 4d5964639add9facf32520f5f558d1a50afbb144
Author: 烤化の初雪 <[email protected]>
AuthorDate: Tue Sep 8 15:57:03 2026 +0800
feat(cluster): show per-broker daily message counters (#3997)
Co-authored-by: unbridled-41
<[email protected]>
---
.../rocketmq/studio/cluster/broker/BrokerVO.java | 4 +++
.../cluster/broker/ClusterRepositoryImpl.java | 8 +++++
.../studio/common/util/BrokerRuntimeStats.java | 25 +++++++++++++
.../provider/apache/RocketMQClusterProvider.java | 10 ++++++
.../provider/apache/RocketMQDashboardProvider.java | 23 ++----------
.../cluster/broker/ClusterRepositoryImplTest.java | 12 +++++++
.../studio/common/util/BrokerRuntimeStatsTest.java | 36 +++++++++++++++++++
.../apache/RocketMQClusterProviderTest.java | 41 ++++++++++++++++++++++
web/src/api/cluster.ts | 4 +++
web/src/i18n/translations.ts | 4 +++
.../pages/cluster/__tests__/ClusterPage.test.tsx | 18 ++++++++++
web/src/pages/cluster/index.tsx | 36 +++++++++++++++++++
12 files changed, 200 insertions(+), 21 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/BrokerVO.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/BrokerVO.java
index 8992122b0..91e402670 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/BrokerVO.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/BrokerVO.java
@@ -34,6 +34,10 @@ public class BrokerVO {
private double diskUsage;
private long tpsIn;
private long tpsOut;
+ private long putMessagesToday;
+ private long putMessagesYesterday;
+ private long getMessagesToday;
+ private long getMessagesYesterday;
@Builder.Default
private boolean runtimeStatsAvailable = true;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterRepositoryImpl.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterRepositoryImpl.java
index a07d0c5ed..3a5731c05 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterRepositoryImpl.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/ClusterRepositoryImpl.java
@@ -121,6 +121,10 @@ public class ClusterRepositoryImpl implements
ClusterRepository {
.diskUsage(broker.getDiskUsage())
.tpsIn(broker.getTpsIn())
.tpsOut(broker.getTpsOut())
+ .putMessagesToday(broker.getPutMessagesToday())
+ .putMessagesYesterday(broker.getPutMessagesYesterday())
+ .getMessagesToday(broker.getGetMessagesToday())
+ .getMessagesYesterday(broker.getGetMessagesYesterday())
.runtimeStatsAvailable(broker.isRuntimeStatsAvailable())
.build();
}
@@ -196,6 +200,10 @@ public class ClusterRepositoryImpl implements
ClusterRepository {
.diskUsage(45.2)
.tpsIn(1200)
.tpsOut(800)
+ .putMessagesToday(1500)
+ .putMessagesYesterday(1200)
+ .getMessagesToday(1300)
+ .getMessagesYesterday(1000)
.build(),
BrokerVO.builder()
.name("broker-b")
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/common/util/BrokerRuntimeStats.java
b/server/src/main/java/org/apache/rocketmq/studio/common/util/BrokerRuntimeStats.java
index 925440cbe..f7d44d3e0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/common/util/BrokerRuntimeStats.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/common/util/BrokerRuntimeStats.java
@@ -51,4 +51,29 @@ public final class BrokerRuntimeStats {
}
return runtimeStats.get(OUTBOUND_TPS_4X);
}
+
+ /**
+ * Computes a daily counter from the morning-snapshot pair a broker
publishes in its runtime
+ * stats (e.g. {@code msgPutTotalTodayMorning}/{@code
msgPutTotalTodayNow}): the delta is the
+ * current value minus the morning snapshot, clamped at zero so a broker
restart (which resets
+ * the counters) yields 0 instead of a negative number. Returns 0 when
either key is missing,
+ * negative, or unparseable.
+ */
+ public static long dailyCounterDelta(Map<String, String> runtimeStats,
String morningKey, String nowKey) {
+ String morningValue = runtimeStats == null ? null :
runtimeStats.get(morningKey);
+ String nowValue = runtimeStats == null ? null :
runtimeStats.get(nowKey);
+ if (morningValue == null || nowValue == null) {
+ return 0L;
+ }
+ try {
+ long morning = Long.parseLong(morningValue.trim());
+ long now = Long.parseLong(nowValue.trim());
+ if (morning < 0 || now < 0) {
+ return 0L;
+ }
+ return Math.max(0L, now - morning);
+ } catch (NumberFormatException exception) {
+ return 0L;
+ }
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProvider.java
index 1e6332a02..6bd320ebf 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProvider.java
@@ -286,6 +286,16 @@ public class RocketMQClusterProvider implements
ClusterProvider {
// keep default
}
}
+
+ // Daily message counters, mirroring the runtime stat keys the
broker publishes.
+
builder.putMessagesToday(BrokerRuntimeStats.dailyCounterDelta(table,
+ "msgPutTotalTodayMorning", "msgPutTotalTodayNow"));
+
builder.putMessagesYesterday(BrokerRuntimeStats.dailyCounterDelta(table,
+ "msgPutTotalYesterdayMorning", "msgPutTotalTodayMorning"));
+
builder.getMessagesToday(BrokerRuntimeStats.dailyCounterDelta(table,
+ "msgGetTotalTodayMorning", "msgGetTotalTodayNow"));
+
builder.getMessagesYesterday(BrokerRuntimeStats.dailyCounterDelta(table,
+ "msgGetTotalYesterdayMorning", "msgGetTotalTodayMorning"));
return true;
} catch (Exception e) {
log.warn("Failed to get runtime info for broker at {}: {}",
brokerAddr, e.getMessage());
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 00979692c..58b47ec22 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
@@ -345,7 +345,8 @@ public class RocketMQDashboardProvider implements
DashboardProvider {
tpsIn += parseTps(table.get("putTps"));
tpsOut +=
parseTps(BrokerRuntimeStats.outboundTps(table));
- messagesToday += parseMessagesToday(table);
+ messagesToday +=
BrokerRuntimeStats.dailyCounterDelta(table,
+ "msgPutTotalTodayMorning",
"msgPutTotalTodayNow");
}
} catch (Exception e) {
log.warn("Failed to get runtime info from broker {}: {}",
brokerAddr, e.getMessage());
@@ -508,26 +509,6 @@ public class RocketMQDashboardProvider implements
DashboardProvider {
return (long) parsed;
}
- private long parseMessagesToday(Map<String, String> runtimeStats) {
- String morningValue = runtimeStats.get("msgPutTotalTodayMorning");
- String currentValue = runtimeStats.get("msgPutTotalTodayNow");
- if (morningValue == null || currentValue == null) {
- return 0;
- }
- try {
- long morning = Long.parseLong(morningValue.trim());
- long current = Long.parseLong(currentValue.trim());
- if (morning < 0 || current < 0) {
- return 0;
- }
- return Math.max(0, current - morning);
- } catch (NumberFormatException exception) {
- log.debug("Failed to parse today's message counters: morning={},
current={}",
- morningValue, currentValue);
- return 0;
- }
- }
-
private boolean isSystemTopic(String topic) {
return SystemTopicFilter.isSystem(topic);
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterRepositoryImplTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterRepositoryImplTest.java
index 721b9a1c5..fd63f41b7 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterRepositoryImplTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/ClusterRepositoryImplTest.java
@@ -69,6 +69,18 @@ class ClusterRepositoryImplTest {
assertThat(second.getConfig().getFileReservedTime()).isEqualTo(72);
}
+ @Test
+ void findByIdShouldPreserveBrokerDailyMessageCounters() {
+ ClusterRepositoryImpl repository = new ClusterRepositoryImpl(true);
+
+ BrokerVO broker =
repository.findById("cluster-001").orElseThrow().getBrokers().get(0);
+
+ assertThat(broker.getPutMessagesToday()).isEqualTo(1500);
+ assertThat(broker.getPutMessagesYesterday()).isEqualTo(1200);
+ assertThat(broker.getGetMessagesToday()).isEqualTo(1300);
+ assertThat(broker.getGetMessagesYesterday()).isEqualTo(1000);
+ }
+
@Test
void updateConfigShouldNotRetainCallerOwnedObject() {
ClusterRepositoryImpl repository = new ClusterRepositoryImpl(true);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/common/util/BrokerRuntimeStatsTest.java
b/server/src/test/java/org/apache/rocketmq/studio/common/util/BrokerRuntimeStatsTest.java
index 0ea1c4c78..2f33b8b56 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/common/util/BrokerRuntimeStatsTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/common/util/BrokerRuntimeStatsTest.java
@@ -53,4 +53,40 @@ class BrokerRuntimeStatsTest {
assertThat(BrokerRuntimeStats.outboundTps(new HashMap<>())).isNull();
assertThat(BrokerRuntimeStats.outboundTps(null)).isNull();
}
+
+ @Test
+ void dailyCounterDeltaSubtractsMorningSnapshotTest() {
+ Map<String, String> stats = new HashMap<>();
+ stats.put("msgPutTotalTodayMorning", "1400");
+ stats.put("msgPutTotalTodayNow", "2000");
+ assertThat(BrokerRuntimeStats.dailyCounterDelta(stats,
+ "msgPutTotalTodayMorning",
"msgPutTotalTodayNow")).isEqualTo(600L);
+ }
+
+ @Test
+ void dailyCounterDeltaClampsRestartResetToZeroTest() {
+ Map<String, String> stats = new HashMap<>();
+ stats.put("msgPutTotalTodayMorning", "1400");
+ stats.put("msgPutTotalTodayNow", "300");
+ assertThat(BrokerRuntimeStats.dailyCounterDelta(stats,
+ "msgPutTotalTodayMorning", "msgPutTotalTodayNow")).isZero();
+ }
+
+ @Test
+ void dailyCounterDeltaReturnsZeroWhenMissingOrUnparseableTest() {
+ assertThat(BrokerRuntimeStats.dailyCounterDelta(new HashMap<>(),
+ "msgPutTotalTodayMorning", "msgPutTotalTodayNow")).isZero();
+ assertThat(BrokerRuntimeStats.dailyCounterDelta(null,
+ "msgPutTotalTodayMorning", "msgPutTotalTodayNow")).isZero();
+ Map<String, String> stats = new HashMap<>();
+ stats.put("msgPutTotalTodayMorning", "1400");
+ stats.put("msgPutTotalTodayNow", "not-a-number");
+ assertThat(BrokerRuntimeStats.dailyCounterDelta(stats,
+ "msgPutTotalTodayMorning", "msgPutTotalTodayNow")).isZero();
+ Map<String, String> negative = new HashMap<>();
+ negative.put("msgPutTotalTodayMorning", "-5");
+ negative.put("msgPutTotalTodayNow", "10");
+ assertThat(BrokerRuntimeStats.dailyCounterDelta(negative,
+ "msgPutTotalTodayMorning", "msgPutTotalTodayNow")).isZero();
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
index 679bc5a0d..9fbbeade0 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
@@ -106,6 +106,47 @@ class RocketMQClusterProviderTest {
.isEqualTo(62.5D);
}
+ @Test
+ void discoverClustersShouldParseDailyMessageCounters() throws Exception {
+ DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+ RocketMQClusterProvider provider = newProvider(adminExt);
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
+ KVTable runtime = runtimeStats();
+ runtime.getTable().put("msgPutTotalYesterdayMorning", "1000");
+ runtime.getTable().put("msgPutTotalTodayMorning", "1400");
+ runtime.getTable().put("msgPutTotalTodayNow", "2000");
+ runtime.getTable().put("msgGetTotalYesterdayMorning", "800");
+ runtime.getTable().put("msgGetTotalTodayMorning", "1200");
+ runtime.getTable().put("msgGetTotalTodayNow", "1500");
+
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtime);
+
+ ClusterVO cluster = provider.discoverClusters().get(0);
+
+ assertThat(cluster.getBrokers()).singleElement()
+
.extracting(org.apache.rocketmq.studio.cluster.broker.BrokerVO::getPutMessagesToday,
+
org.apache.rocketmq.studio.cluster.broker.BrokerVO::getPutMessagesYesterday,
+
org.apache.rocketmq.studio.cluster.broker.BrokerVO::getGetMessagesToday,
+
org.apache.rocketmq.studio.cluster.broker.BrokerVO::getGetMessagesYesterday)
+ .containsExactly(600L, 400L, 300L, 400L);
+ }
+
+ @Test
+ void discoverClustersShouldDefaultDailyMessageCountersWhenKeysAreMissing()
throws Exception {
+ DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
+ RocketMQClusterProvider provider = newProvider(adminExt);
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo());
+
when(adminExt.fetchBrokerRuntimeStats("10.0.0.11:10911")).thenReturn(runtimeStats());
+
+ ClusterVO cluster = provider.discoverClusters().get(0);
+
+ assertThat(cluster.getBrokers()).singleElement()
+
.extracting(org.apache.rocketmq.studio.cluster.broker.BrokerVO::getPutMessagesToday,
+
org.apache.rocketmq.studio.cluster.broker.BrokerVO::getPutMessagesYesterday,
+
org.apache.rocketmq.studio.cluster.broker.BrokerVO::getGetMessagesToday,
+
org.apache.rocketmq.studio.cluster.broker.BrokerVO::getGetMessagesYesterday)
+ .containsExactly(0L, 0L, 0L, 0L);
+ }
+
@Test
void discoverClustersShouldPopulateSafeDefaults() throws Exception {
DefaultMQAdminExt adminExt = mock(DefaultMQAdminExt.class);
diff --git a/web/src/api/cluster.ts b/web/src/api/cluster.ts
index e83f09782..806591ede 100644
--- a/web/src/api/cluster.ts
+++ b/web/src/api/cluster.ts
@@ -44,6 +44,10 @@ export interface BrokerInfo {
tpsIn: number;
tpsOut: number;
diskUsage: number;
+ putMessagesToday?: number;
+ putMessagesYesterday?: number;
+ getMessagesToday?: number;
+ getMessagesYesterday?: number;
version?: string | null;
runtimeStatsAvailable?: boolean;
}
diff --git a/web/src/i18n/translations.ts b/web/src/i18n/translations.ts
index 353df823b..a66f67fba 100644
--- a/web/src/i18n/translations.ts
+++ b/web/src/i18n/translations.ts
@@ -122,6 +122,10 @@ const translations: Record<string, Record<Lang, string>> =
{
'cluster.brokerClusterName': { zh: 'Broker 集群名称', en: 'Broker Cluster Name'
},
'cluster.brokerName': { zh: 'Broker 名称', en: 'Broker Name' },
'cluster.diskUsage': { zh: '磁盘使用', en: 'Disk Usage' },
+ 'cluster.putMessagesToday': { zh: '今日写入', en: 'Put Today' },
+ 'cluster.putMessagesYesterday': { zh: '昨日写入', en: 'Put Yesterday' },
+ 'cluster.getMessagesToday': { zh: '今日消费', en: 'Get Today' },
+ 'cluster.getMessagesYesterday': { zh: '昨日消费', en: 'Get Yesterday' },
'cluster.proxyAddr': { zh: 'Proxy 地址', en: 'Proxy Address' },
'cluster.connections': { zh: '连接数', en: 'Connections' },
'cluster.grpcPort': { zh: 'gRPC 端口', en: 'gRPC Port' },
diff --git a/web/src/pages/cluster/__tests__/ClusterPage.test.tsx
b/web/src/pages/cluster/__tests__/ClusterPage.test.tsx
index 2d501c6f3..9b0c5f49c 100644
--- a/web/src/pages/cluster/__tests__/ClusterPage.test.tsx
+++ b/web/src/pages/cluster/__tests__/ClusterPage.test.tsx
@@ -113,6 +113,10 @@ const buildCluster = ({
diskUsage: 62,
tpsIn,
tpsOut,
+ putMessagesToday: 1234,
+ putMessagesYesterday: 1100,
+ getMessagesToday: 980,
+ getMessagesYesterday: 900,
},
{
name: 'rocketmq-prod-1',
@@ -411,6 +415,20 @@ describe('Cluster page', () => {
expect(within(dialog).getByRole('row', { name: /写队列数/
})).toHaveTextContent('16');
});
+ it('renders per-broker daily message counters in the broker tab', async ()
=> {
+ renderWithProviders(<ClusterPage />);
+
+ expect(await screen.findByText('rocketmq-prod-0')).toBeInTheDocument();
+ expect(screen.getByRole('columnheader', { name: '今日写入'
})).toBeInTheDocument();
+ expect(screen.getByRole('columnheader', { name: '昨日写入'
})).toBeInTheDocument();
+ expect(screen.getByRole('columnheader', { name: '今日消费'
})).toBeInTheDocument();
+ expect(screen.getByRole('columnheader', { name: '昨日消费'
})).toBeInTheDocument();
+ expect(screen.getByText('1,234')).toBeInTheDocument();
+ expect(screen.getByText('1,100')).toBeInTheDocument();
+ expect(screen.getByText('980')).toBeInTheDocument();
+ expect(screen.getByText('900')).toBeInTheDocument();
+ });
+
it('keeps cluster tabs usable when address fields are missing', async () => {
const user = userEvent.setup();
const submitSearch = async (placeholder: string, value: string) => {
diff --git a/web/src/pages/cluster/index.tsx b/web/src/pages/cluster/index.tsx
index accac39ac..021fcc67a 100644
--- a/web/src/pages/cluster/index.tsx
+++ b/web/src/pages/cluster/index.tsx
@@ -1197,6 +1197,42 @@ const ClusterPage = () => {
sorter: (a, b) => a.tpsOut - b.tpsOut,
render: (v: number) => v.toLocaleString(),
},
+ {
+ title: t('cluster.putMessagesToday'),
+ dataIndex: 'putMessagesToday',
+ key: 'putMessagesToday',
+ width: 90,
+ align: 'right',
+ sorter: (a, b) => (a.putMessagesToday ?? -1) - (b.putMessagesToday ??
-1),
+ render: (v?: number) => (v ?? 0).toLocaleString(),
+ },
+ {
+ title: t('cluster.putMessagesYesterday'),
+ dataIndex: 'putMessagesYesterday',
+ key: 'putMessagesYesterday',
+ width: 90,
+ align: 'right',
+ sorter: (a, b) => (a.putMessagesYesterday ?? -1) -
(b.putMessagesYesterday ?? -1),
+ render: (v?: number) => (v ?? 0).toLocaleString(),
+ },
+ {
+ title: t('cluster.getMessagesToday'),
+ dataIndex: 'getMessagesToday',
+ key: 'getMessagesToday',
+ width: 90,
+ align: 'right',
+ sorter: (a, b) => (a.getMessagesToday ?? -1) - (b.getMessagesToday ??
-1),
+ render: (v?: number) => (v ?? 0).toLocaleString(),
+ },
+ {
+ title: t('cluster.getMessagesYesterday'),
+ dataIndex: 'getMessagesYesterday',
+ key: 'getMessagesYesterday',
+ width: 90,
+ align: 'right',
+ sorter: (a, b) => (a.getMessagesYesterday ?? -1) -
(b.getMessagesYesterday ?? -1),
+ render: (v?: number) => (v ?? 0).toLocaleString(),
+ },
{
title: t('common.actions'),
key: 'action',