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;

Reply via email to