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 8e9b8a85 fix(cluster): average broker TPS and use long counters (#1656)
8e9b8a85 is described below

commit 8e9b8a85e0c73d0d3d5fcba794bfcc54a36752ff
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 11 20:45:34 2026 +0800

    fix(cluster): average broker TPS and use long counters (#1656)
---
 .../rocketmq/studio/cluster/broker/BrokerVO.java   |  4 +--
 .../provider/apache/RocketMQClusterProvider.java   | 11 +++++---
 .../apache/RocketMQClusterProviderTest.java        | 30 +++++++++++++++++++---
 3 files changed, 35 insertions(+), 10 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 815838b8..8992122b 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
@@ -32,8 +32,8 @@ public class BrokerVO {
     private String version;
     private BrokerStatus status;
     private double diskUsage;
-    private int tpsIn;
-    private int tpsOut;
+    private long tpsIn;
+    private long tpsOut;
     @Builder.Default
     private boolean runtimeStatsAvailable = true;
 }
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 413538a1..479bb0ec 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
@@ -251,13 +251,16 @@ public class RocketMQClusterProvider implements 
ClusterProvider {
 
     /**
      * TPS properties are space-separated values representing 10s/1min/10min 
averages.
-     * Parse the first value (10s average).
+     * Use the 1-minute average to match the Dashboard overview.
      */
-    private int parseTpsValue(String tpsStr) {
+    private long parseTpsValue(String tpsStr) {
         try {
             String[] parts = tpsStr.trim().split("\\s+");
-            if (parts.length > 0) {
-                return (int) Double.parseDouble(parts[0]);
+            if (parts.length >= 2) {
+                return (long) Double.parseDouble(parts[1]);
+            }
+            if (parts.length == 1) {
+                return (long) Double.parseDouble(parts[0]);
             }
         } catch (NumberFormatException ignored) {
             // fall through
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 7cbf83b2..83ad8f74 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
@@ -52,8 +52,8 @@ class RocketMQClusterProviderTest {
 
         assertThat(clusters).hasSize(1);
         assertThat(clusters.get(0).getBrokers()).hasSize(1);
-        
assertThat(clusters.get(0).getBrokers().get(0).getTpsIn()).isEqualTo(12);
-        
assertThat(clusters.get(0).getBrokers().get(0).getTpsOut()).isEqualTo(34);
+        
assertThat(clusters.get(0).getBrokers().get(0).getTpsIn()).isEqualTo(10);
+        
assertThat(clusters.get(0).getBrokers().get(0).getTpsOut()).isEqualTo(30);
         
assertThat(clusters.get(0).getBrokers().get(0).isRuntimeStatsAvailable()).isTrue();
         
assertThat(clusters.get(0).getStatus()).isEqualTo(ClusterStatus.healthy);
     }
@@ -108,6 +108,24 @@ class RocketMQClusterProviderTest {
                 .isEqualTo(false);
     }
 
+    @Test
+    void discoverClustersShouldKeepOneMinuteTpsAboveIntegerRange() 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("1.0 3000000000.0 3.0", "4.0 
4000000000.0 6.0"));
+
+        List<ClusterVO> clusters = provider.discoverClusters();
+
+        assertThat(clusters).singleElement().satisfies(cluster -> {
+            assertThat(cluster.getBrokers()).singleElement().satisfies(broker 
-> {
+                assertThat(broker.getTpsIn()).isEqualTo(3_000_000_000L);
+                assertThat(broker.getTpsOut()).isEqualTo(4_000_000_000L);
+            });
+        });
+    }
+
     private RocketMQClusterProvider newProvider(DefaultMQAdminExt adminExt) {
         MqAdminExtFactory adminFactory = mock(MqAdminExtFactory.class);
         when(adminFactory.execute(anyString(), any(), 
any())).thenAnswer(invocation ->
@@ -146,10 +164,14 @@ class RocketMQClusterProviderTest {
     }
 
     private KVTable runtimeStats() {
+        return runtimeStats("  12.7   10.0  9.0", "\\t34.2  30.0  29.0");
+    }
+
+    private KVTable runtimeStats(String putTps, String getTransferredTps) {
         KVTable table = new KVTable();
         HashMap<String, String> stats = new HashMap<>();
-        stats.put("putTps", "  12.7   10.0  9.0");
-        stats.put("getTransferredTps", "\t34.2  30.0  29.0");
+        stats.put("putTps", putTps);
+        stats.put("getTransferredTps", getTransferredTps);
         table.setTable(stats);
         return table;
     }

Reply via email to