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',

Reply via email to