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 856b97fd7 fix(dlq): report failed resend message details (#4817)
856b97fd7 is described below
commit 856b97fd7e8cf0e66d6533c74b8bd13b07adca4a
Author: beautyarbutin <[email protected]>
AuthorDate: Thu Sep 24 10:53:07 2026 +0800
fix(dlq): report failed resend message details (#4817)
Both DLQ resend paths now report bounded per-message failure details: each
failed resend carries its DLQ message id, the resolved target topic and a
concise normalised reason. The responses keep the existing aggregate fields and
cap the details at 100 entries with a truncation flag, and the DLQ page renders
the same structured details. The `DLQResendResult` response table in
docs/api-spec.md is completed in the same commit.
---
docs/api-spec.md | 6 +
...ResendResultVO.java => DLQResendFailureVO.java} | 11 +-
.../studio/instance/dlq/DLQResendResultVO.java | 4 +
.../provider/apache/RocketMQDLQProvider.java | 183 +++++++++++++++------
.../studio/instance/dlq/DLQControllerTest.java | 21 ++-
.../provider/apache/RocketMQDLQProviderTest.java | 74 ++++++++-
web/src/api/message.ts | 8 +
web/src/i18n/translations.ts | 16 ++
web/src/pages/instance/__tests__/DLQPage.test.tsx | 79 +++++++++
web/src/pages/instance/dlq.tsx | 70 +++++++-
10 files changed, 406 insertions(+), 66 deletions(-)
diff --git a/docs/api-spec.md b/docs/api-spec.md
index 8fc071c0b..0f8dbc2f9 100644
--- a/docs/api-spec.md
+++ b/docs/api-spec.md
@@ -1571,6 +1571,12 @@ POST /api/dlq/resend
| `outcome` | `string` | 结果: `SUCCESS` / `PARTIAL` / `FAILED` / `NO_MESSAGES` |
| `scanIncomplete` | `boolean` | 是否有部分队列扫描失败 |
| `failedQueueCount` | `number` | 扫描失败的队列数 |
+| `failures` | `DLQResendFailure[]` | 逐条失败明细,最多 100 条;无失败时为空数组 |
+| `failuresTruncated` | `boolean` | 失败明细是否被截断:失败条数超过 100 时为 `true`,此时 `failed`
仍是真实的失败总数 |
+
+`DLQResendFailure` 的字段:`msgId`(`string`,重投失败的死信消息
ID)、`targetTopic`(`string`,解析出的目标
+Topic,未指定 `targetTopic` 时即原 Topic)、`reason`(`string`,归一化后的简短失败原因,最长 256 字符)。
+§9.4「重发选中的死信消息」返回同一结构,两条重投路径共用该 VO。
### 9.3 分页获取死信消息明细
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/DLQResendFailureVO.java
similarity index 85%
copy from
server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendResultVO.java
copy to
server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendFailureVO.java
index 1a7026e91..4aa053fc6 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/DLQResendFailureVO.java
@@ -21,11 +21,8 @@ import lombok.Value;
@Value
@Builder
-public class DLQResendResultVO {
- int matched;
- int resent;
- int failed;
- String outcome;
- boolean scanIncomplete;
- int failedQueueCount;
+public class DLQResendFailureVO {
+ String msgId;
+ String targetTopic;
+ String reason;
}
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 1a7026e91..b35dc5c31 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
@@ -16,6 +16,7 @@
*/
package org.apache.rocketmq.studio.instance.dlq;
+import java.util.List;
import lombok.Builder;
import lombok.Value;
@@ -28,4 +29,7 @@ public class DLQResendResultVO {
String outcome;
boolean scanIncomplete;
int failedQueueCount;
+ @Builder.Default
+ List<DLQResendFailureVO> failures = List.of();
+ boolean failuresTruncated;
}
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 a3edd6fd1..2e2c058e4 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
@@ -45,6 +45,7 @@ import org.apache.rocketmq.studio.instance.dlq.DLQGroupVO;
import org.apache.rocketmq.studio.instance.dlq.DLQMessageExcelRow;
import org.apache.rocketmq.studio.instance.dlq.DLQMessageVO;
import org.apache.rocketmq.studio.instance.dlq.DLQProvider;
+import org.apache.rocketmq.studio.instance.dlq.DLQResendFailureVO;
import org.apache.rocketmq.studio.instance.dlq.DLQResendResultVO;
import org.apache.rocketmq.studio.ops.audit.AuditService;
import org.apache.rocketmq.tools.admin.MQAdminExt;
@@ -83,6 +84,8 @@ public class RocketMQDLQProvider implements DLQProvider {
private static final int RESEND_HARD_CAP = 5000;
private static final int MAX_PAGE_SIZE = 100;
private static final int MAX_CONSECUTIVE_OFFSET_ILLEGAL = 3;
+ private static final int MAX_REPORTED_RESEND_FAILURES = 100;
+ private static final int MAX_FAILURE_REASON_LENGTH = 256;
private static final String ORIGIN_MESSAGE_ID_PROPERTY =
"studio_dlq_origin_message_id";
private static final String ORIGIN_TOPIC_PROPERTY =
"studio_dlq_origin_topic";
@@ -212,33 +215,17 @@ public class RocketMQDLQProvider implements DLQProvider {
throw new BusinessException(404, "No dead-letter queue found for
consumer group: " + groupName);
}
List<MessageExt> deadLetters = scanResult.messages();
- int[] counts = {0, 0};
- if (!deadLetters.isEmpty()) {
- try {
- runtimeAdminClientResolver.executeProducer(instanceId,
producer -> {
- for (MessageExt deadLetter : deadLetters) {
- if (resendOne(producer, deadLetter, targetTopic)) {
- counts[0]++;
- } else {
- counts[1]++;
- }
- }
- return null;
- });
- } catch (Exception e) {
- log.warn("Failed to resend dead letters for group {}: {}",
groupName, e.getMessage());
- counts[1] += deadLetters.size() - counts[0];
- }
- }
- int resent = counts[0];
- int failed = counts[1];
+ ResendBatch resendBatch = resendAll(instanceId, groupName,
deadLetters, targetTopic);
+ int resent = resendBatch.resent();
+ int failed = resendBatch.failed();
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, scanTruncated=%s,
scanFailedQueues=%d",
+ + "failed=%d, scanIncomplete=%s, scanTruncated=%s,
scanFailedQueues=%d, "
+ + "reportedFailures=%d, failuresTruncated=%s",
instanceId, groupName, dlqTopic,
StringUtils.hasText(targetTopic) ? targetTopic : "<original>",
deadLetters.size(), resent, failed,
scanResult.scanIncomplete(), scanResult.truncated(),
- scanResult.failedQueueCount());
+ scanResult.failedQueueCount(), resendBatch.failures().size(),
resendBatch.failuresTruncated());
recordAudit(groupName, detail, outcome);
log.info("DLQ resend completed: {}", detail);
return DLQResendResultVO.builder()
@@ -248,6 +235,8 @@ public class RocketMQDLQProvider implements DLQProvider {
.outcome(outcome)
.scanIncomplete(scanResult.scanIncomplete())
.failedQueueCount(scanResult.failedQueueCount())
+ .failures(resendBatch.failures())
+ .failuresTruncated(resendBatch.failuresTruncated())
.build();
}
@@ -294,31 +283,16 @@ public class RocketMQDLQProvider implements DLQProvider {
}
return resolved;
});
- int[] counts = {0, 0};
- if (!deadLetters.isEmpty()) {
- try {
- runtimeAdminClientResolver.executeProducer(instanceId,
producer -> {
- for (MessageExt deadLetter : deadLetters) {
- if (resendOne(producer, deadLetter, targetTopic)) {
- counts[0]++;
- } else {
- counts[1]++;
- }
- }
- return null;
- });
- } catch (Exception e) {
- log.warn("Failed to resend selected dead letters for group {}:
{}", groupName, e.getMessage());
- counts[1] += deadLetters.size() - counts[0];
- }
- }
- int resent = counts[0];
- int failed = counts[1];
+ ResendBatch resendBatch = resendAll(instanceId, groupName,
deadLetters, targetTopic);
+ int resent = resendBatch.resent();
+ int failed = resendBatch.failed();
boolean foundAll = deadLetters.size() == selected.size();
String outcome = classifyOutcome(deadLetters.size(), resent, failed,
!foundAll);
- String detail = String.format("instanceId=%s, group=%s, selected=%d,
matched=%d, resent=%d, failed=%d",
- instanceId, groupName, selected.size(), deadLetters.size(),
resent, failed);
+ String detail = String.format("instanceId=%s, group=%s, selected=%d,
matched=%d, resent=%d, failed=%d, "
+ + "reportedFailures=%d, failuresTruncated=%s",
+ instanceId, groupName, selected.size(), deadLetters.size(),
resent, failed,
+ resendBatch.failures().size(),
resendBatch.failuresTruncated());
recordAudit(groupName, detail, outcome);
log.info("Selected DLQ resend completed: {}", detail);
return DLQResendResultVO.builder()
@@ -327,6 +301,8 @@ public class RocketMQDLQProvider implements DLQProvider {
.failed(failed)
.outcome(outcome)
.scanIncomplete(!foundAll)
+ .failures(resendBatch.failures())
+ .failuresTruncated(resendBatch.failuresTruncated())
.build();
}
@@ -573,11 +549,38 @@ public class RocketMQDLQProvider implements DLQProvider {
return false;
}
- private boolean resendOne(DefaultMQProducer producer, MessageExt
deadLetter, String targetTopic) {
+ private ResendBatch resendAll(String instanceId, String groupName,
List<MessageExt> deadLetters,
+ String targetTopic) {
+ ResendBatch batch = new ResendBatch();
+ if (deadLetters.isEmpty()) {
+ return batch;
+ }
+ try {
+ runtimeAdminClientResolver.executeProducer(instanceId, producer ->
{
+ for (MessageExt deadLetter : deadLetters) {
+ batch.record(resendOne(producer, deadLetter, targetTopic));
+ }
+ return null;
+ });
+ } catch (Exception e) {
+ log.warn("Failed to use a producer while resending dead letters
for group {}: {}",
+ groupName, e.getMessage());
+ String reason = failureReason("Producer session failed", e);
+ while (batch.processed() < deadLetters.size()) {
+ MessageExt deadLetter = deadLetters.get(batch.processed());
+ batch.record(ResendAttempt.failure(buildFailure(
+ deadLetter, resolveTargetTopic(deadLetter,
targetTopic), reason)));
+ }
+ }
+ return batch;
+ }
+
+ private ResendAttempt resendOne(DefaultMQProducer producer, MessageExt
deadLetter, String targetTopic) {
String destination = resolveTargetTopic(deadLetter, targetTopic);
if (!StringUtils.hasText(destination)) {
log.warn("Skip resend of msgId={}: no target topic resolvable",
deadLetter.getMsgId());
- return false;
+ return ResendAttempt.failure(buildFailure(
+ deadLetter, null, "No target topic could be resolved"));
}
try {
Message message = new Message(destination, deadLetter.getBody());
@@ -608,16 +611,46 @@ public class RocketMQDLQProvider implements DLQProvider {
log.warn("DLQ resend was not accepted: msgId={} topic={}
sendStatus={}",
deadLetter.getMsgId(), destination,
sendResult == null ? "<null>" :
sendResult.getSendStatus());
- return false;
+ String status = sendResult == null ? "no result" :
sendResult.getSendStatus().name();
+ return ResendAttempt.failure(buildFailure(
+ deadLetter, destination, "Producer returned " +
status));
}
log.debug("Resent dead letter msgId={} to topic={}, sendStatus={}",
deadLetter.getMsgId(), destination,
sendResult.getSendStatus());
- return true;
+ return ResendAttempt.success();
} catch (Exception e) {
log.warn("Failed to resend dead letter msgId={} to topic={}: {}",
deadLetter.getMsgId(), destination, e.getMessage());
- return false;
+ return ResendAttempt.failure(buildFailure(
+ deadLetter, destination, failureReason("Producer send
failed", e)));
+ }
+ }
+
+ private DLQResendFailureVO buildFailure(MessageExt deadLetter, String
targetTopic, String reason) {
+ String msgId = StringUtils.hasText(deadLetter.getMsgId()) ?
deadLetter.getMsgId().trim() : "<unknown>";
+ return DLQResendFailureVO.builder()
+ .msgId(msgId)
+ .targetTopic(StringUtils.hasText(targetTopic) ?
targetTopic.trim() : null)
+ .reason(normalizeFailureReason(reason))
+ .build();
+ }
+
+ private static String failureReason(String prefix, Exception exception) {
+ String detail = exception.getMessage();
+ if (!StringUtils.hasText(detail)) {
+ detail = exception.getClass().getSimpleName();
}
+ return normalizeFailureReason(prefix + ": " + detail);
+ }
+
+ private static String normalizeFailureReason(String reason) {
+ String normalized = StringUtils.hasText(reason)
+ ? reason.replaceAll("\\s+", " ").trim()
+ : "Unknown resend failure";
+ if (normalized.length() <= MAX_FAILURE_REASON_LENGTH) {
+ return normalized;
+ }
+ return normalized.substring(0, MAX_FAILURE_REASON_LENGTH - 3) + "...";
}
/**
@@ -704,4 +737,56 @@ public class RocketMQDLQProvider implements DLQProvider {
return failedQueueCount > 0 || truncated;
}
}
+
+ private record ResendAttempt(boolean successful, DLQResendFailureVO
failure) {
+ private static ResendAttempt success() {
+ return new ResendAttempt(true, null);
+ }
+
+ private static ResendAttempt failure(DLQResendFailureVO failure) {
+ return new ResendAttempt(false, failure);
+ }
+ }
+
+ private static final class ResendBatch {
+ private int processed;
+ private int resent;
+ private int failed;
+ private boolean failuresTruncated;
+ private final List<DLQResendFailureVO> failures = new ArrayList<>();
+
+ private void record(ResendAttempt attempt) {
+ processed++;
+ if (attempt.successful()) {
+ resent++;
+ return;
+ }
+ failed++;
+ if (failures.size() < MAX_REPORTED_RESEND_FAILURES) {
+ failures.add(attempt.failure());
+ } else {
+ failuresTruncated = true;
+ }
+ }
+
+ private int processed() {
+ return processed;
+ }
+
+ private int resent() {
+ return resent;
+ }
+
+ private int failed() {
+ return failed;
+ }
+
+ private List<DLQResendFailureVO> failures() {
+ return List.copyOf(failures);
+ }
+
+ private boolean failuresTruncated() {
+ return failuresTruncated;
+ }
+ }
}
diff --git
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
index 0be329af2..7ed1d9b24 100644
---
a/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/instance/dlq/DLQControllerTest.java
@@ -127,13 +127,32 @@ class DLQControllerTest extends WebMvcAuthTestSupport {
"endTime", 2000,
"targetTopic", "target-topic"
);
+ when(dlqService.resendMessages(
+ "instance-1", "test-group", 1000L, 2000L, "target-topic"))
+ .thenReturn(DLQResendResultVO.builder()
+ .matched(2)
+ .resent(1)
+ .failed(1)
+ .outcome("PARTIAL")
+ .failures(List.of(DLQResendFailureVO.builder()
+ .msgId("failed-msg")
+ .targetTopic("target-topic")
+ .reason("Producer returned FLUSH_DISK_TIMEOUT")
+ .build()))
+ .failuresTruncated(false)
+ .build());
mockMvc.perform(post("/api/dlq/resend")
.contentType(MediaType.APPLICATION_JSON)
.content(objectMapper.writeValueAsString(body)))
.andExpect(status().isOk())
.andExpect(jsonPath("$.code").value(200))
- .andExpect(jsonPath("$.message").value("success"));
+ .andExpect(jsonPath("$.message").value("success"))
+
.andExpect(jsonPath("$.data.failures[0].msgId").value("failed-msg"))
+
.andExpect(jsonPath("$.data.failures[0].targetTopic").value("target-topic"))
+ .andExpect(jsonPath("$.data.failures[0].reason")
+ .value("Producer returned FLUSH_DISK_TIMEOUT"))
+ .andExpect(jsonPath("$.data.failuresTruncated").value(false));
verify(dlqService).resendMessages(
eq("instance-1"), eq("test-group"), eq(1000L), eq(2000L),
eq("target-topic"));
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 01b780a3b..a721a9d34 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
@@ -25,6 +25,7 @@ import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.client.producer.SendStatus;
import org.apache.rocketmq.common.MixAll;
import org.apache.rocketmq.common.message.Message;
+import org.apache.rocketmq.common.message.MessageAccessor;
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageDecoder;
import org.apache.rocketmq.common.message.MessageExt;
@@ -715,9 +716,17 @@ class RocketMQDLQProviderTest {
TopicList existingTargets = new TopicList();
existingTargets.setTopicList(Set.of("target-topic"));
when(adminExt.fetchAllTopicList()).thenReturn(existingTargets);
- assertThat(provider.resendMessages("instance-a", "group-a", 100L,
200L, "target-topic"))
+ DLQResendResultVO result = provider.resendMessages(
+ "instance-a", "group-a", 100L, 200L, "target-topic");
+ assertThat(result)
.extracting("matched", "resent", "failed", "outcome")
.containsExactly(1, 0, 1, "FAILED");
+ assertThat(result.getFailures()).singleElement().satisfies(failure -> {
+ assertThat(failure.getMsgId()).isEqualTo("msg-1");
+ assertThat(failure.getTargetTopic()).isEqualTo("target-topic");
+ assertThat(failure.getReason()).contains("FLUSH_DISK_TIMEOUT");
+ });
+ assertThat(result.isFailuresTruncated()).isFalse();
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
verify(runtimeAdminClientResolver).executeProducer(eq("instance-a"),
any());
@@ -823,6 +832,69 @@ class RocketMQDLQProviderTest {
verify(adminExt).viewMessage(dlqTopic, "found-msg");
}
+ @Test
+ void resendSelectedMessagesReportsProducerFailureDetailsTest() throws
Exception {
+ String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+ MessageExt deadLetter = new MessageExt();
+ deadLetter.setMsgId("selected-msg");
+ deadLetter.setTopic(dlqTopic);
+ deadLetter.setBody(new byte[] {1});
+ MessageAccessor.putProperty(deadLetter,
MessageConst.PROPERTY_DLQ_ORIGIN_TOPIC, "orders");
+
+ when(adminExt.viewMessage(dlqTopic,
"selected-msg")).thenReturn(deadLetter);
+ when(dlqProducer.send(any(Message.class))).thenThrow(new
IllegalStateException("broker unavailable"));
+
+ DLQResendResultVO result = provider.resendMessages(
+ "instance-a", "group-a", List.of("selected-msg"), null);
+
+ assertThat(result)
+ .extracting("matched", "resent", "failed", "outcome")
+ .containsExactly(1, 0, 1, "FAILED");
+ assertThat(result.getFailures()).singleElement().satisfies(failure -> {
+ assertThat(failure.getMsgId()).isEqualTo("selected-msg");
+ assertThat(failure.getTargetTopic()).isEqualTo("orders");
+ assertThat(failure.getReason()).contains("broker unavailable");
+ });
+ assertThat(result.isFailuresTruncated()).isFalse();
+ }
+
+ @Test
+ void resendMessagesCapsReportedFailureDetailsTest() throws Exception {
+ String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
+ MessageQueue queue = new MessageQueue(dlqTopic, "broker-a", 0);
+ List<MessageExt> deadLetters = IntStream.range(0, 101)
+ .mapToObj(index -> {
+ MessageExt deadLetter = new MessageExt();
+ deadLetter.setMsgId("msg-" + index);
+ deadLetter.setTopic(dlqTopic);
+ deadLetter.setBody(new byte[] {1});
+ deadLetter.setStoreTimestamp(150L);
+ return deadLetter;
+ })
+ .toList();
+ PullResult pullResult = new PullResult(PullStatus.FOUND, 101L, 0L,
101L, deadLetters);
+ SendResult sendResult = new SendResult();
+ sendResult.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT);
+
+
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+ when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+ when(pullConsumer.searchOffset(queue, 200L)).thenReturn(100L);
+ when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+ when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+ TopicList existingTargets = new TopicList();
+ existingTargets.setTopicList(Set.of("target-topic"));
+ when(adminExt.fetchAllTopicList()).thenReturn(existingTargets);
+
+ DLQResendResultVO result = provider.resendMessages(
+ "instance-a", "group-a", 100L, 200L, "target-topic");
+
+ assertThat(result.getFailed()).isEqualTo(101);
+ assertThat(result.getFailures()).hasSize(100);
+ assertThat(result.getFailures()).extracting("msgId")
+ .containsExactlyElementsOf(IntStream.range(0,
100).mapToObj(index -> "msg-" + index).toList());
+ assertThat(result.isFailuresTruncated()).isTrue();
+ }
+
@Test
void resendMessagesMarksAResultPartialWhenScanReachesHardCap() throws
Exception {
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
diff --git a/web/src/api/message.ts b/web/src/api/message.ts
index b8b4c4ea9..bb31baabf 100644
--- a/web/src/api/message.ts
+++ b/web/src/api/message.ts
@@ -111,6 +111,14 @@ export interface DLQResendResult {
outcome: 'SUCCESS' | 'PARTIAL' | 'FAILED' | 'NO_MESSAGES';
scanIncomplete?: boolean;
failedQueueCount?: number;
+ failures?: DLQResendFailure[];
+ failuresTruncated?: boolean;
+}
+
+export interface DLQResendFailure {
+ msgId: string;
+ targetTopic?: string | null;
+ reason: string;
}
export interface DLQMessage {
diff --git a/web/src/i18n/translations.ts b/web/src/i18n/translations.ts
index b78d1c285..bb0261ad1 100644
--- a/web/src/i18n/translations.ts
+++ b/web/src/i18n/translations.ts
@@ -845,6 +845,22 @@ const translations: Record<string, Record<Lang, string>> =
{
// ─── Dead Letter Queue ───
'dlq.title': { zh: '死信队列', en: 'Dead Letter Queue' },
+ 'dlq.resendPartialSummary': {
+ zh: '重投部分完成:成功 {resent},失败 {failed}',
+ en: 'Resend partially completed: {resent} succeeded, {failed} failed',
+ },
+ 'dlq.resendFailedSummary': {
+ zh: '重投失败:成功 {resent},失败 {failed}',
+ en: 'Resend failed: {resent} succeeded, {failed} failed',
+ },
+ 'dlq.failureDetails': { zh: '失败消息明细', en: 'Failed message details' },
+ 'dlq.failureMessageId': { zh: '消息 ID', en: 'Message ID' },
+ 'dlq.failureTargetTopic': { zh: '目标 Topic', en: 'Target topic' },
+ 'dlq.failureReason': { zh: '失败原因', en: 'Failure reason' },
+ 'dlq.failureDetailsTruncated': {
+ zh: '失败明细较多,仅显示前 100 条。',
+ en: 'Only the first 100 failure details are shown.',
+ },
// ─── Client Connections ───
'clients.title': { zh: '客户端连接', en: 'Client Connections' },
diff --git a/web/src/pages/instance/__tests__/DLQPage.test.tsx
b/web/src/pages/instance/__tests__/DLQPage.test.tsx
index 791a6c799..b84e2a421 100644
--- a/web/src/pages/instance/__tests__/DLQPage.test.tsx
+++ b/web/src/pages/instance/__tests__/DLQPage.test.tsx
@@ -610,6 +610,85 @@ describe('DLQ page', () => {
).toBeInTheDocument();
});
+ it('shows per-message failure details for a partial range resend', async ()
=> {
+ vi.mocked(messageService.resendDLQ).mockResolvedValue({
+ matched: 2,
+ resent: 1,
+ failed: 1,
+ outcome: 'PARTIAL',
+ failures: [
+ {
+ msgId: 'failed-range-msg',
+ targetTopic: 'orders-retry',
+ reason: 'Producer returned FLUSH_DISK_TIMEOUT',
+ },
+ ],
+ failuresTruncated: false,
+ });
+ 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('failed-range-msg')).toBeInTheDocument();
+ expect(screen.getByText('orders-retry')).toBeInTheDocument();
+ expect(screen.getByText('Producer returned
FLUSH_DISK_TIMEOUT')).toBeInTheDocument();
+ });
+
+ it('shows per-message failure details for selected-message resend', async ()
=> {
+ vi.mocked(messageService.listDLQMessages).mockResolvedValue({
+ items: [
+ {
+ msgId: 'msg-1',
+ topic: '%DLQ%cg-order',
+ queueId: 0,
+ offset: 1,
+ storeTime: 2,
+ reconsumeTimes: 3,
+ keys: null,
+ body: 'payload',
+ bodyBase64: null,
+ properties: {},
+ },
+ ],
+ total: 1,
+ page: 1,
+ size: 20,
+ });
+ vi.mocked(messageService.resendDLQSelected).mockResolvedValue({
+ matched: 1,
+ resent: 0,
+ failed: 1,
+ outcome: 'FAILED',
+ failures: [
+ {
+ msgId: 'msg-1',
+ targetTopic: 'orders',
+ reason: 'Producer send failed: broker unavailable',
+ },
+ ],
+ failuresTruncated: false,
+ });
+ 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: /消息明细/ }));
+ const messageRow = (await screen.findByText('msg-1')).closest('tr');
+ if (!messageRow) throw new Error('DLQ message row not found');
+ await user.click(within(messageRow).getByRole('checkbox'));
+ await user.click(screen.getByRole('button', { name: /批量重发选中/ }));
+
+ expect(await screen.findByText('Producer send failed: broker
unavailable')).toBeInTheDocument();
+ expect(screen.getAllByText('msg-1').length).toBeGreaterThanOrEqual(2);
+ expect(screen.getByText('orders')).toBeInTheDocument();
+ });
+
it('clears retry state before loading groups for a newly selected instance',
async () => {
let resolveSecondInstance!: (page: DLQGroupPage) => void;
vi.mocked(messageService.listDLQGroups)
diff --git a/web/src/pages/instance/dlq.tsx b/web/src/pages/instance/dlq.tsx
index ee6472242..21da027ee 100644
--- a/web/src/pages/instance/dlq.tsx
+++ b/web/src/pages/instance/dlq.tsx
@@ -39,7 +39,7 @@ import PageHeader from '../../components/PageHeader';
import InfoBanner from '../../components/InfoBanner';
import { InstanceSelect } from '../../components/InstanceSelect';
import { useLang } from '../../i18n/LangContext';
-import type { DLQGroup, DLQMessage } from '../../api/message';
+import type { DLQGroup, DLQMessage, DLQResendResult } from '../../api/message';
import {
exportDLQExcel,
listDLQGroups,
@@ -237,6 +237,57 @@ const DLQPage = () => {
setRetryModalOpen(true);
};
+ const showResendFailures = (result: DLQResendResult) => {
+ const summaryKey =
+ result.outcome === 'FAILED' ? 'dlq.resendFailedSummary' :
'dlq.resendPartialSummary';
+ const summary = t(summaryKey, { resent: result.resent, failed:
result.failed });
+ const failures = result.failures ?? [];
+ if (failures.length === 0) {
+ if (result.outcome === 'FAILED') {
+ message.error(summary);
+ } else {
+ message.warning(summary);
+ }
+ return;
+ }
+
+ const openFailureModal = result.outcome === 'FAILED' ? Modal.error :
Modal.warning;
+ openFailureModal({
+ title: summary,
+ width: 720,
+ content: (
+ <Space direction="vertical" size="middle" style={{ width: '100%' }}>
+ <Text strong>{t('dlq.failureDetails')}</Text>
+ <div role="list" style={{ maxHeight: 360, overflowY: 'auto' }}>
+ {failures.map((failure, index) => (
+ <div
+ role="listitem"
+ key={`${failure.msgId}-${index}`}
+ style={{ borderBottom: '1px solid #f0f0f0', padding: '8px 0' }}
+ >
+ <div>
+ <Text type="secondary">{t('dlq.failureMessageId')}: </Text>
+ <Text code>{failure.msgId || '-'}</Text>
+ </div>
+ <div>
+ <Text type="secondary">{t('dlq.failureTargetTopic')}: </Text>
+ <Text code>{failure.targetTopic || '-'}</Text>
+ </div>
+ <div>
+ <Text type="secondary">{t('dlq.failureReason')}: </Text>
+ <Text>{failure.reason}</Text>
+ </div>
+ </div>
+ ))}
+ </div>
+ {result.failuresTruncated && (
+ <Alert type="warning" showIcon
message={t('dlq.failureDetailsTruncated')} />
+ )}
+ </Space>
+ ),
+ });
+ };
+
const handleRetry = async () => {
if (!retryTargetTopic) {
message.warning('请输入目标 Topic');
@@ -262,12 +313,17 @@ const DLQPage = () => {
});
if (retryRequestIdRef.current !== requestId) return;
setRefreshKey((key) => key + 1);
- if (result.scanIncomplete) {
+ if (result.failed > 0) {
+ showResendFailures(result);
+ if (result.scanIncomplete) {
+ message.warning(
+ `重投扫描不完整:${result.failedQueueCount ?? 0} 个队列无法扫描,已重投
${result.resent} 条`,
+ );
+ }
+ } else if (result.scanIncomplete) {
message.warning(
`重投扫描不完整:${result.failedQueueCount ?? 0} 个队列无法扫描,已重投
${result.resent} 条`,
);
- } else if (result.failed > 0) {
- message.warning(`重投部分完成:成功 ${result.resent},失败 ${result.failed}`);
} else {
message.success(`重投完成:${groupName} → ${targetTopic}(${result.resent}
条)`);
}
@@ -375,10 +431,8 @@ const DLQPage = () => {
msgIds,
});
if (detailResendRequestIdRef.current !== requestId) return;
- if (result.outcome === 'FAILED' && result.failed > 0) {
- message.error(`重发失败:成功 ${result.resent},失败 ${result.failed}`);
- } else if (result.resent > 0 && result.failed > 0) {
- message.warning(`重发部分完成:成功 ${result.resent},失败 ${result.failed}`);
+ if (result.failed > 0) {
+ showResendFailures(result);
} else {
message.success(`重发完成:成功 ${result.resent} 条`);
}