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);