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 83aba4500 fix(consumer): report unknown lag instead of summing
sentinel queues away (#3988)
83aba4500 is described below
commit 83aba4500645e1589e892f813af33f280adf20ab
Author: Zhao Jianing <[email protected]>
AuthorDate: Mon Sep 7 19:25:33 2026 +0800
fix(consumer): report unknown lag instead of summing sentinel queues away
(#3988)
enrichGroupLiveStats and fillConsumeStats accumulated queue diffs with
an 'if (diff > 0)' clamp, so a queue whose brokerOffset-consumerOffset
diff is negative (the -1 unknown sentinel RocketMQ 5.0 gRPC consumers
report) silently contributed zero to totalLag. A group that also had
healthy queues showed a partial sum presented as the total, and a group
whose queues were all unknown showed a fabricated zero lag with
consumeStatsAvailable=true.
ConsumerLagResolver and the -1 UNKNOWN sentinel were introduced exactly
to keep this state visible (getTopicConsumersPage and getGroupProgress
already resolve every diff through it), and the web UI already renders
the sentinel as an explicit unavailable state (web/src/utils/
consumerLag.ts isLagAvailable/formatLag). The two group-list/group-
detail aggregation sites were left on the legacy clamping behavior.
Route both loops through ConsumerLagResolver: any unknown queue marks
the totalLag unknown (matching getTopicConsumersPage's diffTotal
semantics) while the delay-seconds scan still covers every queue.
Signed-off-by: zjncs <[email protected]>
---
.../ApacheRocketMqBusinessMetricsCollector.java | 10 ++++++-
.../studio/instance/topic/MetadataService.java | 7 ++++-
.../ops/ai/tool/ConsumerGroupListToolHandler.java | 3 ++-
.../provider/apache/RocketMQAdminClientImpl.java | 18 +++++++++----
.../provider/apache/RocketMQMetadataProvider.java | 13 ++++++---
...ApacheRocketMqBusinessMetricsCollectorTest.java | 22 +++++++++++++++
.../apache/RocketMQAdminClientImplTest.java | 19 +++++++++++++
.../apache/RocketMQMetadataProviderTest.java | 31 ++++++++++++++++++++++
8 files changed, 111 insertions(+), 12 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/collectors/ApacheRocketMqBusinessMetricsCollector.java
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/collectors/ApacheRocketMqBusinessMetricsCollector.java
index e4cfe9af9..f7ad0e365 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/collectors/ApacheRocketMqBusinessMetricsCollector.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/cluster/metrics/collectors/ApacheRocketMqBusinessMetricsCollector.java
@@ -27,6 +27,7 @@ import
org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
import org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
+import org.apache.rocketmq.studio.provider.apache.ConsumerLagResolver;
import org.apache.rocketmq.studio.ops.alert.AlertDomain;
import org.springframework.stereotype.Component;
@@ -84,7 +85,14 @@ public class ApacheRocketMqBusinessMetricsCollector
implements BusinessMetricsCo
"CONSUMER_STATS_UNAVAILABLE"));
continue;
}
- samples.add(totalLagSample(instance, group, collectedAt));
+ if (group.getTotalLag() == ConsumerLagResolver.UNKNOWN) {
+ // totalLag now carries the -1 unknown sentinel; do not
clamp it into a fabricated
+ // zero-lag AVAILABLE sample that would feed
consumer.lag.total alerts a fake 0.
+ samples.add(unavailable(CONSUMER_LAG_TOTAL, instance,
+ Map.of("consumerGroup", group.getName()),
collectedAt, "CONSUMER_LAG_UNKNOWN"));
+ } else {
+ samples.add(totalLagSample(instance, group, collectedAt));
+ }
if (group.isConsumptionTimestampAvailable()) {
samples.add(delaySample(instance, group, collectedAt));
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
index 04d347c8b..b2a552c91 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/topic/MetadataService.java
@@ -21,6 +21,7 @@ import
org.apache.rocketmq.studio.audit.OperationAuditConstants.ResourceType;
import org.apache.rocketmq.studio.audit.OperationAuditConstants.Result;
import org.apache.rocketmq.studio.audit.OperationAuditService;
import org.apache.rocketmq.studio.provider.apache.AdminClient;
+import org.apache.rocketmq.studio.provider.apache.ConsumerLagResolver;
import org.apache.rocketmq.studio.provider.apache.MetadataProvider;
import org.apache.rocketmq.studio.common.domain.PageResult;
import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
@@ -454,7 +455,7 @@ public class MetadataService {
for (ConsumerGroupVO group : groups) {
CsvUtil.appendRow(csv, group.getName(), group.getNamespace(),
group.getClusterId(),
toText(group.getSubscriptionMode()),
toText(group.getConsumeType()),
- group.getOnlineInstances(), group.getTotalLag(),
group.getDelaySeconds(),
+ group.getOnlineInstances(), lagText(group.getTotalLag()),
group.getDelaySeconds(),
group.getSubscriptionDataType(),
group.getDeliveryOrderType(), group.getRetryMaxTimes(),
String.join(";", group.getSubscribedTopics() == null ?
List.of() : group.getSubscribedTopics()),
group.getGmtCreate(), group.getGmtModified());
@@ -462,6 +463,10 @@ public class MetadataService {
return csv.toString();
}
+ private static String lagText(long totalLag) {
+ return totalLag == ConsumerLagResolver.UNKNOWN ? "unknown" :
String.valueOf(totalLag);
+ }
+
private String toText(Object value) {
if (value == null) {
return "";
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/ConsumerGroupListToolHandler.java
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/ConsumerGroupListToolHandler.java
index 0caf1e785..fcaf58381 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/ConsumerGroupListToolHandler.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/ops/ai/tool/ConsumerGroupListToolHandler.java
@@ -18,6 +18,7 @@ package org.apache.rocketmq.studio.ops.ai.tool;
import org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.instance.topic.MetadataService;
+import org.apache.rocketmq.studio.provider.apache.ConsumerLagResolver;
import org.apache.rocketmq.studio.common.domain.PageResult;
import lombok.RequiredArgsConstructor;
import org.springframework.stereotype.Component;
@@ -60,7 +61,7 @@ public class ConsumerGroupListToolHandler implements
ToolHandler {
result.put("consumeType", requiredEnumName(
group.getConsumeType(), "consumeType", group.getName()));
result.put("onlineInstances", group.getOnlineInstances());
- result.put("totalLag", group.getTotalLag());
+ result.put("totalLag", group.getTotalLag() ==
ConsumerLagResolver.UNKNOWN ? null : group.getTotalLag());
result.put("subscribedTopics", copyList(group.getSubscribedTopics()));
result.put("retryMaxTimes", group.getRetryMaxTimes());
return result;
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
index 1baa1ee3a..63bcad7cb 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImpl.java
@@ -158,7 +158,9 @@ public class RocketMQAdminClientImpl implements AdminClient
{
* Fills totalLag and delaySeconds from the broker consume stats.
Proxy-connected groups
* still maintain broker-side offset tables (the proxy forwards offset
updates), so this
* works even when the connection lookup reports the group offline; groups
without any
- * offset table (e.g. pure POP) simply keep the zero defaults.
+ * offset table (e.g. pure POP) simply keep the zero defaults. A queue
whose offsets
+ * resolve to the unknown sentinel marks totalLag unknown instead of
summing it away as
+ * zero lag.
*
* <p>delaySeconds is derived from the newest consumed-message timestamp
(the consumption
* frontier). Using the oldest timestamp is misleading for POP groups,
where untouched
@@ -175,18 +177,24 @@ public class RocketMQAdminClientImpl implements
AdminClient {
return;
}
long totalLag = 0;
+ boolean lagUnknown = false;
long newestConsumedTimestamp = 0;
for (OffsetWrapper wrapper : stats.getOffsetTable().values()) {
- long diff = wrapper.getBrokerOffset() -
wrapper.getConsumerOffset();
- if (diff > 0) {
- totalLag += diff;
+ long queueDiff = ConsumerLagResolver.resolve(
+ wrapper.getBrokerOffset() -
wrapper.getConsumerOffset(), null);
+ if (queueDiff == ConsumerLagResolver.UNKNOWN) {
+ // a queue with the -1 sentinel (5.0 gRPC consumers) must
not be summed
+ // away as zero lag; report the whole total as unknown
instead
+ lagUnknown = true;
+ } else {
+ totalLag += queueDiff;
}
long lastTimestamp = wrapper.getLastTimestamp();
if (lastTimestamp > newestConsumedTimestamp) {
newestConsumedTimestamp = lastTimestamp;
}
}
- vo.setTotalLag(totalLag);
+ vo.setTotalLag(lagUnknown ? ConsumerLagResolver.UNKNOWN :
totalLag);
if (newestConsumedTimestamp > 0) {
long delaySeconds = (System.currentTimeMillis() -
newestConsumedTimestamp) / 1000;
vo.setDelaySeconds((int) Math.max(delaySeconds, 0));
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 8070d5655..6d77fc91d 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
@@ -330,18 +330,23 @@ public class RocketMQMetadataProvider implements
MetadataProvider {
return;
}
long totalLag = 0;
+ boolean lagUnknown = false;
long newestConsumedTimestamp = 0;
for (OffsetWrapper wrapper : stats.getOffsetTable().values()) {
- long diff = wrapper.getBrokerOffset() -
wrapper.getConsumerOffset();
- if (diff > 0) {
- totalLag += diff;
+ long queueDiff = resolveDiff(wrapper.getBrokerOffset(),
wrapper.getConsumerOffset());
+ if (queueDiff == ConsumerLagResolver.UNKNOWN) {
+ // a queue with the -1 sentinel (5.0 gRPC consumers) must
not be summed
+ // away as zero lag; report the whole total as unknown
instead
+ lagUnknown = true;
+ } else {
+ totalLag += queueDiff;
}
long lastTimestamp = wrapper.getLastTimestamp();
if (lastTimestamp > newestConsumedTimestamp) {
newestConsumedTimestamp = lastTimestamp;
}
}
- vo.setTotalLag(totalLag);
+ vo.setTotalLag(lagUnknown ? ConsumerLagResolver.UNKNOWN :
totalLag);
if (newestConsumedTimestamp > 0) {
long delaySeconds = (System.currentTimeMillis() -
newestConsumedTimestamp) / 1000;
vo.setDelaySeconds((int) Math.max(delaySeconds, 0));
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/collectors/ApacheRocketMqBusinessMetricsCollectorTest.java
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/collectors/ApacheRocketMqBusinessMetricsCollectorTest.java
index ffb5d5cba..ad2b4c813 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/collectors/ApacheRocketMqBusinessMetricsCollectorTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/cluster/metrics/collectors/ApacheRocketMqBusinessMetricsCollectorTest.java
@@ -24,6 +24,7 @@ import
org.apache.rocketmq.studio.instance.group.ConsumerGroupVO;
import org.apache.rocketmq.studio.instance.group.QueueProgressVO;
import org.apache.rocketmq.studio.provider.InstanceProvider;
import org.apache.rocketmq.studio.provider.InstanceProviderRegistry;
+import org.apache.rocketmq.studio.provider.apache.ConsumerLagResolver;
import org.junit.jupiter.api.Test;
import java.util.List;
@@ -128,6 +129,27 @@ class ApacheRocketMqBusinessMetricsCollectorTest {
});
}
+ @Test
+ void reportsUnavailableTotalLagWhenSentinelUnknownTest() {
+ InstanceProviderRegistry registry =
mock(InstanceProviderRegistry.class);
+ InstanceProvider provider = mock(InstanceProvider.class);
+ ConsumerGroupVO orders = group("orders", "cluster-a",
ConsumerLagResolver.UNKNOWN);
+ when(registry.byInstanceId("local")).thenReturn(Optional.of(provider));
+ when(provider.listConsumerGroups("local",
null)).thenReturn(List.of(orders));
+ when(provider.getGroupProgress("local",
"orders")).thenReturn(List.of(QueueProgressVO.builder()
+ .topic("orders-topic").diffTotal(17).build()));
+
+ List<MetricSample> samples = new
ApacheRocketMqBusinessMetricsCollector(registry).collect(apacheInstance());
+
+ assertThat(samples).filteredOn(sample -> sample.metricKey().equals(
+ ApacheRocketMqBusinessMetricsCollector.CONSUMER_LAG_TOTAL))
+ .singleElement().satisfies(sample -> {
+
assertThat(sample.availability()).isEqualTo(MetricAvailability.UNAVAILABLE);
+ assertThat(sample.value()).isNull();
+
assertThat(sample.unavailableReason()).isEqualTo("CONSUMER_LAG_UNKNOWN");
+ });
+ }
+
private static ConsumerGroupVO group(String name, String clusterId, long
lag) {
ConsumerGroupVO group = new ConsumerGroupVO();
group.setName(name);
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
index 4e0993def..741149d9f 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQAdminClientImplTest.java
@@ -235,6 +235,25 @@ class RocketMQAdminClientImplTest {
assertThat(group.getDelaySeconds()).isBetween(4, 30);
}
+ @Test
+ void
getConsumerGroupReportsUnknownTotalLagWhenAnyQueueOffsetIsUnknownTest() throws
Exception {
+ org.apache.rocketmq.remoting.protocol.body.ConsumerConnection
connection =
+ new
org.apache.rocketmq.remoting.protocol.body.ConsumerConnection();
+ connection.setConnectionSet(new java.util.HashSet<>());
+
when(adminExt.examineConsumerConnectionInfo("orders")).thenReturn(connection);
+
+ org.apache.rocketmq.remoting.protocol.admin.ConsumeStats stats =
+ new org.apache.rocketmq.remoting.protocol.admin.ConsumeStats();
+ stats.getOffsetTable().put(new MessageQueue("orders-topic",
"broker-a", 0), offsetWrapper(100L, 60L));
+ stats.getOffsetTable().put(new MessageQueue("orders-topic",
"broker-b", 1), offsetWrapper(0L, 1L));
+ when(adminExt.examineConsumeStats("orders")).thenReturn(stats);
+
+ ConsumerGroupVO group = adminClient.getConsumerGroup(null, "orders");
+
+ assertThat(group.isConsumeStatsAvailable()).isTrue();
+ assertThat(group.getTotalLag()).isEqualTo(ConsumerLagResolver.UNKNOWN);
+ }
+
@Test
void getConsumerGroupFillsOnlineInstanceListFromConnectionsTest() throws
Exception {
org.apache.rocketmq.remoting.protocol.body.ConsumerConnection
connection =
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 3fdd47f6b..4bb63e688 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
@@ -726,6 +726,37 @@ class RocketMQMetadataProviderTest {
assertThat(groups.get(0).getDelaySeconds()).isBetween(4, 30);
}
+ @Test
+ void
listConsumerGroupsShouldReportUnknownTotalLagWhenAnyQueueOffsetIsUnknownTest()
throws Exception {
+ RmqGroup entity = new RmqGroup();
+ entity.setName("cg-unknown");
+ entity.setInstanceId("instance-a");
+ when(groupMapper.selectList(any())).thenReturn(List.of(entity));
+
+ DefaultMQAdminExt admin =
org.mockito.Mockito.mock(DefaultMQAdminExt.class);
+ org.apache.rocketmq.remoting.protocol.body.ConsumerConnection
connection =
+ new
org.apache.rocketmq.remoting.protocol.body.ConsumerConnection();
+ connection.setConnectionSet(new java.util.HashSet<>());
+
when(admin.examineConsumerConnectionInfo("cg-unknown")).thenReturn(connection);
+
+ org.apache.rocketmq.remoting.protocol.admin.ConsumeStats stats =
+ new org.apache.rocketmq.remoting.protocol.admin.ConsumeStats();
+ stats.getOffsetTable().put(new MessageQueue("studio-normal",
"broker-a", 0), offset(100, 60));
+ stats.getOffsetTable().put(new MessageQueue("studio-normal",
"broker-b", 1), offset(0, 1));
+ when(admin.examineConsumeStats("cg-unknown")).thenReturn(stats);
+
when(runtimeAdminClientResolver.execute(org.mockito.ArgumentMatchers.eq("instance-a"),
any()))
+ .thenAnswer(invocation ->
+
invocation.<MqAdminExtFactory.AdminAction<Object>>getArgument(1).apply(admin));
+
+ RocketMQMetadataProvider provider = newLiveProvider(admin);
+
+ List<ConsumerGroupVO> groups =
provider.listConsumerGroups("instance-a", null, null);
+
+ assertThat(groups).hasSize(1);
+ assertThat(groups.get(0).isConsumeStatsAvailable()).isTrue();
+
assertThat(groups.get(0).getTotalLag()).isEqualTo(ConsumerLagResolver.UNKNOWN);
+ }
+
private RocketMQMetadataProvider newLiveProvider(MQAdminExt admin) throws
Exception {
MqAdminExtFactory factory = mock(MqAdminExtFactory.class);
RocketMQProperties liveProperties = new RocketMQProperties();