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;
+ }
}