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 a913ce7ae fix(group): preserve unknown lag in reset preview (#4540)
a913ce7ae is described below
commit a913ce7ae2c8dc7ad1cc6f473a7d1ce770b905a8
Author: zmuxuny <[email protected]>
AuthorDate: Mon Sep 21 17:30:16 2026 +0800
fix(group): preserve unknown lag in reset preview (#4540)
Reset-offset preview computed queue lag with Math.max(0, brokerOffset -
consumerOffset), which turned an undeterminable lag into a healthy-looking 0.
It now reuses ConsumerLagResolver.resolve so a negative difference yields the
established UNKNOWN (-1) sentinel, the aggregate totals propagate it instead of
summing a fabricated zero, and the preview warnings state that the affected
backlog totals are unavailable.
Fixes #4539
---
.../provider/apache/RocketMQAdminClientImpl.java | 24 +++++++---
.../apache/RocketMQAdminClientImplTest.java | 52 ++++++++++++++++++++++
2 files changed, 71 insertions(+), 5 deletions(-)
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 0b34bd39f..6257eb2ec 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
@@ -860,10 +860,8 @@ public class RocketMQAdminClientImpl implements
AdminClient {
"No consume offset data found for topic " + topic);
}
- long currentTotalLag =
queues.stream().mapToLong(ResetConsumerOffsetQueuePreviewVO::getCurrentLag).sum();
- long projectedTotalLag = queues.stream()
-
.mapToLong(ResetConsumerOffsetQueuePreviewVO::getProjectedLag)
- .sum();
+ long currentTotalLag = aggregateResetPreviewLag(queues, false);
+ long projectedTotalLag = aggregateResetPreviewLag(queues, true);
long totalOffsetDelta =
queues.stream().mapToLong(ResetConsumerOffsetQueuePreviewVO::getOffsetDelta).sum();
int rewindQueueCount = (int) queues.stream().filter(queue ->
queue.getOffsetDelta() < 0).count();
int fastForwardQueueCount = (int) queues.stream().filter(queue ->
queue.getOffsetDelta() > 0).count();
@@ -999,6 +997,10 @@ public class RocketMQAdminClientImpl implements
AdminClient {
if (rewindQueueCount > 0) {
warnings.add(rewindQueueCount + " queue(s) will move backward and
may replay consumed messages");
}
+ if (queues.stream().anyMatch(queue -> queue.getCurrentLag() ==
ConsumerLagResolver.UNKNOWN
+ || queue.getProjectedLag() == ConsumerLagResolver.UNKNOWN)) {
+ warnings.add("At least one queue has unavailable lag; affected
backlog totals are unavailable");
+ }
if (queues.stream().anyMatch(queue -> queue.getMinOffset() >= 0
&& queue.getTargetOffset() == queue.getMinOffset())) {
warnings.add("At least one queue will reset to the minimum
retained offset");
@@ -1010,8 +1012,20 @@ public class RocketMQAdminClientImpl implements
AdminClient {
return warnings;
}
+ private long
aggregateResetPreviewLag(List<ResetConsumerOffsetQueuePreviewVO> queues,
boolean projected) {
+ long total = 0L;
+ for (ResetConsumerOffsetQueuePreviewVO queue : queues) {
+ long lag = projected ? queue.getProjectedLag() :
queue.getCurrentLag();
+ if (lag == ConsumerLagResolver.UNKNOWN) {
+ return ConsumerLagResolver.UNKNOWN;
+ }
+ total += lag;
+ }
+ return total;
+ }
+
private long resolveLag(long brokerOffset, long consumerOffset) {
- return Math.max(0L, brokerOffset - consumerOffset);
+ return ConsumerLagResolver.resolve(brokerOffset - consumerOffset,
null);
}
private long clampOffset(long offset, long minOffset, long maxOffset) {
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 0a0fde32e..91a5735ad 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
@@ -408,6 +408,58 @@ class RocketMQAdminClientImplTest {
anyLong(), anyBoolean());
}
+ @Test
+ void previewResetOffsetShouldPreserveUnknownCurrentLagTest() throws
Exception {
+ long timestamp = 1784246400000L;
+ ConsumeStats stats = new ConsumeStats();
+ MessageQueue queue = new MessageQueue("orders", "broker-a", 0);
+ stats.getOffsetTable().put(queue, offsetWrapper(100L, 120L));
+ when(adminExt.examineConsumeStats("cg-orders")).thenReturn(stats);
+ ClusterInfo clusterInfo = clusterInfoWithMaster();
+
clusterInfo.getBrokerAddrTable().values().iterator().next().setBrokerName("broker-a");
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo);
+ when(adminExt.minOffset(queue)).thenReturn(0L);
+ when(adminExt.maxOffset(queue)).thenReturn(200L);
+ when(adminExt.searchOffset("10.0.0.1:10911", "orders", 0, timestamp,
3_000L)).thenReturn(80L);
+
+ ResetConsumerOffsetPreviewVO preview = adminClient.previewResetOffset(
+ null, "cg-orders", timestamp, "orders");
+
+
assertThat(preview.getCurrentTotalLag()).isEqualTo(ConsumerLagResolver.UNKNOWN);
+ assertThat(preview.getProjectedTotalLag()).isEqualTo(20L);
+ assertThat(preview.getQueues()).singleElement().satisfies(row -> {
+
assertThat(row.getCurrentLag()).isEqualTo(ConsumerLagResolver.UNKNOWN);
+ assertThat(row.getProjectedLag()).isEqualTo(20L);
+ });
+ assertThat(preview.getWarnings())
+ .contains("At least one queue has unavailable lag; affected
backlog totals are unavailable");
+ }
+
+ @Test
+ void previewResetOffsetShouldPreserveUnknownProjectedLagTest() throws
Exception {
+ long timestamp = 1784246400000L;
+ ConsumeStats stats = new ConsumeStats();
+ MessageQueue queue = new MessageQueue("orders", "broker-a", 0);
+ stats.getOffsetTable().put(queue, offsetWrapper(100L, 80L));
+ when(adminExt.examineConsumeStats("cg-orders")).thenReturn(stats);
+ ClusterInfo clusterInfo = clusterInfoWithMaster();
+
clusterInfo.getBrokerAddrTable().values().iterator().next().setBrokerName("broker-a");
+ when(adminExt.examineBrokerClusterInfo()).thenReturn(clusterInfo);
+ when(adminExt.minOffset(queue)).thenReturn(0L);
+ when(adminExt.maxOffset(queue)).thenReturn(200L);
+ when(adminExt.searchOffset("10.0.0.1:10911", "orders", 0, timestamp,
3_000L)).thenReturn(120L);
+
+ ResetConsumerOffsetPreviewVO preview = adminClient.previewResetOffset(
+ null, "cg-orders", timestamp, "orders");
+
+ assertThat(preview.getCurrentTotalLag()).isEqualTo(20L);
+
assertThat(preview.getProjectedTotalLag()).isEqualTo(ConsumerLagResolver.UNKNOWN);
+ assertThat(preview.getQueues()).singleElement().satisfies(row -> {
+ assertThat(row.getCurrentLag()).isEqualTo(20L);
+
assertThat(row.getProjectedLag()).isEqualTo(ConsumerLagResolver.UNKNOWN);
+ });
+ }
+
@Test
void previewResetOffsetShouldMarkFailedQueuesIncomplete() throws Exception
{
long timestamp = 1784246400000L;