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 dc0e049a fix(dlq): distinguish scan failures from partial resend 
failures (#1653)
dc0e049a is described below

commit dc0e049a77d94c4be8f3f302dbe5612a4feaa508
Author: aias00 <[email protected]>
AuthorDate: Tue Aug 11 20:44:56 2026 +0800

    fix(dlq): distinguish scan failures from partial resend failures (#1653)
---
 .../studio/instance/dlq/DLQResendResultVO.java     |  2 +
 .../provider/apache/RocketMQDLQProvider.java       | 53 ++++++++++++++++----
 .../provider/apache/RocketMQDLQProviderTest.java   | 57 ++++++++++++++++++++++
 web/src/api/message.ts                             |  4 +-
 web/src/pages/instance/__tests__/DLQPage.test.tsx  | 24 +++++++++
 web/src/pages/instance/dlq.tsx                     |  6 ++-
 6 files changed, 135 insertions(+), 11 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendResultVO.java
 
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendResultVO.java
index 36cf13b0..1a7026e9 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendResultVO.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendResultVO.java
@@ -26,4 +26,6 @@ public class DLQResendResultVO {
     int resent;
     int failed;
     String outcome;
+    boolean scanIncomplete;
+    int failedQueueCount;
 }
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
index 2338d4c7..87ad7901 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProvider.java
@@ -31,6 +31,7 @@ import 
org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
 import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
 import org.apache.rocketmq.remoting.protocol.body.TopicList;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.instance.dlq.DLQGroupVO;
 import org.apache.rocketmq.studio.instance.dlq.DLQProvider;
 import org.apache.rocketmq.studio.instance.dlq.DLQResendResultVO;
@@ -142,7 +143,18 @@ public class RocketMQDLQProvider implements DLQProvider {
         long end = endTime != null ? endTime : System.currentTimeMillis();
         long begin = startTime != null ? startTime : end - ONE_HOUR_MILLIS;
 
-        List<MessageExt> deadLetters = collectDeadLetters(endpoint, dlqTopic, 
begin, end);
+        DeadLetterScanResult scanResult;
+        try {
+            scanResult = collectDeadLetters(endpoint, dlqTopic, begin, end);
+        } catch (BusinessException e) {
+            String detail = String.format("instanceId=%s, group=%s, 
dlqTopic=%s, targetTopic=%s, "
+                            + "matched=0, resent=0, failed=0, 
scanIncomplete=true, scanFailedQueues=all",
+                    instanceId, groupName, dlqTopic,
+                    StringUtils.hasText(targetTopic) ? targetTopic : 
"<original>");
+            recordAudit(groupName, detail, "FAILED");
+            throw e;
+        }
+        List<MessageExt> deadLetters = scanResult.messages();
         int resent = 0;
         int failed = 0;
         if (!deadLetters.isEmpty()) {
@@ -164,27 +176,32 @@ public class RocketMQDLQProvider implements DLQProvider {
             }
         }
 
-        String detail = String.format("instanceId=%s, group=%s, dlqTopic=%s, 
targetTopic=%s, matched=%d, resent=%d, failed=%d",
+        String outcome = classifyOutcome(deadLetters.size(), resent, failed, 
scanResult.scanIncomplete());
+        String detail = String.format("instanceId=%s, group=%s, dlqTopic=%s, 
targetTopic=%s, matched=%d, resent=%d, "
+                        + "failed=%d, scanIncomplete=%s, scanFailedQueues=%d",
                 instanceId, groupName, dlqTopic, 
StringUtils.hasText(targetTopic) ? targetTopic : "<original>",
-                deadLetters.size(), resent, failed);
-        recordAudit(groupName, detail, classifyOutcome(deadLetters.size(), 
resent, failed));
+                deadLetters.size(), resent, failed, 
scanResult.scanIncomplete(), scanResult.failedQueueCount());
+        recordAudit(groupName, detail, outcome);
         log.info("DLQ resend completed: {}", detail);
         return DLQResendResultVO.builder()
                 .matched(deadLetters.size())
                 .resent(resent)
                 .failed(failed)
-                .outcome(classifyOutcome(deadLetters.size(), resent, failed))
+                .outcome(outcome)
+                .scanIncomplete(scanResult.scanIncomplete())
+                .failedQueueCount(scanResult.failedQueueCount())
                 .build();
     }
 
-    private List<MessageExt> collectDeadLetters(String endpoint, String 
dlqTopic, long begin, long end) {
+    private DeadLetterScanResult collectDeadLetters(String endpoint, String 
dlqTopic, long begin, long end) {
         DefaultMQPullConsumer consumer = newPullConsumer(endpoint);
         List<MessageExt> result = new ArrayList<>();
+        int failedQueueCount = 0;
         try {
             consumer.start();
             Set<MessageQueue> queues = 
consumer.fetchSubscribeMessageQueues(dlqTopic);
             if (queues == null || queues.isEmpty()) {
-                return result;
+                return new DeadLetterScanResult(result, 0);
             }
             outer:
             for (MessageQueue queue : queues) {
@@ -201,6 +218,7 @@ public class RocketMQDLQProvider implements DLQProvider {
                         PullResult pullResult = consumer.pull(queue, "*", 
offset, 32);
                         if (pullResult == null) {
                             log.warn("Stop DLQ scan for {} because queue {} 
returned no pull result", dlqTopic, queue);
+                            failedQueueCount++;
                             break;
                         }
                         long nextOffset = pullResult.getNextBeginOffset();
@@ -243,15 +261,23 @@ public class RocketMQDLQProvider implements DLQProvider {
                         }
                     }
                 } catch (Exception e) {
+                    failedQueueCount++;
                     log.warn("Failed to scan DLQ queue {} in {}: {}", queue, 
dlqTopic, e.getMessage());
                 }
             }
+            if (failedQueueCount == queues.size()) {
+                throw new BusinessException(502, "Failed to scan DLQ topic " + 
dlqTopic);
+            }
         } catch (Exception e) {
+            if (e instanceof BusinessException businessException) {
+                throw businessException;
+            }
             log.warn("Failed to collect dead letters from {}: {}", dlqTopic, 
e.getMessage());
+            throw new BusinessException(502, "Failed to scan DLQ topic " + 
dlqTopic + ": " + e.getMessage());
         } finally {
             consumer.shutdown();
         }
-        return result;
+        return new DeadLetterScanResult(result, failedQueueCount);
     }
 
     private boolean resendOne(DefaultMQProducer producer, MessageExt 
deadLetter, String targetTopic) {
@@ -335,7 +361,10 @@ public class RocketMQDLQProvider implements DLQProvider {
         return ShortLivedClientName.next("studio-dlq-resend");
     }
 
-    private String classifyOutcome(int matched, int resent, int failed) {
+    private String classifyOutcome(int matched, int resent, int failed, 
boolean scanIncomplete) {
+        if (scanIncomplete) {
+            return "PARTIAL";
+        }
         if (matched == 0) {
             return "NO_MESSAGES";
         }
@@ -355,4 +384,10 @@ public class RocketMQDLQProvider implements DLQProvider {
             log.warn("Failed to record DLQ resend audit: {}", e.getMessage());
         }
     }
+
+    private record DeadLetterScanResult(List<MessageExt> messages, int 
failedQueueCount) {
+        boolean scanIncomplete() {
+            return failedQueueCount > 0;
+        }
+    }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
index 38c347e4..a566b048 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQDLQProviderTest.java
@@ -30,6 +30,7 @@ import 
org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
 import org.apache.rocketmq.remoting.protocol.body.TopicList;
 import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.instance.dlq.DLQGroupVO;
 import org.apache.rocketmq.studio.ops.audit.AuditService;
 import org.apache.rocketmq.tools.admin.MQAdminExt;
@@ -46,6 +47,7 @@ import java.util.Set;
 import java.util.concurrent.TimeUnit;
 
 import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyInt;
 import static org.mockito.ArgumentMatchers.anyLong;
@@ -53,6 +55,7 @@ import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.ArgumentMatchers.contains;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.doNothing;
+import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.lenient;
 import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.never;
@@ -146,6 +149,60 @@ class RocketMQDLQProviderTest {
         verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
     }
 
+    @Test
+    void resendMessagesRejectsAnAllFailedDlqScan() throws Exception {
+        String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
+                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
+                         doThrow(new IllegalStateException("broker 
unavailable")).when(consumer).start();
+                         doNothing().when(consumer).shutdown();
+                     })) {
+            assertThatThrownBy(() -> provider.resendMessages("instance-a", 
"group-a", 100L, 200L, "target-topic"))
+                    .isInstanceOf(BusinessException.class)
+                    .hasMessageContaining("Failed to scan DLQ topic " + 
dlqTopic)
+                    .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(502));
+
+            verify(mockedConsumers.constructed().get(0)).shutdown();
+        }
+        verify(auditService).record(
+                eq("RESEND_DLQ"),
+                eq("group-a"),
+                contains("scanFailedQueues=all"),
+                eq("FAILED"));
+    }
+
+    @Test
+    void resendMessagesMarksAResultPartialWhenOneDlqQueueCannotBeScanned() 
throws Exception {
+        String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+        MessageQueue unavailableQueue = new MessageQueue(dlqTopic, "broker-a", 
0);
+        MessageQueue emptyQueue = new MessageQueue(dlqTopic, "broker-b", 0);
+        PullResult emptyResult = new PullResult(PullStatus.NO_NEW_MSG, 1L, 0L, 
0L, List.of());
+        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
+                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
+                         doNothing().when(consumer).start();
+                         when(consumer.fetchSubscribeMessageQueues(dlqTopic))
+                                 .thenReturn(Set.of(unavailableQueue, 
emptyQueue));
+                         when(consumer.searchOffset(eq(unavailableQueue), 
anyLong()))
+                                 .thenThrow(new IllegalStateException("broker 
unavailable"));
+                         when(consumer.searchOffset(eq(emptyQueue), 
anyLong())).thenReturn(0L);
+                         when(consumer.pull(eq(emptyQueue), eq("*"), eq(0L), 
eq(32))).thenReturn(emptyResult);
+                         doNothing().when(consumer).shutdown();
+                     });
+             MockedConstruction<DefaultMQProducer> mockedProducers =
+                     mockConstruction(DefaultMQProducer.class)) {
+            assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "target-topic"))
+                    .extracting("matched", "resent", "failed", "outcome", 
"scanIncomplete", "failedQueueCount")
+                    .containsExactly(0, 0, 0, "PARTIAL", true, 1);
+
+            assertThat(mockedProducers.constructed()).isEmpty();
+        }
+        verify(auditService).record(
+                eq("RESEND_DLQ"),
+                eq("group-a"),
+                contains("scanFailedQueues=1"),
+                eq("PARTIAL"));
+    }
+
     @Test
     @Timeout(value = 1, unit = TimeUnit.SECONDS)
     void resendMessagesStopsWhenPullOffsetDoesNotAdvance() throws Exception {
diff --git a/web/src/api/message.ts b/web/src/api/message.ts
index cffbd5c3..f70d15d2 100644
--- a/web/src/api/message.ts
+++ b/web/src/api/message.ts
@@ -69,7 +69,9 @@ export interface DLQResendResult {
   matched: number;
   resent: number;
   failed: number;
-  outcome: 'SUCCESS' | 'PARTIAL';
+  outcome: 'SUCCESS' | 'PARTIAL' | 'FAILED' | 'NO_MESSAGES';
+  scanIncomplete?: boolean;
+  failedQueueCount?: number;
 }
 
 // ─── Messages ───────────────────────────────────────────────────
diff --git a/web/src/pages/instance/__tests__/DLQPage.test.tsx 
b/web/src/pages/instance/__tests__/DLQPage.test.tsx
index 5d0e8d6c..3bbc7891 100644
--- a/web/src/pages/instance/__tests__/DLQPage.test.tsx
+++ b/web/src/pages/instance/__tests__/DLQPage.test.tsx
@@ -291,6 +291,30 @@ describe('DLQ page', () => {
     expect(await screen.findByText('DLQ provider is not 
configured')).toBeInTheDocument();
   });
 
+  it('warns when DLQ resend scans only part of the available queues', async () 
=> {
+    vi.mocked(messageService.resendDLQ).mockResolvedValue({
+      matched: 3,
+      resent: 3,
+      failed: 0,
+      outcome: 'PARTIAL',
+      scanIncomplete: true,
+      failedQueueCount: 1,
+    });
+    const user = userEvent.setup();
+    renderWithProviders(<DLQPage />);
+
+    const orderRow = (await screen.findByText('cg-order')).closest('tr');
+    if (!orderRow) throw new Error('DLQ group row not found');
+
+    await user.click(within(orderRow).getByRole('button', { name: '重投消息' }));
+    await user.type(screen.getByPlaceholderText('输入目标 Topic 名称'), 
'orders-retry');
+    await user.click(screen.getByRole('button', { name: '确认重投' }));
+
+    expect(
+      await screen.findByText('重投扫描不完整:1 个队列无法扫描,已重投 3 条'),
+    ).toBeInTheDocument();
+  });
+
   it('clears retry state before loading groups for a newly selected instance', 
async () => {
     let resolveSecondInstance!: (groups: DLQGroup[]) => void;
     vi.mocked(messageService.listDLQGroups)
diff --git a/web/src/pages/instance/dlq.tsx b/web/src/pages/instance/dlq.tsx
index 91f990d8..94a95b87 100644
--- a/web/src/pages/instance/dlq.tsx
+++ b/web/src/pages/instance/dlq.tsx
@@ -223,7 +223,11 @@ const DLQPage = () => {
         targetTopic: retryTargetTopic,
       });
       setRefreshKey((key) => key + 1);
-      if (result.failed > 0) {
+      if (result.scanIncomplete) {
+        message.warning(
+          `重投扫描不完整:${result.failedQueueCount ?? 0} 个队列无法扫描,已重投 
${result.resent} 条`,
+        );
+      } else if (result.failed > 0) {
         message.warning(`重投部分完成:成功 ${result.resent},失败 ${result.failed}`);
       } else {
         message.success(

Reply via email to