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 01f629f1 fix(consumer): keep aggregated lag unknown when any queue is 
unknown (#1652)
01f629f1 is described below

commit 01f629f102c2560160dbfb33ea0da5f6c4716bdc
Author: youngkermit8-coder <[email protected]>
AuthorDate: Tue Aug 11 21:19:49 2026 +0800

    fix(consumer): keep aggregated lag unknown when any queue is unknown (#1652)
    
    * fix: preserve unknown consumer lag summaries
    
    Signed-off-by: youngkermit8-coder <[email protected]>
    
    * fix: narrow unknown lag aggregation scope
    
    Signed-off-by: youngkermit8-coder <[email protected]>
    
    ---------
    
    Signed-off-by: youngkermit8-coder <[email protected]>
---
 .../provider/apache/RocketMQMetadataProvider.java  |  7 ++-
 .../apache/RocketMQMetadataProviderTest.java       | 55 +++++++++++++++++++++-
 2 files changed, 60 insertions(+), 2 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
index 9db74b72..6db0abb7 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProvider.java
@@ -309,7 +309,12 @@ public class RocketMQMetadataProvider implements 
MetadataProvider {
                     if (stats != null && stats.getOffsetTable() != null) {
                         for (Map.Entry<MessageQueue, OffsetWrapper> entry : 
stats.getOffsetTable().entrySet()) {
                             OffsetWrapper ow = entry.getValue();
-                            diffTotal += resolveDiff(ow.getBrokerOffset(), 
ow.getConsumerOffset());
+                            long queueDiff = resolveDiff(ow.getBrokerOffset(), 
ow.getConsumerOffset());
+                            if (queueDiff == ConsumerLagResolver.UNKNOWN) {
+                                diffTotal = ConsumerLagResolver.UNKNOWN;
+                                break;
+                            }
+                            diffTotal += queueDiff;
                         }
                         consumeTps = stats.getConsumeTps();
                     }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
index e4c204c6..4c5d07ba 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMetadataProviderTest.java
@@ -17,11 +17,14 @@
 package org.apache.rocketmq.studio.provider.apache;
 
 import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.common.message.MessageQueue;
+import org.apache.rocketmq.remoting.protocol.admin.ConsumeStats;
+import org.apache.rocketmq.remoting.protocol.admin.OffsetWrapper;
+import org.apache.rocketmq.remoting.protocol.body.GroupList;
 import org.apache.rocketmq.studio.common.domain.enums.ConsumeType;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
 import org.apache.rocketmq.tools.admin.MQAdminExt;
-import org.apache.rocketmq.remoting.protocol.body.GroupList;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
 import org.apache.rocketmq.studio.instance.group.SubscriptionEntryVO;
@@ -39,6 +42,8 @@ import org.mockito.junit.jupiter.MockitoExtension;
 
 import java.util.HashSet;
 import java.util.List;
+import java.util.LinkedHashMap;
+import java.util.Map;
 
 import static org.assertj.core.api.Assertions.assertThat;
 import static org.assertj.core.api.Assertions.assertThatThrownBy;
@@ -199,6 +204,30 @@ class RocketMQMetadataProviderTest {
         });
     }
 
+    @Test
+    void getTopicConsumersKeepsUnknownWhenAnyQueueLagIsUnknown() throws 
Exception {
+        DefaultMQAdminExt admin = 
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+        mockTopicConsumeStats(admin, offset(20, 10), offset(0, 1));
+
+        List<TopicConsumerVO> consumers = 
newLiveProvider(admin).getTopicConsumers(null, "TopicA");
+
+        assertThat(consumers).singleElement()
+                .extracting(TopicConsumerVO::getDiffTotal)
+                .isEqualTo(ConsumerLagResolver.UNKNOWN);
+    }
+
+    @Test
+    void getTopicConsumersStillSumsKnownQueueLags() throws Exception {
+        DefaultMQAdminExt admin = 
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+        mockTopicConsumeStats(admin, offset(20, 10), offset(7, 4));
+
+        List<TopicConsumerVO> consumers = 
newLiveProvider(admin).getTopicConsumers(null, "TopicA");
+
+        assertThat(consumers).singleElement()
+                .extracting(TopicConsumerVO::getDiffTotal)
+                .isEqualTo(13L);
+    }
+
     @Test
     void getGroupProgressSurfacesAdminFailure() throws Exception {
         DefaultMQAdminExt admin = 
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
@@ -231,4 +260,28 @@ class RocketMQMetadataProviderTest {
         return new RocketMQMetadataProvider(factory, liveProperties, 
topicMapper, groupMapper,
                 runtimeAdminClientResolver);
     }
+
+    private void mockTopicConsumeStats(DefaultMQAdminExt admin, 
OffsetWrapper... queueOffsets) throws Exception {
+        mockTopicGroup(admin);
+        Map<MessageQueue, OffsetWrapper> offsets = new LinkedHashMap<>();
+        for (int queueId = 0; queueId < queueOffsets.length; queueId++) {
+            offsets.put(new MessageQueue("TopicA", "broker-a", queueId), 
queueOffsets[queueId]);
+        }
+        ConsumeStats stats = new ConsumeStats();
+        stats.setOffsetTable(offsets);
+        when(admin.examineConsumeStats("group-a", "TopicA")).thenReturn(stats);
+    }
+
+    private void mockTopicGroup(DefaultMQAdminExt admin) throws Exception {
+        GroupList groupList = new GroupList();
+        groupList.setGroupList(new HashSet<>(List.of("group-a")));
+        when(admin.queryTopicConsumeByWho("TopicA")).thenReturn(groupList);
+    }
+
+    private OffsetWrapper offset(long brokerOffset, long consumerOffset) {
+        OffsetWrapper offset = new OffsetWrapper();
+        offset.setBrokerOffset(brokerOffset);
+        offset.setConsumerOffset(consumerOffset);
+        return offset;
+    }
 }

Reply via email to