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 8cd898a2 fix(dlq): report resend outcomes and isolate short-lived
clients (#993)
8cd898a2 is described below
commit 8cd898a23522c6880c81ac6c0daf96b7cef55a49
Author: aias00 <[email protected]>
AuthorDate: Wed Aug 5 02:40:11 2026 -0700
fix(dlq): report resend outcomes and isolate short-lived clients (#993)
---
.../rocketmq/studio/instance/dlq/DLQController.java | 7 +++----
.../rocketmq/studio/instance/dlq/DLQProvider.java | 2 +-
.../studio/instance/dlq/DLQProviderStub.java | 2 +-
.../{DLQProvider.java => DLQResendResultVO.java} | 14 +++++++++-----
.../rocketmq/studio/instance/dlq/DLQService.java | 4 ++--
.../studio/rocketmq/RocketMQDLQProvider.java | 13 ++++++++++---
.../studio/rocketmq/RocketMQMessageProvider.java | 2 +-
.../ShortLivedClientName.java} | 14 +++++++++-----
.../studio/rocketmq/ShortLivedClientNameTest.java} | 21 ++++++++++++++++-----
web/src/api/dlq.test.ts | 5 +++--
web/src/api/message.ts | 16 +++++++++++++---
web/src/pages/instance/dlq.tsx | 10 ++++++++--
web/src/services/messageService.ts | 12 +++++++++---
13 files changed, 85 insertions(+), 37 deletions(-)
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQController.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQController.java
index 5b5517f7..872cdd6a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQController.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQController.java
@@ -42,11 +42,10 @@ public class DLQController {
}
@PostMapping("/resend")
- public Result<Void> resendMessages(@Valid @RequestBody(required = false)
DLQResendRequestDTO request) {
+ public Result<DLQResendResultVO> resendMessages(@Valid
@RequestBody(required = false) DLQResendRequestDTO request) {
requireRequest(request);
- dlqService.resendMessages(
- request.getGroupName(), request.getStartTime(),
request.getEndTime(), request.getTargetTopic());
- return Result.ok();
+ return Result.ok(dlqService.resendMessages(
+ request.getGroupName(), request.getStartTime(),
request.getEndTime(), request.getTargetTopic()));
}
private void requireRequest(DLQResendRequestDTO request) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
index bb732aa5..f143a68c 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
@@ -21,5 +21,5 @@ import java.util.List;
public interface DLQProvider {
List<DLQGroupVO> listDLQGroups(String clusterId);
- void resendMessages(String groupName, Long startTime, Long endTime, String
targetTopic);
+ DLQResendResultVO resendMessages(String groupName, Long startTime, Long
endTime, String targetTopic);
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
index 9a301a3b..ad6e3bb3 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProviderStub.java
@@ -36,7 +36,7 @@ public class DLQProviderStub implements DLQProvider {
}
@Override
- public void resendMessages(String groupName, Long startTime, Long endTime,
String targetTopic) {
+ public DLQResendResultVO resendMessages(String groupName, Long startTime,
Long endTime, String targetTopic) {
log.warn("DLQProviderStub.resendMessages called but no real DLQ
provider is configured. group={}, targetTopic={}",
groupName, targetTopic);
throw unsupported();
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendResultVO.java
similarity index 80%
copy from
server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
copy to
server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendResultVO.java
index bb732aa5..36cf13b0 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQResendResultVO.java
@@ -16,10 +16,14 @@
*/
package org.apache.rocketmq.studio.instance.dlq;
+import lombok.Builder;
+import lombok.Value;
-import java.util.List;
-
-public interface DLQProvider {
- List<DLQGroupVO> listDLQGroups(String clusterId);
- void resendMessages(String groupName, Long startTime, Long endTime, String
targetTopic);
+@Value
+@Builder
+public class DLQResendResultVO {
+ int matched;
+ int resent;
+ int failed;
+ String outcome;
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
index a829c07f..a18e4e5a 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQService.java
@@ -37,10 +37,10 @@ public class DLQService {
return dlqProvider.listDLQGroups(clusterId);
}
- public void resendMessages(String groupName, Long startTime, Long endTime,
String targetTopic) {
+ public DLQResendResultVO resendMessages(String groupName, Long startTime,
Long endTime, String targetTopic) {
validateResendRequest(groupName, startTime, endTime);
log.info("Resending DLQ messages: group={}, targetTopic={}",
groupName, targetTopic);
- dlqProvider.resendMessages(groupName, startTime, endTime, targetTopic);
+ return dlqProvider.resendMessages(groupName, startTime, endTime,
targetTopic);
}
private void validateResendRequest(String groupName, Long startTime, Long
endTime) {
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
index 426c08bd..c4fdeaf5 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQDLQProvider.java
@@ -31,6 +31,7 @@ import
org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
import org.apache.rocketmq.remoting.protocol.body.TopicList;
import org.apache.rocketmq.studio.instance.dlq.DLQGroupVO;
import org.apache.rocketmq.studio.instance.dlq.DLQProvider;
+import org.apache.rocketmq.studio.instance.dlq.DLQResendResultVO;
import org.apache.rocketmq.studio.ops.audit.AuditService;
import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
import org.slf4j.Logger;
@@ -142,7 +143,7 @@ public class RocketMQDLQProvider implements DLQProvider {
}
@Override
- public void resendMessages(String groupName, Long startTime, Long endTime,
String targetTopic) {
+ public DLQResendResultVO resendMessages(String groupName, Long startTime,
Long endTime, String targetTopic) {
DefaultMQAdminExt adminExt = adminExtProvider.getIfAvailable();
String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + groupName;
@@ -176,6 +177,12 @@ public class RocketMQDLQProvider implements DLQProvider {
deadLetters.size(), resent, failed);
recordAudit(groupName, detail, failed == 0 ? "SUCCESS" : "PARTIAL");
log.info("DLQ resend completed: {}", detail);
+ return DLQResendResultVO.builder()
+ .matched(deadLetters.size())
+ .resent(resent)
+ .failed(failed)
+ .outcome(failed == 0 ? "SUCCESS" : "PARTIAL")
+ .build();
}
private List<MessageExt> collectDeadLetters(String dlqTopic, long begin,
long end) {
@@ -267,7 +274,7 @@ public class RocketMQDLQProvider implements DLQProvider {
private DefaultMQPullConsumer newPullConsumer() {
DefaultMQPullConsumer consumer = new
DefaultMQPullConsumer("studio-dlq-query-group");
- consumer.setInstanceName("studio-dlq-query-" +
System.currentTimeMillis());
+
consumer.setInstanceName(ShortLivedClientName.next("studio-dlq-query"));
if (StringUtils.hasText(properties.getNamesrvAddr())) {
consumer.setNamesrvAddr(properties.getNamesrvAddr());
}
@@ -276,7 +283,7 @@ public class RocketMQDLQProvider implements DLQProvider {
private DefaultMQProducer newProducer(String groupName) {
DefaultMQProducer producer = new
DefaultMQProducer("studio-dlq-resend-" + groupName);
- producer.setInstanceName("studio-dlq-resend-" +
System.currentTimeMillis());
+
producer.setInstanceName(ShortLivedClientName.next("studio-dlq-resend"));
producer.setRetryTimesWhenSendFailed(2);
if (StringUtils.hasText(properties.getNamesrvAddr())) {
producer.setNamesrvAddr(properties.getNamesrvAddr());
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
index b361fa2e..1657d4e6 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/RocketMQMessageProvider.java
@@ -410,7 +410,7 @@ public class RocketMQMessageProvider implements
MessageProvider {
private DefaultMQPullConsumer newPullConsumer(String groupPrefix) {
DefaultMQPullConsumer consumer = new DefaultMQPullConsumer(groupPrefix
+ "-group");
- consumer.setInstanceName(groupPrefix + "-" +
System.currentTimeMillis());
+ consumer.setInstanceName(ShortLivedClientName.next(groupPrefix));
if (StringUtils.hasText(properties.getNamesrvAddr())) {
consumer.setNamesrvAddr(properties.getNamesrvAddr());
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/ShortLivedClientName.java
similarity index 76%
copy from
server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
copy to
server/src/main/java/org/apache/rocketmq/studio/rocketmq/ShortLivedClientName.java
index bb732aa5..b2ee385d 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
+++
b/server/src/main/java/org/apache/rocketmq/studio/rocketmq/ShortLivedClientName.java
@@ -14,12 +14,16 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.instance.dlq;
+package org.apache.rocketmq.studio.rocketmq;
+import java.util.UUID;
-import java.util.List;
+final class ShortLivedClientName {
-public interface DLQProvider {
- List<DLQGroupVO> listDLQGroups(String clusterId);
- void resendMessages(String groupName, Long startTime, Long endTime, String
targetTopic);
+ private ShortLivedClientName() {
+ }
+
+ static String next(String prefix) {
+ return prefix + "-" + UUID.randomUUID();
+ }
}
diff --git
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/ShortLivedClientNameTest.java
similarity index 57%
copy from
server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
copy to
server/src/test/java/org/apache/rocketmq/studio/rocketmq/ShortLivedClientNameTest.java
index bb732aa5..1449a186 100644
---
a/server/src/main/java/org/apache/rocketmq/studio/instance/dlq/DLQProvider.java
+++
b/server/src/test/java/org/apache/rocketmq/studio/rocketmq/ShortLivedClientNameTest.java
@@ -14,12 +14,23 @@
* See the License for the specific language governing permissions and
* limitations under the License.
*/
-package org.apache.rocketmq.studio.instance.dlq;
+package org.apache.rocketmq.studio.rocketmq;
+import static org.assertj.core.api.Assertions.assertThat;
-import java.util.List;
+import java.util.Set;
+import java.util.stream.Collectors;
+import java.util.stream.IntStream;
+import org.junit.jupiter.api.Test;
-public interface DLQProvider {
- List<DLQGroupVO> listDLQGroups(String clusterId);
- void resendMessages(String groupName, Long startTime, Long endTime, String
targetTopic);
+class ShortLivedClientNameTest {
+
+ @Test
+ void generatesUniqueNamesWithoutClockDelays() {
+ Set<String> names = IntStream.range(0, 100)
+ .mapToObj(ignored -> ShortLivedClientName.next("studio-query"))
+ .collect(Collectors.toSet());
+
+ assertThat(names).hasSize(100).allMatch(name ->
name.startsWith("studio-query-"));
+ }
}
diff --git a/web/src/api/dlq.test.ts b/web/src/api/dlq.test.ts
index de4e0b69..6adca0af 100644
--- a/web/src/api/dlq.test.ts
+++ b/web/src/api/dlq.test.ts
@@ -55,11 +55,12 @@ describe('DLQ API', () => {
endTime: 1784332800000,
targetTopic: 'orders-retry',
};
+ const result = { matched: 2, resent: 2, failed: 0, outcome: 'SUCCESS' };
mock.onPost('/dlq/resend').reply((config) => {
expect(JSON.parse(config.data)).toEqual(payload);
- return [200, { code: 200, data: null }];
+ return [200, { code: 200, data: result }];
});
- await expect(resendDLQ(payload)).resolves.toBeUndefined();
+ await expect(resendDLQ(payload)).resolves.toEqual(result);
});
});
diff --git a/web/src/api/message.ts b/web/src/api/message.ts
index 0c879543..b103745c 100644
--- a/web/src/api/message.ts
+++ b/web/src/api/message.ts
@@ -63,6 +63,13 @@ export interface DLQGroup {
status: string;
}
+export interface DLQResendResult {
+ matched: number;
+ resent: number;
+ failed: number;
+ outcome: 'SUCCESS' | 'PARTIAL';
+}
+
// ─── Messages ───────────────────────────────────────────────────
export async function queryMessages(params: MessageQuery) {
const res = await client.get<{ data: MessageRecord[] }>('/messages', {
params });
@@ -70,7 +77,9 @@ export async function queryMessages(params: MessageQuery) {
}
export async function getMessageTrace(msgId: string) {
- const res = await client.get<{ data: TraceRecord
}>(`/messages/${encodeURIComponent(msgId)}/trace`);
+ const res = await client.get<{ data: TraceRecord }>(
+ `/messages/${encodeURIComponent(msgId)}/trace`,
+ );
return res.data.data;
}
@@ -85,6 +94,7 @@ export async function resendDLQ(data: {
startTime: number;
endTime: number;
targetTopic?: string;
-}) {
- await client.post('/dlq/resend', data);
+}): Promise<DLQResendResult> {
+ const res = await client.post<{ data: DLQResendResult }>('/dlq/resend',
data);
+ return res.data.data;
}
diff --git a/web/src/pages/instance/dlq.tsx b/web/src/pages/instance/dlq.tsx
index 8f3cdf33..0f02bddf 100644
--- a/web/src/pages/instance/dlq.tsx
+++ b/web/src/pages/instance/dlq.tsx
@@ -181,14 +181,20 @@ const DLQPage = () => {
setRetrySubmitting(true);
setRetryError(null);
try {
- await resendDLQ({
+ const result = await resendDLQ({
groupName: retryGroup.groupName,
startTime: retryRange[0].valueOf(),
endTime: retryRange[1].valueOf(),
targetTopic: retryTargetTopic,
});
setRefreshKey((key) => key + 1);
- message.success(`已提交重投任务:${retryGroup.groupName} → ${retryTargetTopic}`);
+ if (result.failed > 0) {
+ message.warning(`重投部分完成:成功 ${result.resent},失败 ${result.failed}`);
+ } else {
+ message.success(
+ `重投完成:${retryGroup.groupName} → ${retryTargetTopic}(${result.resent}
条)`,
+ );
+ }
setRetryModalOpen(false);
setRetryGroup(null);
setRetryError(null);
diff --git a/web/src/services/messageService.ts
b/web/src/services/messageService.ts
index 80decb21..d4abcf84 100644
--- a/web/src/services/messageService.ts
+++ b/web/src/services/messageService.ts
@@ -1,7 +1,13 @@
import { isMockMode } from './dataMode';
import * as messageApi from '../api/message';
import { sortMessagesByStoreTimeDesc } from '../api/message';
-import type { MessageQuery, MessageRecord, TraceRecord, DLQGroup } from
'../api/message';
+import type {
+ MessageQuery,
+ MessageRecord,
+ TraceRecord,
+ DLQGroup,
+ DLQResendResult,
+} from '../api/message';
import { mockMessages, mockMessageTraces } from '../mock/messages';
import { mockDLQGroups } from '../mock/dlq';
@@ -60,7 +66,7 @@ export async function resendDLQ(data: {
startTime: number;
endTime: number;
targetTopic?: string;
-}): Promise<void> {
- if (isMockMode()) return;
+}): Promise<DLQResendResult> {
+ if (isMockMode()) return { matched: 0, resent: 0, failed: 0, outcome:
'SUCCESS' };
return messageApi.resendDLQ(data);
}