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 cb2ca635b fix(metrics): preserve unknown Apache queue lag (#4517)
cb2ca635b is described below

commit cb2ca635ba0de410f6e47739c25d298945ce684b
Author: zmuxuny <[email protected]>
AuthorDate: Mon Sep 21 17:27:52 2026 +0800

    fix(metrics): preserve unknown Apache queue lag (#4517)
    
    ApacheRocketMqBusinessMetricsCollector no longer clamps 
ConsumerLagResolver.UNKNOWN (-1) queue rows to zero when aggregating 
consumer.lag.max_queue and topic.backlog.total: an unknown queue now yields an 
UNAVAILABLE sample with reason CONSUMER_LAG_UNKNOWN instead of a fabricated 
zero lag. UNAVAILABLE samples also carry the group's clusterId so 
cluster-scoped alert rules no longer mis-resolve them.
    
    Fixes #4516
---
 .../ApacheRocketMqBusinessMetricsCollector.java    | 56 +++++++++++++++-------
 ...ApacheRocketMqBusinessMetricsCollectorTest.java | 37 ++++++++++++++
 2 files changed, 75 insertions(+), 18 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 f7ad0e365..c0e10eaa5 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
@@ -75,20 +75,20 @@ public class ApacheRocketMqBusinessMetricsCollector 
implements BusinessMetricsCo
                 }
                 if (!group.isConsumeStatsAvailable()) {
                     Map<String, String> labels = Map.of("consumerGroup", 
group.getName());
-                    samples.add(unavailable(CONSUMER_LAG_TOTAL, instance, 
labels, collectedAt,
+                    samples.add(unavailable(CONSUMER_LAG_TOTAL, instance, 
group.getClusterId(), labels, collectedAt,
                             "CONSUMER_STATS_UNAVAILABLE"));
-                    samples.add(unavailable(CONSUMER_LAG_MAX_QUEUE, instance, 
labels, collectedAt,
+                    samples.add(unavailable(CONSUMER_LAG_MAX_QUEUE, instance, 
group.getClusterId(), labels, collectedAt,
                             "CONSUMER_STATS_UNAVAILABLE"));
-                    samples.add(unavailable(CONSUMER_DELAY_SECONDS, instance, 
labels, collectedAt,
+                    samples.add(unavailable(CONSUMER_DELAY_SECONDS, instance, 
group.getClusterId(), labels, collectedAt,
                             "CONSUMER_STATS_UNAVAILABLE"));
-                    samples.add(unavailable(TOPIC_BACKLOG_TOTAL, instance, 
labels, collectedAt,
+                    samples.add(unavailable(TOPIC_BACKLOG_TOTAL, instance, 
group.getClusterId(), labels, collectedAt,
                             "CONSUMER_STATS_UNAVAILABLE"));
                     continue;
                 }
                 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,
+                    samples.add(unavailable(CONSUMER_LAG_TOTAL, instance, 
group.getClusterId(),
                             Map.of("consumerGroup", group.getName()), 
collectedAt, "CONSUMER_LAG_UNKNOWN"));
                 } else {
                     samples.add(totalLagSample(instance, group, collectedAt));
@@ -129,32 +129,52 @@ public class ApacheRocketMqBusinessMetricsCollector 
implements BusinessMetricsCo
         Map<String, String> labels = Map.of("consumerGroup", group.getName());
         try {
             List<QueueProgressVO> progress = 
provider.getGroupProgress(instance.getName(), group.getName());
-            long maxLag = 
progress.stream().mapToLong(QueueProgressVO::getDiffTotal)
-                    .map(value -> Math.max(0, value)).max().orElse(0);
+            boolean lagUnknown = progress.stream()
+                    .anyMatch(row -> row.getDiffTotal() == 
ConsumerLagResolver.UNKNOWN);
             List<MetricSample> samples = new ArrayList<>();
-            samples.add(new MetricSample(CONSUMER_LAG_MAX_QUEUE, 
AlertDomain.BUSINESS, instance.getName(),
-                    group.getClusterId(), labels, (double) maxLag, 
MetricAvailability.AVAILABLE, collectedAt));
+            if (lagUnknown) {
+                samples.add(unavailable(CONSUMER_LAG_MAX_QUEUE, instance, 
group.getClusterId(), labels, collectedAt,
+                        "CONSUMER_LAG_UNKNOWN"));
+            } else {
+                long maxLag = 
progress.stream().mapToLong(QueueProgressVO::getDiffTotal).max().orElse(0);
+                samples.add(new MetricSample(CONSUMER_LAG_MAX_QUEUE, 
AlertDomain.BUSINESS, instance.getName(),
+                        group.getClusterId(), labels, (double) maxLag, 
MetricAvailability.AVAILABLE, collectedAt));
+            }
             progress.stream().filter(row -> row.getTopic() != null && 
!row.getTopic().isBlank())
-                    
.collect(java.util.stream.Collectors.groupingBy(QueueProgressVO::getTopic,
-                            java.util.stream.Collectors.summingLong(
-                                    row -> Math.max(0, row.getDiffTotal()))))
-                    .forEach((topic, lag) -> samples.add(new 
MetricSample(TOPIC_BACKLOG_TOTAL, AlertDomain.BUSINESS,
-                            instance.getName(), group.getClusterId(), 
Map.of("consumerGroup", group.getName(),
-                            "topic", topic), (double) lag, 
MetricAvailability.AVAILABLE, collectedAt)));
+                    
.collect(java.util.stream.Collectors.groupingBy(QueueProgressVO::getTopic))
+                    .forEach((topic, rows) -> addTopicBacklogSample(samples, 
instance, group, topic, rows, collectedAt));
             return samples;
         } catch (RuntimeException error) {
             log.warn("Failed to collect queue lag for group {} on instance {}: 
{}", group.getName(),
                     instance.getName(), error.getMessage());
-            return List.of(unavailable(CONSUMER_LAG_MAX_QUEUE, instance, 
labels, collectedAt,
+            return List.of(unavailable(CONSUMER_LAG_MAX_QUEUE, instance, 
group.getClusterId(), labels, collectedAt,
                     "CONSUMER_PROGRESS_UNAVAILABLE"),
-                    unavailable(TOPIC_BACKLOG_TOTAL, instance, labels, 
collectedAt,
+                    unavailable(TOPIC_BACKLOG_TOTAL, instance, 
group.getClusterId(), labels, collectedAt,
                             "CONSUMER_PROGRESS_UNAVAILABLE"));
         }
     }
 
+    private static void addTopicBacklogSample(List<MetricSample> samples, 
InstanceVO instance,
+            ConsumerGroupVO group, String topic, List<QueueProgressVO> rows, 
Instant collectedAt) {
+        Map<String, String> topicLabels = Map.of("consumerGroup", 
group.getName(), "topic", topic);
+        if (rows.stream().anyMatch(row -> row.getDiffTotal() == 
ConsumerLagResolver.UNKNOWN)) {
+            samples.add(unavailable(TOPIC_BACKLOG_TOTAL, instance, 
group.getClusterId(), topicLabels, collectedAt,
+                    "CONSUMER_LAG_UNKNOWN"));
+            return;
+        }
+        long lag = 
rows.stream().mapToLong(QueueProgressVO::getDiffTotal).sum();
+        samples.add(new MetricSample(TOPIC_BACKLOG_TOTAL, 
AlertDomain.BUSINESS, instance.getName(),
+                group.getClusterId(), topicLabels, (double) lag, 
MetricAvailability.AVAILABLE, collectedAt));
+    }
+
     private static MetricSample unavailable(String metric, InstanceVO 
instance, Map<String, String> labels,
             Instant collectedAt, String reason) {
-        return new MetricSample(metric, AlertDomain.BUSINESS, 
instance.getName(), null, labels, null,
+        return unavailable(metric, instance, null, labels, collectedAt, 
reason);
+    }
+
+    private static MetricSample unavailable(String metric, InstanceVO 
instance, String clusterId,
+            Map<String, String> labels, Instant collectedAt, String reason) {
+        return new MetricSample(metric, AlertDomain.BUSINESS, 
instance.getName(), clusterId, labels, null,
                 MetricAvailability.UNAVAILABLE, collectedAt, reason);
     }
 }
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 ad2b4c813..f84483319 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
@@ -150,6 +150,43 @@ class ApacheRocketMqBusinessMetricsCollectorTest {
                 });
     }
 
+    @Test
+    void preservesUnknownQueueLagInMetricsAvailabilityTest() {
+        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("known-topic").diffTotal(12).build(),
+                
QueueProgressVO.builder().topic("unknown-topic").diffTotal(ConsumerLagResolver.UNKNOWN).build()));
+
+        List<MetricSample> samples = new 
ApacheRocketMqBusinessMetricsCollector(registry).collect(apacheInstance());
+
+        assertThat(samples).filteredOn(sample -> sample.metricKey().equals(
+                ApacheRocketMqBusinessMetricsCollector.CONSUMER_LAG_MAX_QUEUE))
+                .singleElement().satisfies(sample -> {
+                    
assertThat(sample.availability()).isEqualTo(MetricAvailability.UNAVAILABLE);
+                    assertThat(sample.value()).isNull();
+                    assertThat(sample.clusterId()).isEqualTo("cluster-a");
+                });
+        assertThat(samples).filteredOn(sample -> sample.metricKey().equals(
+                        
ApacheRocketMqBusinessMetricsCollector.TOPIC_BACKLOG_TOTAL)
+                        && "known-topic".equals(sample.labels().get("topic")))
+                .singleElement().satisfies(sample -> {
+                    
assertThat(sample.availability()).isEqualTo(MetricAvailability.AVAILABLE);
+                    assertThat(sample.value()).isEqualTo(12D);
+                });
+        assertThat(samples).filteredOn(sample -> sample.metricKey().equals(
+                        
ApacheRocketMqBusinessMetricsCollector.TOPIC_BACKLOG_TOTAL)
+                        && 
"unknown-topic".equals(sample.labels().get("topic")))
+                .singleElement().satisfies(sample -> {
+                    
assertThat(sample.availability()).isEqualTo(MetricAvailability.UNAVAILABLE);
+                    assertThat(sample.value()).isNull();
+                    assertThat(sample.clusterId()).isEqualTo("cluster-a");
+                });
+    }
+
     private static ConsumerGroupVO group(String name, String clusterId, long 
lag) {
         ConsumerGroupVO group = new ConsumerGroupVO();
         group.setName(name);

Reply via email to