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 b5a7d6c0b refactor(server): replace short-lived admin clients with 
pooled long-lived instances (#2544)
b5a7d6c0b is described below

commit b5a7d6c0b2263b790730b6ff2b01c25d9f470f02
Author: lizhimins <[email protected]>
AuthorDate: Mon Aug 24 14:58:31 2026 +0800

    refactor(server): replace short-lived admin clients with pooled long-lived 
instances (#2544)
    
    Introduce MqClientPool to cache pull-consumer and producer instances by
    (endpoint, credential) pair, replacing the previous pattern of creating
    and shutting down clients on every request. RuntimeAdminClientResolver
    now exposes executePullConsumer/executeProducer entry points backed by
    the pool. Remove ShortLivedClientName utility which is no longer needed.
    
    All provider implementations (RocketMQAdminClientImpl, 
RocketMQMessageProvider,
    RocketMQDLQProvider) are migrated to use pooled clients. Tests updated
    accordingly.
---
 .../studio/cluster/broker/MqClientPool.java        | 219 +++++++++++
 .../cluster/broker/RuntimeAdminClientResolver.java |  28 ++
 .../provider/apache/RocketMQAdminClientImpl.java   |  32 +-
 .../provider/apache/RocketMQDLQProvider.java       |  68 ++--
 .../provider/apache/RocketMQMessageProvider.java   | 240 +++++++-----
 .../provider/apache/ShortLivedClientName.java      |  29 --
 .../broker/RuntimeAdminClientResolverTest.java     |  19 +-
 .../apache/RocketMQAdminClientImplTest.java        | 216 +++++------
 .../apache/RocketMQClusterProviderTest.java        |   3 +-
 .../provider/apache/RocketMQDLQProviderTest.java   | 421 ++++++++-------------
 .../apache/RocketMQMessageProviderTest.java        | 312 ++++++---------
 .../provider/apache/ShortLivedClientNameTest.java  |  37 --
 12 files changed, 803 insertions(+), 821 deletions(-)

diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqClientPool.java
 
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqClientPool.java
new file mode 100644
index 000000000..d9fcf292f
--- /dev/null
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/MqClientPool.java
@@ -0,0 +1,219 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.rocketmq.studio.cluster.broker;
+
+import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
+import org.apache.rocketmq.client.producer.DefaultMQProducer;
+import org.apache.rocketmq.remoting.RPCHook;
+import org.apache.rocketmq.studio.common.exception.BusinessException;
+
+import jakarta.annotation.PreDestroy;
+import lombok.extern.slf4j.Slf4j;
+import org.springframework.stereotype.Component;
+
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.atomic.AtomicInteger;
+
+/**
+ * Long-lived pool of {@link DefaultMQPullConsumer} and {@link 
DefaultMQProducer} clients, keyed by
+ * normalized NameServer address plus a non-secret authentication identity. 
Mirrors the lifecycle
+ * rules of {@link MqAdminExtFactory}: clients are created lazily, reused 
across requests, and shut
+ * down only on application shutdown or explicit endpoint release. Per-request 
client creation is
+ * forbidden because each {@code start()} registers with NameServer/brokers 
and tears down
+ * connections and threads again on {@code shutdown()}.
+ */
+@Slf4j
+@Component
+public class MqClientPool {
+
+    private static final String PULL_CONSUMER_GROUP = 
"studio-pool-pull-consumer";
+    private static final String PRODUCER_GROUP = "studio-pool-producer";
+    private static final long PRODUCER_SEND_TIMEOUT_MILLIS = 5000L;
+
+    private enum Kind { PULL_CONSUMER, PRODUCER }
+
+    private record ClientKey(String namesrvAddr, String 
authenticationIdentity, Kind kind) {
+    }
+
+    private final Map<ClientKey, Object> cache = new ConcurrentHashMap<>();
+    private final AtomicInteger instanceCounter = new AtomicInteger();
+    private volatile boolean closed = false;
+
+    @FunctionalInterface
+    public interface ClientAction<C, T> {
+        T apply(C client) throws Exception;
+    }
+
+    public <T> T withPullConsumer(String namesrvAddr, RPCHook rpcHook, String 
authenticationIdentity,
+                                  ClientAction<DefaultMQPullConsumer, T> 
action) {
+        return execute(new ClientKey(normalize(namesrvAddr), 
identity(authenticationIdentity), Kind.PULL_CONSUMER),
+                rpcHook, this::createPullConsumer, action);
+    }
+
+    public <T> T withProducer(String namesrvAddr, RPCHook rpcHook, String 
authenticationIdentity,
+                              ClientAction<DefaultMQProducer, T> action) {
+        return execute(new ClientKey(normalize(namesrvAddr), 
identity(authenticationIdentity), Kind.PRODUCER),
+                rpcHook, this::createProducer, action);
+    }
+
+    private <C, T> T execute(ClientKey cacheKey, RPCHook rpcHook,
+                             ClientCreator<C> creator, ClientAction<C, T> 
action) {
+        if (cacheKey.namesrvAddr().isEmpty()) {
+            throw new BusinessException(400, "NameServer address is required");
+        }
+        if (closed) {
+            throw new BusinessException(503, "RocketMQ client pool is shutting 
down");
+        }
+        @SuppressWarnings("unchecked")
+        C client = (C) cache.computeIfAbsent(cacheKey, key -> {
+            // Re-check under the cache lock so a request that passed the 
initial closed check
+            // cannot create a fresh connection while the pool is shutting 
down.
+            if (closed) {
+                throw new BusinessException(503, "RocketMQ client pool is 
shutting down");
+            }
+            return creator.create(key.namesrvAddr(), rpcHook);
+        });
+        try {
+            return action.apply(client);
+        } catch (BusinessException ex) {
+            throw ex;
+        } catch (Exception ex) {
+            log.warn("RocketMQ client action failed against namesrv {}: {}", 
cacheKey.namesrvAddr(), ex.getMessage());
+            throw new BusinessException(502, "RocketMQ client call failed: " + 
rootMessage(ex));
+        }
+    }
+
+    /** Stops and removes pooled clients for an endpoint that is no longer 
referenced by Studio. */
+    public void release(String namesrvAddr) {
+        String normalized = normalize(namesrvAddr);
+        if (normalized.isEmpty()) {
+            return;
+        }
+        cache.entrySet().removeIf(entry -> {
+            if (!entry.getKey().namesrvAddr().equals(normalized)) {
+                return false;
+            }
+            safeShutdown(entry.getValue());
+            return true;
+        });
+        log.info("Released pooled RocketMQ clients for namesrv {}", 
normalized);
+    }
+
+    /** Stops one credential-scoped client without interrupting other 
identities on the endpoint. */
+    public void release(String namesrvAddr, String authenticationIdentity) {
+        String normalized = normalize(namesrvAddr);
+        if (normalized.isEmpty()) {
+            return;
+        }
+        for (Kind kind : Kind.values()) {
+            ClientKey key = new ClientKey(normalized, 
identity(authenticationIdentity), kind);
+            Object client = cache.remove(key);
+            if (client != null) {
+                safeShutdown(client);
+            }
+        }
+        log.info("Released pooled RocketMQ clients for namesrv {} and identity 
{}", normalized,
+                identity(authenticationIdentity));
+    }
+
+    @FunctionalInterface
+    private interface ClientCreator<C> {
+        C create(String namesrvAddr, RPCHook rpcHook);
+    }
+
+    private DefaultMQPullConsumer createPullConsumer(String namesrvAddr, 
RPCHook rpcHook) {
+        DefaultMQPullConsumer consumer = new 
DefaultMQPullConsumer(PULL_CONSUMER_GROUP, rpcHook);
+        consumer.setNamesrvAddr(namesrvAddr);
+        consumer.setInstanceName(buildInstanceName(namesrvAddr));
+        try {
+            consumer.start();
+            log.info("Started pooled RocketMQ pull consumer for namesrv {}", 
namesrvAddr);
+            return consumer;
+        } catch (Exception ex) {
+            safeShutdown(consumer);
+            throw new BusinessException(502,
+                    "Failed to connect NameServer " + namesrvAddr + ": " + 
rootMessage(ex));
+        }
+    }
+
+    private DefaultMQProducer createProducer(String namesrvAddr, RPCHook 
rpcHook) {
+        DefaultMQProducer producer = new DefaultMQProducer(PRODUCER_GROUP, 
rpcHook);
+        producer.setNamesrvAddr(namesrvAddr);
+        producer.setInstanceName(buildInstanceName(namesrvAddr));
+        producer.setSendMsgTimeout((int) PRODUCER_SEND_TIMEOUT_MILLIS);
+        producer.setRetryTimesWhenSendFailed(2);
+        // Clusters without a trace topic drop trace dispatch asynchronously, 
so enabling
+        // it unconditionally lets Studio-sent messages produce traces 
wherever trace is on.
+        producer.setEnableTrace(true);
+        try {
+            producer.start();
+            log.info("Started pooled RocketMQ producer for namesrv {}", 
namesrvAddr);
+            return producer;
+        } catch (Exception ex) {
+            safeShutdown(producer);
+            throw new BusinessException(502,
+                    "Failed to connect NameServer " + namesrvAddr + ": " + 
rootMessage(ex));
+        }
+    }
+
+    private String buildInstanceName(String namesrvAddr) {
+        return "rmq-studio-pool-" + Integer.toHexString(namesrvAddr.hashCode())
+                + "-" + instanceCounter.incrementAndGet();
+    }
+
+    private static String normalize(String namesrvAddr) {
+        if (namesrvAddr == null || namesrvAddr.isBlank()) {
+            return "";
+        }
+        return MqAdminExtFactory.normalizeNamesrvAddr(namesrvAddr);
+    }
+
+    private static String identity(String authenticationIdentity) {
+        return authenticationIdentity == null || 
authenticationIdentity.isBlank()
+                ? "anonymous" : authenticationIdentity.trim();
+    }
+
+    private void safeShutdown(Object client) {
+        try {
+            if (client instanceof DefaultMQPullConsumer pullConsumer) {
+                pullConsumer.shutdown();
+            } else if (client instanceof DefaultMQProducer producer) {
+                producer.shutdown();
+            }
+        } catch (Exception ex) {
+            log.debug("Ignoring client shutdown error: {}", ex.getMessage());
+        }
+    }
+
+    private String rootMessage(Throwable ex) {
+        Throwable cause = ex;
+        while (cause.getCause() != null && cause.getCause() != cause) {
+            cause = cause.getCause();
+        }
+        String message = cause.getMessage();
+        return message == null ? cause.getClass().getSimpleName() : message;
+    }
+
+    @PreDestroy
+    public void shutdown() {
+        closed = true;
+        cache.values().forEach(this::safeShutdown);
+        cache.clear();
+        log.info("Shut down all pooled RocketMQ clients");
+    }
+}
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
 
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
index 37cdf1c75..f689e3acd 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolver.java
@@ -9,6 +9,8 @@ package org.apache.rocketmq.studio.cluster.broker;
 import lombok.RequiredArgsConstructor;
 import org.apache.rocketmq.acl.common.AclClientRPCHook;
 import org.apache.rocketmq.acl.common.SessionCredentials;
+import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
+import org.apache.rocketmq.client.producer.DefaultMQProducer;
 import org.apache.rocketmq.remoting.RPCHook;
 import org.apache.rocketmq.studio.common.domain.enums.InstanceVendor;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
@@ -24,6 +26,7 @@ public class RuntimeAdminClientResolver {
     private final InstanceRepository instanceRepository;
     private final MqAdminExtFactory adminFactory;
     private final MqAdminProperties adminProperties;
+    private final MqClientPool clientPool;
 
     public InstanceVO resolveInstance(String instanceId) {
         if (!StringUtils.hasText(instanceId)) {
@@ -64,6 +67,31 @@ public class RuntimeAdminClientResolver {
                 credentialRef, action);
     }
 
+    /** Runs an action on a pooled, long-lived pull consumer bound to the 
instance endpoint. */
+    public <T> T executePullConsumer(String instanceId,
+                                     
MqClientPool.ClientAction<DefaultMQPullConsumer, T> action) {
+        InstanceVO instance = 
requireApacheInstance(resolveInstance(instanceId));
+        String credentialRef = credentialRef(instance);
+        return clientPool.withPullConsumer(requireEndpoint(instance), 
resolveCredential(credentialRef),
+                credentialRef, action);
+    }
+
+    /** Runs an action on a pooled, long-lived producer bound to the instance 
endpoint. */
+    public <T> T executeProducer(String instanceId,
+                                 MqClientPool.ClientAction<DefaultMQProducer, 
T> action) {
+        InstanceVO instance = 
requireApacheInstance(resolveInstance(instanceId));
+        String credentialRef = credentialRef(instance);
+        return clientPool.withProducer(requireEndpoint(instance), 
resolveCredential(credentialRef),
+                credentialRef, action);
+    }
+
+    private String requireEndpoint(InstanceVO instance) {
+        if (instance == null || !StringUtils.hasText(instance.getEndpoint())) {
+            throw new BusinessException(400, "Instance endpoint is required");
+        }
+        return instance.getEndpoint().trim();
+    }
+
     private String credentialRef(InstanceVO instance) {
         return StringUtils.hasText(instance.getAdminCredentialRef())
                 ? instance.getAdminCredentialRef().trim() : null;
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 8c223a3d8..6448be000 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
@@ -26,12 +26,12 @@ import org.apache.rocketmq.client.producer.SendStatus;
 import org.apache.rocketmq.common.TopicConfig;
 import org.apache.rocketmq.common.TopicAttributes;
 import org.apache.rocketmq.common.message.Message;
-import org.apache.rocketmq.remoting.RPCHook;
 import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
 import org.apache.rocketmq.remoting.protocol.ResponseCode;
 import org.apache.rocketmq.remoting.protocol.route.BrokerData;
 import 
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
 import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.common.domain.enums.TopicPerm;
@@ -71,7 +71,6 @@ import java.util.Set;
 @RequiredArgsConstructor
 public class RocketMQAdminClientImpl implements AdminClient {
 
-    private static final String MESSAGE_SENDER_GROUP_PREFIX = 
"studio-msg-sender";
     private static final int MAX_MESSAGE_SIZE = 4 * 1024 * 1024; // 4 MB 
default broker limit
     private static final String LEGACY_METADATA_SCOPE = "";
 
@@ -81,6 +80,7 @@ public class RocketMQAdminClientImpl implements AdminClient {
     private final RmqGroupMapper groupMapper;
     private final AuditService auditService;
     private final RuntimeAdminClientResolver runtimeAdminClientResolver;
+    private final MqClientPool clientPool;
 
     @org.springframework.beans.factory.annotation.Autowired(required = false)
     private ProxyConsumerResolver proxyConsumerResolver;
@@ -441,16 +441,7 @@ public class RocketMQAdminClientImpl implements 
AdminClient {
             throw new BusinessException(400, message);
         }
 
-        String namesrvAddr = namesrvAddr(request.getInstanceId());
-        RPCHook credentialHook = credentialHook(request.getInstanceId());
-
-        DefaultMQProducer producer = new 
DefaultMQProducer(nextMessageSenderGroup(), credentialHook);
-        producer.setNamesrvAddr(namesrvAddr);
-        producer.setSendMsgTimeout(5000);
-
-        try {
-            producer.start();
-
+        MqClientPool.ClientAction<DefaultMQProducer, SendMessageVO> sendAction 
= producer -> {
             Message msg = new Message(topic, tag, key, bodyBytes);
 
             // Add custom properties
@@ -476,21 +467,20 @@ public class RocketMQAdminClientImpl implements 
AdminClient {
                     .sendTime(System.currentTimeMillis())
                     .offsetMsgId(sendResult.getOffsetMsgId())
                     .build();
+        };
+        try {
+            return StringUtils.hasText(request.getInstanceId())
+                    ? 
runtimeAdminClientResolver.executeProducer(request.getInstanceId(), sendAction)
+                    : clientPool.withProducer(namesrvAddr(), null, null, 
sendAction);
         } catch (BusinessException e) {
             recordAudit("SEND_MESSAGE", request.getTopic(), e.getMessage(), 
"FAILED");
             throw e;
         } catch (Exception e) {
             recordAudit("SEND_MESSAGE", request.getTopic(), e.getMessage(), 
"FAILED");
             throw new BusinessException(500, "Failed to send message: " + 
e.getMessage());
-        } finally {
-            producer.shutdown();
         }
     }
 
-    static String nextMessageSenderGroup() {
-        return ShortLivedClientName.next(MESSAGE_SENDER_GROUP_PREFIX);
-    }
-
     @Override
     public ConsumerGroupVO createConsumerGroup(ConsumerGroupVO group) {
         if (group != null && group.getInstanceId() != null) {
@@ -749,12 +739,6 @@ public class RocketMQAdminClientImpl implements 
AdminClient {
                 : namesrvAddr();
     }
 
-    private RPCHook credentialHook(String instanceId) {
-        return StringUtils.hasText(instanceId)
-                ? runtimeAdminClientResolver.resolveCredentialHook(instanceId)
-                : null;
-    }
-
     private <T> T executeForInstance(String instanceId, 
MqAdminExtFactory.AdminAction<T> action) {
         if (StringUtils.hasText(instanceId)) {
             return runtimeAdminClientResolver.execute(instanceId, action);
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 7c4aca641..2058a77a0 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
@@ -27,7 +27,6 @@ import org.apache.rocketmq.common.message.Message;
 import org.apache.rocketmq.common.message.MessageConst;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.remoting.RPCHook;
 import org.apache.rocketmq.remoting.protocol.admin.TopicOffset;
 import org.apache.rocketmq.remoting.protocol.admin.TopicStatsTable;
 import org.apache.rocketmq.remoting.protocol.body.TopicList;
@@ -174,13 +173,11 @@ public class RocketMQDLQProvider implements DLQProvider {
             throw new BusinessException(400, "DLQ resend start time must be 
before end time");
         }
 
-        String endpoint = 
runtimeAdminClientResolver.resolveEndpoint(instanceId);
-        RPCHook credentialHook = 
runtimeAdminClientResolver.resolveCredentialHook(instanceId);
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + groupName;
 
         DeadLetterScanResult scanResult;
         try {
-            scanResult = collectDeadLetters(endpoint, credentialHook, 
dlqTopic, begin, end, RESEND_HARD_CAP);
+            scanResult = collectDeadLetters(instanceId, dlqTopic, begin, end, 
RESEND_HARD_CAP);
         } 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",
@@ -190,26 +187,26 @@ public class RocketMQDLQProvider implements DLQProvider {
             throw e;
         }
         List<MessageExt> deadLetters = scanResult.messages();
-        int resent = 0;
-        int failed = 0;
+        int[] counts = {0, 0};
         if (!deadLetters.isEmpty()) {
-            DefaultMQProducer producer = newProducer(endpoint, credentialHook);
             try {
-                producer.start();
-                for (MessageExt deadLetter : deadLetters) {
-                    if (resendOne(producer, deadLetter, targetTopic)) {
-                        resent++;
-                    } else {
-                        failed++;
+                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 start resend producer for group {}: {}", 
groupName, e.getMessage());
-                failed += deadLetters.size() - resent;
-            } finally {
-                producer.shutdown();
+                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];
 
         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, "
@@ -232,13 +229,11 @@ public class RocketMQDLQProvider implements DLQProvider {
     @Override
     public List<DLQMessageVO> exportMessages(String instanceId, String 
groupName, Long startTime, Long endTime,
                                              Integer maxCount) {
-        String endpoint = 
runtimeAdminClientResolver.resolveEndpoint(instanceId);
-        RPCHook credentialHook = 
runtimeAdminClientResolver.resolveCredentialHook(instanceId);
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + groupName;
         long end = endTime != null ? endTime : System.currentTimeMillis();
         long begin = startTime != null ? startTime : end - ONE_HOUR_MILLIS;
         int cap = maxCount == null || maxCount <= 0 ? RESEND_HARD_CAP : 
Math.min(maxCount, RESEND_HARD_CAP);
-        DeadLetterScanResult scanResult = collectDeadLetters(endpoint, 
credentialHook, dlqTopic, begin, end, cap);
+        DeadLetterScanResult scanResult = collectDeadLetters(instanceId, 
dlqTopic, begin, end, cap);
         return scanResult.messages().stream().map(this::toExportVO).toList();
     }
 
@@ -267,14 +262,18 @@ public class RocketMQDLQProvider implements DLQProvider {
         }
     }
 
-    private DeadLetterScanResult collectDeadLetters(String endpoint, RPCHook 
credentialHook, String dlqTopic,
+    private DeadLetterScanResult collectDeadLetters(String instanceId, String 
dlqTopic,
                                                      long begin, long end, int 
cap) {
-        DefaultMQPullConsumer consumer = newPullConsumer(endpoint, 
credentialHook);
+        return runtimeAdminClientResolver.executePullConsumer(instanceId,
+                consumer -> scanDeadLetters(consumer, dlqTopic, begin, end, 
cap));
+    }
+
+    private DeadLetterScanResult scanDeadLetters(DefaultMQPullConsumer 
consumer, String dlqTopic,
+                                                  long begin, long end, int 
cap) {
         List<MessageExt> result = new ArrayList<>();
         int failedQueueCount = 0;
         boolean truncated = false;
         try {
-            consumer.start();
             Set<MessageQueue> queues = 
consumer.fetchSubscribeMessageQueues(dlqTopic);
             if (queues == null || queues.isEmpty()) {
                 return new DeadLetterScanResult(result, 0, false);
@@ -352,8 +351,6 @@ public class RocketMQDLQProvider implements DLQProvider {
             }
             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 new DeadLetterScanResult(result, failedQueueCount, truncated);
     }
@@ -423,25 +420,6 @@ public class RocketMQDLQProvider implements DLQProvider {
         return null;
     }
 
-    private DefaultMQPullConsumer newPullConsumer(String endpoint, RPCHook 
credentialHook) {
-        DefaultMQPullConsumer consumer = new 
DefaultMQPullConsumer("studio-dlq-query-group", credentialHook);
-        
consumer.setInstanceName(ShortLivedClientName.next("studio-dlq-query"));
-        consumer.setNamesrvAddr(endpoint);
-        return consumer;
-    }
-
-    private DefaultMQProducer newProducer(String endpoint, RPCHook 
credentialHook) {
-        DefaultMQProducer producer = new 
DefaultMQProducer(nextResendProducerGroup(), credentialHook);
-        
producer.setInstanceName(ShortLivedClientName.next("studio-dlq-resend"));
-        producer.setRetryTimesWhenSendFailed(2);
-        producer.setNamesrvAddr(endpoint);
-        return producer;
-    }
-
-    static String nextResendProducerGroup() {
-        return ShortLivedClientName.next("studio-dlq-resend");
-    }
-
     private String classifyOutcome(int matched, int resent, int failed, 
boolean scanIncomplete) {
         if (scanIncomplete) {
             return "PARTIAL";
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
index fc324929b..b1b3c00de 100644
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
+++ 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProvider.java
@@ -17,6 +17,7 @@
 package org.apache.rocketmq.studio.provider.apache;
 
 import org.apache.rocketmq.client.QueryResult;
+import org.apache.rocketmq.client.exception.MQClientException;
 import org.apache.rocketmq.client.consumer.DefaultMQPullConsumer;
 import org.apache.rocketmq.client.consumer.PullResult;
 import org.apache.rocketmq.client.consumer.PullStatus;
@@ -25,15 +26,18 @@ import org.apache.rocketmq.common.message.MessageDecoder;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageId;
 import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.remoting.RPCHook;
+import org.apache.rocketmq.remoting.protocol.ResponseCode;
 import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
 import org.apache.rocketmq.remoting.protocol.route.BrokerData;
+import org.apache.rocketmq.remoting.protocol.route.QueueData;
+import org.apache.rocketmq.remoting.protocol.route.TopicRouteData;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
 import org.apache.rocketmq.studio.instance.message.ConsumerStatusVO;
 import org.apache.rocketmq.studio.instance.message.MessageProvider;
 import org.apache.rocketmq.studio.instance.message.MessageRecordVO;
+import org.apache.rocketmq.studio.instance.message.QueueOffsetVO;
 import org.apache.rocketmq.studio.instance.message.TraceNodeVO;
 import org.apache.rocketmq.studio.instance.message.TraceRecordVO;
 import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
@@ -87,7 +91,7 @@ public class RocketMQMessageProvider implements 
MessageProvider {
     private static final long ONE_HOUR_MILLIS = 3600_000L;
     private static final long ONE_DAY_MILLIS = 24 * ONE_HOUR_MILLIS;
     private static final long MAX_TOPIC_QUERY_WINDOW_MILLIS = 7 * 
ONE_DAY_MILLIS;
-    private static final int MAX_PULLS_PER_QUEUE = 1_000;
+    private static final int MAX_PULLS_PER_QUEUE = 32;
     private static final int MAX_CONSECUTIVE_OFFSET_ILLEGAL = 3;
     private static final int MAX_TOPIC_SCAN_MESSAGES_PER_QUEUE = 
MAX_PULLS_PER_QUEUE * TOPIC_PULL_BATCH_SIZE;
     private static final int MAX_PULL_ATTEMPTS_PER_QUEUE = MAX_PULLS_PER_QUEUE 
+ MAX_CONSECUTIVE_OFFSET_ILLEGAL;
@@ -100,15 +104,12 @@ public class RocketMQMessageProvider implements 
MessageProvider {
     @Override
     public List<MessageRecordVO> queryMessages(String instanceId, String 
topic, String msgId, String tag, String key,
                                                Long startTime, Long endTime) {
-        String endpoint = 
runtimeAdminClientResolver.resolveEndpoint(instanceId);
-        RPCHook credentialHook = 
runtimeAdminClientResolver.resolveCredentialHook(instanceId);
         return runtimeAdminClientResolver.execute(instanceId,
-                adminExt -> queryMessages(instanceId, (DefaultMQAdminExt) 
adminExt, endpoint,
-                        credentialHook, topic, msgId, tag, key, startTime, 
endTime));
+                adminExt -> queryMessages(instanceId, (DefaultMQAdminExt) 
adminExt, topic, msgId, tag, key,
+                        startTime, endTime));
     }
 
-    private List<MessageRecordVO> queryMessages(String instanceId, 
DefaultMQAdminExt adminExt, String endpoint,
-                                                 RPCHook credentialHook,
+    private List<MessageRecordVO> queryMessages(String instanceId, 
DefaultMQAdminExt adminExt,
                                                  String topic, String msgId, 
String tag, String key,
                                                  Long startTime, Long endTime) 
{
 
@@ -127,7 +128,7 @@ public class RocketMQMessageProvider implements 
MessageProvider {
             if (begin >= 0 && end >= 0 && end - begin > 
MAX_TOPIC_QUERY_WINDOW_MILLIS) {
                 throw new BusinessException(400, "Topic message query time 
range must not exceed 7 days");
             }
-            result = queryByTopic(endpoint, credentialHook, topic, tag, begin, 
end, DEFAULT_TOPIC_LIMIT);
+            result = queryByTopic(instanceId, topic, tag, begin, end, 
DEFAULT_TOPIC_LIMIT);
         } else {
             log.warn("queryMessages requires at least one of msgId/topic, 
returning empty list");
             return Collections.emptyList();
@@ -240,91 +241,139 @@ public class RocketMQMessageProvider implements 
MessageProvider {
         }
     }
 
+    @Override
+    public List<QueueOffsetVO> getQueueOffsets(String instanceId, String 
topic) {
+        return runtimeAdminClientResolver.execute(instanceId, adminExt -> {
+            List<QueueOffsetVO> result = new ArrayList<>();
+            try {
+                TopicRouteData route = adminExt.examineTopicRouteInfo(topic);
+                if (route == null || route.getQueueDatas() == null) {
+                    return Collections.emptyList();
+                }
+                for (QueueData queueData : route.getQueueDatas()) {
+                    for (int queueId = 0; queueId < 
queueData.getWriteQueueNums(); queueId++) {
+                        MessageQueue queue = new MessageQueue(topic, 
queueData.getBrokerName(), queueId);
+                        result.add(QueueOffsetVO.builder()
+                                .brokerName(queue.getBrokerName())
+                                .queueId(queue.getQueueId())
+                                .minOffset(adminExt.minOffset(queue))
+                                .maxOffset(adminExt.maxOffset(queue))
+                                .build());
+                    }
+                }
+                result.sort(Comparator.comparing(QueueOffsetVO::getBrokerName)
+                        .thenComparingInt(QueueOffsetVO::getQueueId));
+            } catch (Exception e) {
+                log.warn("getQueueOffsets(topic={}) failed: {}", topic, 
e.getMessage());
+                throw new BusinessException(502, "Failed to get queue offsets: 
" + e.getMessage());
+            }
+            return result;
+        });
+    }
+
+    @Override
+    public MessageRecordVO pullMessageAtOffset(String instanceId, String 
topic, String brokerName,
+                                                int queueId, long offset) {
+        return runtimeAdminClientResolver.executePullConsumer(instanceId, 
consumer -> {
+            try {
+                MessageQueue queue = new MessageQueue(topic, brokerName, 
queueId);
+                PullResult pullResult = consumer.pull(queue, "*", offset, 1);
+                if (pullResult == null || pullResult.getPullStatus() != 
PullStatus.FOUND
+                        || pullResult.getMsgFoundList() == null || 
pullResult.getMsgFoundList().isEmpty()) {
+                    return null;
+                }
+                return toRecordVO(pullResult.getMsgFoundList().get(0), 
brokerName);
+            } catch (Exception e) {
+                log.warn("pullMessageAtOffset(topic={}, broker={}, queue={}, 
offset={}) failed: {}",
+                        topic, brokerName, queueId, offset, e.getMessage());
+                throw new BusinessException(502, "Failed to pull message at 
offset: " + e.getMessage());
+            }
+        });
+    }
+
     /**
-     * Scan a topic within a time range using a short-lived pull consumer, 
mirroring the approach
-     * used by the RocketMQ dashboard for time-range topic queries.
+     * Scan a topic within a time range on the pooled long-lived pull 
consumer, mirroring the
+     * approach used by the RocketMQ dashboard for time-range topic queries.
      */
-    private List<MessageRecordVO> queryByTopic(String endpoint, RPCHook 
credentialHook, String topic, String tag,
+    private List<MessageRecordVO> queryByTopic(String instanceId, String 
topic, String tag,
                                                 long begin, long end, int 
limit) {
-        DefaultMQPullConsumer consumer = newPullConsumer("studio-msg-query", 
endpoint, credentialHook);
         int resultLimit = Math.min(limit, TOPIC_QUERY_HARD_CAP);
-        PriorityQueue<MessageRecordVO> newestMessages = new 
PriorityQueue<>(TOPIC_QUERY_ORDER);
-        try {
-            consumer.start();
-            Set<MessageQueue> queues = 
consumer.fetchSubscribeMessageQueues(topic);
-            if (queues == null || queues.isEmpty()) {
-                return Collections.emptyList();
-            }
-            for (MessageQueue queue : queues) {
-                TopicQueueScanPlan scanPlan = 
buildTopicQueueScanPlan(consumer, queue, begin, end);
-                if (scanPlan.isEmpty()) {
-                    continue;
-                }
-                if (scanPlan.truncated()) {
-                    log.info("Truncate topic query for {} queue {} to offsets 
[{}..{}) within the guarded tail budget",
-                            topic, queue, scanPlan.startOffset(), 
scanPlan.endOffsetExclusive());
+        return runtimeAdminClientResolver.executePullConsumer(instanceId, 
consumer -> {
+            PriorityQueue<MessageRecordVO> newestMessages = new 
PriorityQueue<>(TOPIC_QUERY_ORDER);
+            try {
+                Set<MessageQueue> queues = 
consumer.fetchSubscribeMessageQueues(topic);
+                if (queues == null || queues.isEmpty()) {
+                    return Collections.emptyList();
                 }
-                int consecutiveIllegalOffsets = 0;
-                int pullAttempts = 0;
-                for (long offset = scanPlan.startOffset(); offset < 
scanPlan.endOffsetExclusive(); ) {
-                    if (++pullAttempts > MAX_PULL_ATTEMPTS_PER_QUEUE) {
-                        log.warn("Stop topic query for {} because queue {} 
exhausted the guarded pull budget at offset {}",
-                                topic, queue, offset);
-                        break;
-                    }
-                    PullResult pullResult = consumer.pull(queue, "*", offset, 
TOPIC_PULL_BATCH_SIZE);
-                    if (pullResult == null) {
-                        log.warn("Stop topic query for {} because queue {} 
returned no pull result", topic, queue);
-                        break;
+                for (MessageQueue queue : queues) {
+                    TopicQueueScanPlan scanPlan = 
buildTopicQueueScanPlan(consumer, queue, begin, end);
+                    if (scanPlan.isEmpty()) {
+                        continue;
                     }
-                    long nextOffset = pullResult.getNextBeginOffset();
-                    if (nextOffset <= offset) {
-                        log.warn("Stop topic query for {} because queue {} did 
not advance offset {}", topic, queue, offset);
-                        break;
+                    if (scanPlan.truncated()) {
+                        log.info("Truncate topic query for {} queue {} to 
offsets [{}..{}) within the guarded tail budget",
+                                topic, queue, scanPlan.startOffset(), 
scanPlan.endOffsetExclusive());
                     }
-                    offset = Math.min(nextOffset, 
scanPlan.endOffsetExclusive());
-                    if (pullResult.getPullStatus() == 
PullStatus.OFFSET_ILLEGAL) {
-                        // The broker returned a corrected offset in 
nextBeginOffset because
-                        // the requested offset is no longer valid (expired, 
compacted, or
-                        // before the queue's minimum offset). Retry from the 
corrected
-                        // position instead of abandoning the queue -- 
otherwise messages
-                        // that still exist after the corrected offset are 
silently dropped.
-                        if (++consecutiveIllegalOffsets > 
MAX_CONSECUTIVE_OFFSET_ILLEGAL) {
-                            log.warn("Stop topic query for {} because queue {} 
returned OFFSET_ILLEGAL "
-                                    + "{} times consecutively, giving up at 
offset {}", topic, queue,
-                                    consecutiveIllegalOffsets, offset);
+                    int consecutiveIllegalOffsets = 0;
+                    int pullAttempts = 0;
+                    for (long offset = scanPlan.startOffset(); offset < 
scanPlan.endOffsetExclusive(); ) {
+                        if (++pullAttempts > MAX_PULL_ATTEMPTS_PER_QUEUE) {
+                            log.warn("Stop topic query for {} because queue {} 
exhausted the guarded pull budget at offset {}",
+                                    topic, queue, offset);
                             break;
                         }
-                        log.debug("Offset was illegal for queue {} in topic 
{}, retrying from {}",
-                                queue, topic, offset);
-                        continue;
-                    }
-                    if (pullResult.getPullStatus() != PullStatus.FOUND
-                            || pullResult.getMsgFoundList() == null) {
-                        break;
-                    }
-                    consecutiveIllegalOffsets = 0;
-                    for (MessageExt messageExt : pullResult.getMsgFoundList()) 
{
-                        if (messageExt.getStoreTimestamp() < begin
-                                || messageExt.getStoreTimestamp() > end) {
-                            continue;
+                        PullResult pullResult = consumer.pull(queue, "*", 
offset, TOPIC_PULL_BATCH_SIZE);
+                        if (pullResult == null) {
+                            log.warn("Stop topic query for {} because queue {} 
returned no pull result", topic, queue);
+                            break;
+                        }
+                        long nextOffset = pullResult.getNextBeginOffset();
+                        if (nextOffset <= offset) {
+                            log.warn("Stop topic query for {} because queue {} 
did not advance offset {}", topic, queue, offset);
+                            break;
                         }
-                        if (!matchesTag(messageExt, tag)) {
+                        offset = Math.min(nextOffset, 
scanPlan.endOffsetExclusive());
+                        if (pullResult.getPullStatus() == 
PullStatus.OFFSET_ILLEGAL) {
+                            // The broker returned a corrected offset in 
nextBeginOffset because
+                            // the requested offset is no longer valid 
(expired, compacted, or
+                            // before the queue's minimum offset). Retry from 
the corrected
+                            // position instead of abandoning the queue -- 
otherwise messages
+                            // that still exist after the corrected offset are 
silently dropped.
+                            if (++consecutiveIllegalOffsets > 
MAX_CONSECUTIVE_OFFSET_ILLEGAL) {
+                                log.warn("Stop topic query for {} because 
queue {} returned OFFSET_ILLEGAL "
+                                        + "{} times consecutively, giving up 
at offset {}", topic, queue,
+                                        consecutiveIllegalOffsets, offset);
+                                break;
+                            }
+                            log.debug("Offset was illegal for queue {} in 
topic {}, retrying from {}",
+                                    queue, topic, offset);
                             continue;
                         }
-                        addTopicQueryCandidate(newestMessages, 
toRecordVO(messageExt), resultLimit);
+                        if (pullResult.getPullStatus() != PullStatus.FOUND
+                                || pullResult.getMsgFoundList() == null) {
+                            break;
+                        }
+                        consecutiveIllegalOffsets = 0;
+                        for (MessageExt messageExt : 
pullResult.getMsgFoundList()) {
+                            if (messageExt.getStoreTimestamp() < begin
+                                    || messageExt.getStoreTimestamp() > end) {
+                                continue;
+                            }
+                            if (!matchesTag(messageExt, tag)) {
+                                continue;
+                            }
+                            addTopicQueryCandidate(newestMessages, 
toRecordVO(messageExt, queue.getBrokerName()), resultLimit);
+                        }
                     }
                 }
+            } catch (Exception e) {
+                log.warn("queryByTopic(topic={}) failed: {}", topic, 
e.getMessage());
+                throw new BusinessException(502, "Failed to query messages by 
topic: " + e.getMessage());
             }
-        } catch (Exception e) {
-            log.warn("queryByTopic(topic={}) failed: {}", topic, 
e.getMessage());
-            throw new BusinessException(502, "Failed to query messages by 
topic: " + e.getMessage());
-        } finally {
-            consumer.shutdown();
-        }
-        return newestMessages.stream()
-                .sorted(TOPIC_QUERY_ORDER.reversed())
-                .toList();
+            return newestMessages.stream()
+                    .sorted(TOPIC_QUERY_ORDER.reversed())
+                    .toList();
+        });
     }
 
     private TopicQueueScanPlan buildTopicQueueScanPlan(DefaultMQPullConsumer 
consumer, MessageQueue queue,
@@ -410,6 +459,13 @@ public class RocketMQMessageProvider implements 
MessageProvider {
         } catch (BusinessException e) {
             throw e;
         } catch (Exception e) {
+            if (isTraceTopicAbsent(e)) {
+                // The cluster has no trace topic route (trace dispatch 
disabled): the RPC
+                // succeeded but there is no business data, so return an empty 
trace instead
+                // of surfacing an error (exception-grading convention).
+                log.info("Trace topic not available on this cluster 
(msgId={}), returning empty trace", msgId);
+                return emptyTrace();
+            }
             log.warn("Trace query for msgId={} failed: {}", msgId, 
e.getMessage());
             throw new BusinessException(502, "Failed to query message trace: " 
+ e.getMessage());
         }
@@ -420,6 +476,18 @@ public class RocketMQMessageProvider implements 
MessageProvider {
                 .build();
     }
 
+    private static boolean isTraceTopicAbsent(Throwable error) {
+        Throwable cause = error;
+        while (cause != null) {
+            if (cause instanceof MQClientException clientException
+                    && clientException.getResponseCode() == 
ResponseCode.TOPIC_NOT_EXIST) {
+                return true;
+            }
+            cause = cause.getCause() == cause ? null : cause.getCause();
+        }
+        return false;
+    }
+
     /**
      * Attempts to resolve the store timestamp of the original message so the 
trace
      * query window can be derived from the message's own timeline rather than 
the
@@ -601,6 +669,10 @@ public class RocketMQMessageProvider implements 
MessageProvider {
     }
 
     MessageRecordVO toRecordVO(MessageExt messageExt) {
+        return toRecordVO(messageExt, null);
+    }
+
+    MessageRecordVO toRecordVO(MessageExt messageExt, String brokerName) {
         byte[] body = messageExt.getBody();
         DisplayBody displayBody = displayBody(body);
         Map<String, String> properties = messageExt.getProperties();
@@ -610,6 +682,9 @@ public class RocketMQMessageProvider implements 
MessageProvider {
                 .topic(messageExt.getTopic())
                 .tag(messageExt.getTags())
                 .key(messageExt.getKeys())
+                .brokerName(brokerName)
+                .queueId(messageExt.getQueueId())
+                .queueOffset(messageExt.getQueueOffset())
                 .body(displayBody.value())
                 .bodyEncoding(displayBody.encoding())
                 .bodyTruncated(displayBody.truncated())
@@ -685,13 +760,6 @@ public class RocketMQMessageProvider implements 
MessageProvider {
         return tag.equals(messageExt.getTags());
     }
 
-    private DefaultMQPullConsumer newPullConsumer(String groupPrefix, String 
endpoint, RPCHook credentialHook) {
-        DefaultMQPullConsumer consumer = new DefaultMQPullConsumer(groupPrefix 
+ "-group", credentialHook);
-        consumer.setInstanceName(ShortLivedClientName.next(groupPrefix));
-        consumer.setNamesrvAddr(endpoint);
-        return consumer;
-    }
-
     private static TraceRecordVO emptyTrace() {
         return TraceRecordVO.builder()
                 .nodes(Collections.emptyList())
diff --git 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientName.java
 
b/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientName.java
deleted file mode 100644
index 3259856dc..000000000
--- 
a/server/src/main/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientName.java
+++ /dev/null
@@ -1,29 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements.  See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package org.apache.rocketmq.studio.provider.apache;
-
-import java.util.UUID;
-
-final class ShortLivedClientName {
-
-    private ShortLivedClientName() {
-    }
-
-    static String next(String prefix) {
-        return prefix + "-" + UUID.randomUUID();
-    }
-}
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
index f54831898..f6b18b4de 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/cluster/broker/RuntimeAdminClientResolverTest.java
@@ -50,6 +50,9 @@ class RuntimeAdminClientResolverTest {
     @Mock
     private MqAdminExtFactory adminFactory;
 
+    @Mock
+    private MqClientPool clientPool;
+
     @Test
     void resolvesTrimmedEndpointFromSelectedInstance() {
         InstanceVO instance = InstanceVO.builder().endpoint(" namesrv-a:9876 
").build();
@@ -57,7 +60,7 @@ class RuntimeAdminClientResolverTest {
         
when(instanceRepository.findByIdentifier("instance-a")).thenReturn(Optional.of(instance));
 
         RuntimeAdminClientResolver resolver = new 
RuntimeAdminClientResolver(instanceRepository, adminFactory,
-                new MqAdminProperties());
+                new MqAdminProperties(), clientPool);
 
         
assertThat(resolver.resolveEndpoint("instance-a")).isEqualTo("namesrv-a:9876");
     }
@@ -65,7 +68,7 @@ class RuntimeAdminClientResolverTest {
     @Test
     void rejectsUnknownOrUnconfiguredInstances() {
         RuntimeAdminClientResolver resolver = new 
RuntimeAdminClientResolver(instanceRepository, adminFactory,
-                new MqAdminProperties());
+                new MqAdminProperties(), clientPool);
         
when(instanceRepository.findByIdentifier("missing")).thenReturn(Optional.empty());
         InstanceVO noEndpoint = InstanceVO.builder().endpoint(" ").build();
         
when(instanceRepository.findByIdentifier("no-endpoint")).thenReturn(Optional.of(noEndpoint));
@@ -84,7 +87,7 @@ class RuntimeAdminClientResolverTest {
         
when(instanceRepository.findByIdentifier("instance-b")).thenReturn(Optional.of(instance));
         when(adminFactory.execute(eq("namesrv-b:9876"), isNull(), isNull(), 
any())).thenReturn("done");
         RuntimeAdminClientResolver resolver = new 
RuntimeAdminClientResolver(instanceRepository, adminFactory,
-                new MqAdminProperties());
+                new MqAdminProperties(), clientPool);
 
         String result = resolver.execute("instance-b", admin -> "unused");
         assertThat(result).isEqualTo("done");
@@ -100,7 +103,7 @@ class RuntimeAdminClientResolverTest {
         instance.setId(2L);
         
when(instanceRepository.findByIdentifier("cloud-instance")).thenReturn(Optional.of(instance));
         RuntimeAdminClientResolver resolver = new 
RuntimeAdminClientResolver(instanceRepository, adminFactory,
-                new MqAdminProperties());
+                new MqAdminProperties(), clientPool);
 
         assertThatThrownBy(() -> resolver.resolveEndpoint("cloud-instance"))
                 .isInstanceOf(BusinessException.class)
@@ -127,7 +130,7 @@ class RuntimeAdminClientResolverTest {
         when(adminFactory.execute(eq("namesrv-b:9876"), any(), 
eq("production-admin"), any()))
                 .thenReturn("done");
         RuntimeAdminClientResolver resolver = new 
RuntimeAdminClientResolver(instanceRepository, adminFactory,
-                properties);
+                properties, clientPool);
 
         String result = resolver.execute("instance-b", ignored -> "unused");
 
@@ -155,7 +158,7 @@ class RuntimeAdminClientResolverTest {
         properties.getCredentials().put("production-admin", credential);
         
when(instanceRepository.findByIdentifier("instance-b")).thenReturn(Optional.of(instance));
         RuntimeAdminClientResolver resolver = new 
RuntimeAdminClientResolver(instanceRepository, adminFactory,
-                properties);
+                properties, clientPool);
 
         org.apache.rocketmq.acl.common.AclClientRPCHook hook =
                 (org.apache.rocketmq.acl.common.AclClientRPCHook) 
resolver.resolveCredentialHook("instance-b");
@@ -171,7 +174,7 @@ class RuntimeAdminClientResolverTest {
         instance.setId(5L);
         
when(instanceRepository.findByIdentifier("instance-b")).thenReturn(Optional.of(instance));
         RuntimeAdminClientResolver resolver = new 
RuntimeAdminClientResolver(instanceRepository, adminFactory,
-                new MqAdminProperties());
+                new MqAdminProperties(), clientPool);
 
         assertThat(resolver.resolveCredentialHook("instance-b")).isNull();
         verifyNoInteractions(adminFactory);
@@ -192,7 +195,7 @@ class RuntimeAdminClientResolverTest {
         instance.setId(3L);
         
when(instanceRepository.findByIdentifier("instance-b")).thenReturn(Optional.of(instance));
         RuntimeAdminClientResolver resolver = new 
RuntimeAdminClientResolver(instanceRepository, adminFactory,
-                properties);
+                properties, clientPool);
 
         assertThatThrownBy(() -> resolver.execute("instance-b", ignored -> 
"unused"))
                 .isInstanceOf(BusinessException.class)
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 fba406cac..df80091c7 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
@@ -22,13 +22,13 @@ import 
org.apache.rocketmq.client.producer.DefaultMQProducer;
 import org.apache.rocketmq.client.producer.SendResult;
 import org.apache.rocketmq.client.producer.SendStatus;
 import org.apache.rocketmq.common.message.Message;
-import org.apache.rocketmq.remoting.RPCHook;
 import org.apache.rocketmq.remoting.exception.RemotingTimeoutException;
 import org.apache.rocketmq.remoting.protocol.ResponseCode;
 import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
 import org.apache.rocketmq.remoting.protocol.route.BrokerData;
 import 
org.apache.rocketmq.remoting.protocol.subscription.SubscriptionGroupConfig;
 import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.common.domain.enums.TopicType;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
@@ -48,11 +48,9 @@ import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.mockito.ArgumentCaptor;
-import org.mockito.MockedConstruction;
 import org.mockito.Mock;
 import org.mockito.junit.jupiter.MockitoExtension;
 
-import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.HashSet;
 import java.util.List;
@@ -65,13 +63,13 @@ 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.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.isNull;
 import static org.mockito.Mockito.doNothing;
 import static org.mockito.Mockito.lenient;
-import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.doThrow;
-import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verifyNoInteractions;
 import static org.mockito.Mockito.when;
@@ -93,6 +91,10 @@ class RocketMQAdminClientImplTest {
     private AuditService auditService;
     @Mock
     private RuntimeAdminClientResolver runtimeAdminClientResolver;
+    @Mock
+    private MqClientPool clientPool;
+    @Mock
+    private DefaultMQProducer sendProducer;
 
     private RocketMQAdminClientImpl adminClient;
 
@@ -101,8 +103,12 @@ class RocketMQAdminClientImplTest {
         
lenient().when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
         lenient().when(adminFactory.execute(anyString(), any(), 
any())).thenAnswer(invocation ->
                 
invocation.<MqAdminExtFactory.AdminAction<Object>>getArgument(2).apply(adminExt));
+        lenient().when(clientPool.withProducer(any(), any(), any(), 
any())).thenAnswer(invocation ->
+                invocation.<MqClientPool.ClientAction<DefaultMQProducer, 
Object>>getArgument(3).apply(sendProducer));
+        lenient().when(runtimeAdminClientResolver.executeProducer(any(), 
any())).thenAnswer(invocation ->
+                invocation.<MqClientPool.ClientAction<DefaultMQProducer, 
Object>>getArgument(1).apply(sendProducer));
         adminClient = new RocketMQAdminClientImpl(adminFactory, properties, 
topicMapper, groupMapper, auditService,
-                runtimeAdminClientResolver);
+                runtimeAdminClientResolver, clientPool);
     }
 
     @Test
@@ -729,26 +735,20 @@ class RocketMQAdminClientImplTest {
 
     @Test
     void sendMessageShouldNotFailWhenAuditRecordingFails() throws Exception {
-        when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
         doThrow(new RuntimeException("audit db down")).when(auditService)
                 .record(anyString(), anyString(), anyString(), any(), 
anyString(), anyString());
-        try (MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         doNothing().when(producer).start();
-                         SendResult sendResult = new SendResult();
-                         sendResult.setSendStatus(SendStatus.SEND_OK);
-                         sendResult.setMsgId("msg-1");
-                         sendResult.setOffsetMsgId("offset-1");
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            SendMessageDTO request = new SendMessageDTO();
-            request.setTopic("TopicA");
-            request.setBody("hello");
-            SendMessageVO result = adminClient.sendMessage(request);
-            // The message was already delivered; an audit failure must not 
turn this into an error.
-            assertThat(result.getMsgId()).isEqualTo("msg-1");
-        }
+        SendResult sendResult = new SendResult();
+        sendResult.setSendStatus(SendStatus.SEND_OK);
+        sendResult.setMsgId("msg-1");
+        sendResult.setOffsetMsgId("offset-1");
+        when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+        SendMessageDTO request = new SendMessageDTO();
+        request.setTopic("TopicA");
+        request.setBody("hello");
+        SendMessageVO result = adminClient.sendMessage(request);
+        // The message was already delivered; an audit failure must not turn 
this into an error.
+        assertThat(result.getMsgId()).isEqualTo("msg-1");
     }
 
     @Test
@@ -758,16 +758,13 @@ class RocketMQAdminClientImplTest {
         request.setTopic("TopicA");
         request.setBody("\u754c".repeat((4 * 1024 * 1024 / 3) + 1));
 
-        try (MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class)) {
-            assertThatThrownBy(() -> adminClient.sendMessage(request))
-                    .isInstanceOf(BusinessException.class)
-                    .hasMessageContaining("exceeds the maximum")
-                    .satisfies(exception -> assertThat(((BusinessException) 
exception).getCode()).isEqualTo(400));
+        assertThatThrownBy(() -> adminClient.sendMessage(request))
+                .isInstanceOf(BusinessException.class)
+                .hasMessageContaining("exceeds the maximum")
+                .satisfies(exception -> assertThat(((BusinessException) 
exception).getCode()).isEqualTo(400));
 
-            assertThat(mockedProducers.constructed()).isEmpty();
-        }
         verifyNoInteractions(runtimeAdminClientResolver);
+        verifyNoInteractions(clientPool);
         verify(properties, never()).getNamesrvAddr();
         verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", null,
                 "Message body size 4194306 exceeds the maximum of 4194304 
bytes", "FAILED");
@@ -775,112 +772,75 @@ class RocketMQAdminClientImplTest {
 
     @Test
     void sendMessageShouldAllowBodyAtMaximumSize() throws Exception {
-        when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
-        List<List<?>> constructorArguments = new ArrayList<>();
-        try (MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         constructorArguments.add(context.arguments());
-                         doNothing().when(producer).start();
-                         SendResult sendResult = new SendResult();
-                         sendResult.setSendStatus(SendStatus.SEND_OK);
-                         sendResult.setMsgId("msg-1");
-                         sendResult.setOffsetMsgId("offset-1");
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            SendMessageDTO request = new SendMessageDTO();
-            request.setTopic("TopicA");
-            request.setBody("x".repeat(4 * 1024 * 1024));
-
-            SendMessageVO result = adminClient.sendMessage(request);
-
-            assertThat(result.getMsgId()).isEqualTo("msg-1");
-            DefaultMQProducer producer = 
mockedProducers.constructed().getFirst();
-            ArgumentCaptor<Message> messageCaptor = 
ArgumentCaptor.forClass(Message.class);
-            verify(producer).send(messageCaptor.capture());
-            assertThat(messageCaptor.getValue().getBody()).hasSize(4 * 1024 * 
1024);
-            assertThat(constructorArguments).singleElement();
-            assertThat(constructorArguments.get(0)).hasSize(2);
-            assertThat(constructorArguments.get(0).get(1)).isNull();
-        }
+        SendResult sendResult = new SendResult();
+        sendResult.setSendStatus(SendStatus.SEND_OK);
+        sendResult.setMsgId("msg-1");
+        sendResult.setOffsetMsgId("offset-1");
+        when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+        SendMessageDTO request = new SendMessageDTO();
+        request.setTopic("TopicA");
+        request.setBody("x".repeat(4 * 1024 * 1024));
+
+        SendMessageVO result = adminClient.sendMessage(request);
+
+        assertThat(result.getMsgId()).isEqualTo("msg-1");
+        ArgumentCaptor<Message> messageCaptor = 
ArgumentCaptor.forClass(Message.class);
+        verify(sendProducer).send(messageCaptor.capture());
+        assertThat(messageCaptor.getValue().getBody()).hasSize(4 * 1024 * 
1024);
+        verify(clientPool).withProducer(eq("10.0.0.1:9876"), isNull(), 
isNull(), any());
     }
 
     @Test
-    void sendMessageUsesSelectedInstanceEndpointAndCredentialHook() throws 
Exception {
-        RPCHook credentialHook = mock(RPCHook.class);
-        List<List<?>> constructorArguments = new ArrayList<>();
-        
when(runtimeAdminClientResolver.resolveEndpoint("instance-a")).thenReturn("10.0.0.2:9876");
-        
when(runtimeAdminClientResolver.resolveCredentialHook("instance-a")).thenReturn(credentialHook);
-        try (MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         constructorArguments.add(context.arguments());
-                         doNothing().when(producer).start();
-                         SendResult sendResult = new SendResult();
-                         sendResult.setSendStatus(SendStatus.SEND_OK);
-                         sendResult.setMsgId("msg-1");
-                         sendResult.setOffsetMsgId("offset-1");
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            SendMessageDTO request = new SendMessageDTO();
-            request.setTopic("TopicA");
-            request.setBody("hello");
-            request.setInstanceId("instance-a");
-
-            adminClient.sendMessage(request);
-
-            DefaultMQProducer producer = 
mockedProducers.constructed().getFirst();
-            verify(producer).setNamesrvAddr("10.0.0.2:9876");
-            assertThat(constructorArguments).singleElement();
-            assertThat(constructorArguments.get(0)).hasSize(2);
-            
assertThat(constructorArguments.get(0).get(1)).isSameAs(credentialHook);
-        }
-        verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
+    void sendMessageUsesSelectedInstancePooledProducer() throws Exception {
+        SendResult sendResult = new SendResult();
+        sendResult.setSendStatus(SendStatus.SEND_OK);
+        sendResult.setMsgId("msg-1");
+        sendResult.setOffsetMsgId("offset-1");
+        when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+        SendMessageDTO request = new SendMessageDTO();
+        request.setTopic("TopicA");
+        request.setBody("hello");
+        request.setInstanceId("instance-a");
+
+        adminClient.sendMessage(request);
+
+        verify(runtimeAdminClientResolver).executeProducer(eq("instance-a"), 
any());
+        verify(clientPool, never()).withProducer(any(), any(), any(), any());
     }
 
     @Test
     void sendMessageShouldRejectNonSuccessfulSendStatus() throws Exception {
-        when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
-        try (MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         doNothing().when(producer).start();
-                         SendResult sendResult = new SendResult();
-                         
sendResult.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT);
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            SendMessageDTO request = new SendMessageDTO();
-            request.setTopic("TopicA");
-            request.setBody("hello");
-
-            assertThatThrownBy(() -> adminClient.sendMessage(request))
-                    .isInstanceOf(BusinessException.class)
-                    .hasMessageContaining("FLUSH_DISK_TIMEOUT");
-
-            verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", 
null,
-                    "Message send did not succeed: FLUSH_DISK_TIMEOUT", 
"FAILED");
-        }
+        SendResult sendResult = new SendResult();
+        sendResult.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT);
+        when(sendProducer.send(any(Message.class))).thenReturn(sendResult);
+
+        SendMessageDTO request = new SendMessageDTO();
+        request.setTopic("TopicA");
+        request.setBody("hello");
+
+        assertThatThrownBy(() -> adminClient.sendMessage(request))
+                .isInstanceOf(BusinessException.class)
+                .hasMessageContaining("FLUSH_DISK_TIMEOUT");
+
+        verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", null,
+                "Message send did not succeed: FLUSH_DISK_TIMEOUT", "FAILED");
     }
 
     @Test
     void sendMessageShouldRejectNullSendResult() throws Exception {
-        when(properties.getNamesrvAddr()).thenReturn("10.0.0.1:9876");
-        try (MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         doNothing().when(producer).start();
-                         
when(producer.send(any(Message.class))).thenReturn(null);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            SendMessageDTO request = new SendMessageDTO();
-            request.setTopic("TopicA");
-            request.setBody("hello");
-
-            assertThatThrownBy(() -> adminClient.sendMessage(request))
-                    .isInstanceOf(BusinessException.class)
-                    .hasMessageContaining("null");
-
-            verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", 
null,
-                    "Message send did not succeed: null", "FAILED");
-        }
+        when(sendProducer.send(any(Message.class))).thenReturn(null);
+
+        SendMessageDTO request = new SendMessageDTO();
+        request.setTopic("TopicA");
+        request.setBody("hello");
+
+        assertThatThrownBy(() -> adminClient.sendMessage(request))
+                .isInstanceOf(BusinessException.class)
+                .hasMessageContaining("null");
+
+        verify(auditService).record("SEND_MESSAGE", "MESSAGE", "TopicA", null,
+                "Message send did not succeed: null", "FAILED");
     }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
index 5177df3e9..5ba2fb094 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQClusterProviderTest.java
@@ -300,7 +300,8 @@ class RocketMQClusterProviderTest {
         credential.setSecretKey("admin-sk");
         adminProperties.getCredentials().put("cluster-admin", credential);
         RuntimeAdminClientResolver resolver =
-                new RuntimeAdminClientResolver(instanceRepository, 
adminFactory, adminProperties);
+                new RuntimeAdminClientResolver(instanceRepository, 
adminFactory, adminProperties,
+                        
org.mockito.Mockito.mock(org.apache.rocketmq.studio.cluster.broker.MqClientPool.class));
         when(adminFactory.execute(eq("10.0.0.2:9876"), any(), 
eq("cluster-admin"), any()))
                 .thenAnswer(invocation -> 
invocation.<MqAdminExtFactory.AdminAction<Object>>getArgument(3)
                         .apply(adminExt));
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 607d1bddd..c22b07036 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
@@ -27,10 +27,10 @@ import org.apache.rocketmq.common.message.Message;
 import org.apache.rocketmq.common.message.MessageConst;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.remoting.RPCHook;
 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.MqClientPool;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.common.domain.PageResult;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
@@ -44,11 +44,9 @@ import org.junit.jupiter.api.Timeout;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.mockito.ArgumentCaptor;
 import org.mockito.Mock;
-import org.mockito.MockedConstruction;
 import org.mockito.junit.jupiter.MockitoExtension;
 
 import java.nio.charset.StandardCharsets;
-import java.util.ArrayList;
 import java.util.Base64;
 import java.util.List;
 import java.util.Set;
@@ -64,11 +62,7 @@ import static org.mockito.ArgumentMatchers.anyString;
 import static org.mockito.ArgumentMatchers.contains;
 import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.ArgumentMatchers.isNull;
-import static org.mockito.Mockito.doNothing;
-import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.lenient;
-import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
@@ -87,6 +81,12 @@ class RocketMQDLQProviderTest {
     @Mock
     private MQAdminExt adminExt;
 
+    @Mock
+    private DefaultMQPullConsumer pullConsumer;
+
+    @Mock
+    private DefaultMQProducer dlqProducer;
+
     private RocketMQDLQProvider provider;
 
     @BeforeEach
@@ -96,6 +96,18 @@ class RocketMQDLQProviderTest {
             MqAdminExtFactory.AdminAction<Object> action = 
invocation.getArgument(1);
             return action.apply(adminExt);
         });
+        
lenient().when(runtimeAdminClientResolver.executePullConsumer(anyString(), 
any()))
+                .thenAnswer(invocation -> {
+                    MqClientPool.ClientAction<DefaultMQPullConsumer, Object> 
action =
+                            invocation.getArgument(1);
+                    return action.apply(pullConsumer);
+                });
+        lenient().when(runtimeAdminClientResolver.executeProducer(anyString(), 
any()))
+                .thenAnswer(invocation -> {
+                    MqClientPool.ClientAction<DefaultMQProducer, Object> 
action =
+                            invocation.getArgument(1);
+                    return action.apply(dlqProducer);
+                });
         provider = new RocketMQDLQProvider(runtimeAdminClientResolver, 
auditService);
     }
 
@@ -226,20 +238,12 @@ class RocketMQDLQProviderTest {
     @Test
     void 
resendMessagesShouldNormalizeGroupNameBeforeBuildingDlqTopicAndAuditing() 
throws Exception {
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(null);
-                         doNothing().when(consumer).shutdown();
-                     });
-             MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class)) {
-            provider.resendMessages("instance-a", " group-a ", 100L, 200L, 
"target-topic");
-
-            assertThat(mockedConsumers.constructed()).hasSize(1);
-            
verify(mockedConsumers.constructed().get(0)).fetchSubscribeMessageQueues(dlqTopic);
-            assertThat(mockedProducers.constructed()).isEmpty();
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(null);
+        provider.resendMessages("instance-a", " group-a ", 100L, 200L, 
"target-topic");
+
+        
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
+        verify(pullConsumer).fetchSubscribeMessageQueues(dlqTopic);
+        verify(runtimeAdminClientResolver, 
never()).executeProducer(anyString(), any());
         verify(auditService).record(eq("RESEND_DLQ"), eq("DLQ"), 
eq("group-a"), eq(null),
                 contains("group=group-a, dlqTopic=%DLQ%group-a"), 
eq("NO_MESSAGES"));
     }
@@ -247,30 +251,13 @@ class RocketMQDLQProviderTest {
     @Test
     void resendMessagesDoesNotPullWhenDlqQueueSetIsNull() throws Exception {
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
-        List<List<?>> consumerConstructorArguments = new ArrayList<>();
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         consumerConstructorArguments.add(context.arguments());
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(null);
-                         doNothing().when(consumer).shutdown();
-                     });
-             MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class)) {
-            provider.resendMessages("instance-a", "group-a", 100L, 200L, 
"target-topic");
-
-            assertThat(mockedConsumers.constructed()).hasSize(1);
-            assertThat(consumerConstructorArguments).singleElement();
-            assertThat(consumerConstructorArguments.get(0)).hasSize(2);
-            assertThat(consumerConstructorArguments.get(0).get(1)).isNull();
-            DefaultMQPullConsumer consumer = 
mockedConsumers.constructed().get(0);
-            verify(consumer).setNamesrvAddr("namesrv-a:9876");
-            verify(consumer).start();
-            verify(consumer).fetchSubscribeMessageQueues(dlqTopic);
-            verify(consumer, never()).pull(any(MessageQueue.class), 
anyString(), anyLong(), anyInt());
-            verify(consumer).shutdown();
-            assertThat(mockedProducers.constructed()).isEmpty();
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(null);
+        provider.resendMessages("instance-a", "group-a", 100L, 200L, 
"target-topic");
+
+        
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
+        verify(pullConsumer).fetchSubscribeMessageQueues(dlqTopic);
+        verify(pullConsumer, never()).pull(any(MessageQueue.class), 
anyString(), anyLong(), anyInt());
+        verify(runtimeAdminClientResolver, 
never()).executeProducer(anyString(), any());
         verify(auditService).record(
                 eq("RESEND_DLQ"),
                 eq("DLQ"),
@@ -278,12 +265,10 @@ class RocketMQDLQProviderTest {
                 isNull(),
                 contains("matched=0, resent=0, failed=0"),
                 eq("NO_MESSAGES"));
-        verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
-        verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
     }
 
     @Test
-    void resendMessagesUsesSelectedInstanceCredentialHookForScanAndResend() 
throws Exception {
+    void resendMessagesUsesPooledClientsForScanAndResend() throws Exception {
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
         MessageQueue queue = new MessageQueue(dlqTopic, "broker-a", 0);
         MessageExt deadLetter = new MessageExt();
@@ -294,56 +279,27 @@ class RocketMQDLQProviderTest {
         PullResult pullResult = new PullResult(PullStatus.FOUND, 1L, 0L, 0L, 
List.of(deadLetter));
         SendResult sendResult = new SendResult();
         sendResult.setSendStatus(SendStatus.SEND_OK);
-        RPCHook credentialHook = mock(RPCHook.class);
-        List<List<?>> consumerConstructorArguments = new ArrayList<>();
-        List<List<?>> producerConstructorArguments = new ArrayList<>();
-        
when(runtimeAdminClientResolver.resolveCredentialHook("instance-a")).thenReturn(credentialHook);
-
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         consumerConstructorArguments.add(context.arguments());
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
-                         when(consumer.searchOffset(queue, 
100L)).thenReturn(0L);
-                         when(consumer.searchOffset(queue, 
200L)).thenReturn(0L);
-                         when(consumer.pull(queue, "*", 0L, 
32)).thenReturn(pullResult);
-                         doNothing().when(consumer).shutdown();
-                     });
-             MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         producerConstructorArguments.add(context.arguments());
-                         doNothing().when(producer).start();
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            provider.resendMessages("instance-a", "group-a", 100L, 200L, 
"target-topic");
-
-            assertThat(mockedConsumers.constructed()).singleElement();
-            assertThat(mockedProducers.constructed()).singleElement();
-        }
-
-        assertThat(consumerConstructorArguments).singleElement();
-        
assertThat(consumerConstructorArguments.get(0).get(1)).isSameAs(credentialHook);
-        assertThat(producerConstructorArguments).singleElement();
-        
assertThat(producerConstructorArguments.get(0).get(1)).isSameAs(credentialHook);
-        verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+        when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+        when(pullConsumer.searchOffset(queue, 200L)).thenReturn(0L);
+        when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+        when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+        provider.resendMessages("instance-a", "group-a", 100L, 200L, 
"target-topic");
+
+        
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
+        verify(runtimeAdminClientResolver).executeProducer(eq("instance-a"), 
any());
     }
 
     @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();
-        }
+        when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic))
+                .thenThrow(new IllegalStateException("broker unavailable"));
+        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(auditService).record(
                 eq("RESEND_DLQ"),
                 eq("DLQ"),
@@ -359,25 +315,17 @@ class RocketMQDLQProviderTest {
         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();
-        }
+        when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic))
+                .thenReturn(Set.of(unavailableQueue, emptyQueue));
+        when(pullConsumer.searchOffset(eq(unavailableQueue), anyLong()))
+                .thenThrow(new IllegalStateException("broker unavailable"));
+        when(pullConsumer.searchOffset(eq(emptyQueue), 
anyLong())).thenReturn(0L);
+        when(pullConsumer.pull(eq(emptyQueue), eq("*"), eq(0L), 
eq(32))).thenReturn(emptyResult);
+        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);
+
+        verify(runtimeAdminClientResolver, 
never()).executeProducer(anyString(), any());
         verify(auditService).record(
                 eq("RESEND_DLQ"),
                 eq("DLQ"),
@@ -393,21 +341,13 @@ class RocketMQDLQProviderTest {
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
         MessageQueue queue = new MessageQueue(dlqTopic, "broker-a", 0);
         PullResult stalledResult = new PullResult(PullStatus.FOUND, 10, 0, 10, 
List.of());
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
-                         when(consumer.searchOffset(eq(queue), 
anyLong())).thenReturn(10L);
-                         when(consumer.pull(eq(queue), eq("*"), eq(10L), 
eq(32))).thenReturn(stalledResult);
-                         doNothing().when(consumer).shutdown();
-                     });
-             MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class)) {
-            provider.resendMessages("instance-a", "group-a", 100L, 200L, 
"target-topic");
-
-            verify(mockedConsumers.constructed().get(0), times(1)).pull(queue, 
"*", 10L, 32);
-            assertThat(mockedProducers.constructed()).isEmpty();
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+        when(pullConsumer.searchOffset(eq(queue), anyLong())).thenReturn(10L);
+        when(pullConsumer.pull(eq(queue), eq("*"), eq(10L), 
eq(32))).thenReturn(stalledResult);
+        provider.resendMessages("instance-a", "group-a", 100L, 200L, 
"target-topic");
+
+        verify(pullConsumer, times(1)).pull(queue, "*", 10L, 32);
+        verify(runtimeAdminClientResolver, 
never()).executeProducer(anyString(), any());
     }
 
     @Test
@@ -425,33 +365,22 @@ class RocketMQDLQProviderTest {
         SendResult sendResult = new SendResult();
         sendResult.setSendStatus(SendStatus.SEND_OK);
 
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
-                         when(consumer.searchOffset(queue, 
100L)).thenReturn(10L);
-                         when(consumer.searchOffset(queue, 
200L)).thenReturn(50L);
-                         when(consumer.pull(queue, "*", 10L, 
32)).thenReturn(illegalOffset);
-                         when(consumer.pull(queue, "*", 20L, 
32)).thenReturn(foundAfterCorrection);
-                         when(consumer.pull(queue, "*", 40L, 
32)).thenReturn(endOfQueue);
-                         doNothing().when(consumer).shutdown();
-                     });
-             MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         doNothing().when(producer).start();
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "orders"))
-                    .extracting("matched", "resent", "failed", "outcome")
-                    .containsExactly(1, 1, 0, "SUCCESS");
-
-            DefaultMQPullConsumer consumer = 
mockedConsumers.constructed().get(0);
-            verify(consumer).pull(queue, "*", 10L, 32);
-            verify(consumer).pull(queue, "*", 20L, 32);
-            verify(consumer).pull(queue, "*", 40L, 32);
-            
verify(mockedProducers.constructed().get(0)).send(any(Message.class));
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+        when(pullConsumer.searchOffset(queue, 100L)).thenReturn(10L);
+        when(pullConsumer.searchOffset(queue, 200L)).thenReturn(50L);
+        when(pullConsumer.pull(queue, "*", 10L, 32)).thenReturn(illegalOffset);
+        when(pullConsumer.pull(queue, "*", 20L, 
32)).thenReturn(foundAfterCorrection);
+        when(pullConsumer.pull(queue, "*", 40L, 32)).thenReturn(endOfQueue);
+        when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+        assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "orders"))
+                .extracting("matched", "resent", "failed", "outcome")
+                .containsExactly(1, 1, 0, "SUCCESS");
+
+        DefaultMQPullConsumer consumer = pullConsumer;
+        verify(consumer).pull(queue, "*", 10L, 32);
+        verify(consumer).pull(queue, "*", 20L, 32);
+        verify(consumer).pull(queue, "*", 40L, 32);
+        verify(dlqProducer).send(any(Message.class));
         verify(auditService).record(
                 eq("RESEND_DLQ"),
                 eq("DLQ"),
@@ -474,29 +403,18 @@ class RocketMQDLQProviderTest {
         SendResult sendResult = new SendResult();
         sendResult.setSendStatus(SendStatus.FLUSH_DISK_TIMEOUT);
 
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
-                         when(consumer.searchOffset(queue, 
100L)).thenReturn(0L);
-                         when(consumer.searchOffset(queue, 
200L)).thenReturn(0L);
-                         when(consumer.pull(queue, "*", 0L, 
32)).thenReturn(pullResult);
-                         doNothing().when(consumer).shutdown();
-                     });
-             MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         doNothing().when(producer).start();
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "target-topic"))
-                    .extracting("matched", "resent", "failed", "outcome")
-                    .containsExactly(1, 0, 1, "FAILED");
-
-            assertThat(mockedConsumers.constructed()).hasSize(1);
-            assertThat(mockedProducers.constructed()).hasSize(1);
-            
verify(mockedProducers.constructed().get(0)).send(any(Message.class));
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+        when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+        when(pullConsumer.searchOffset(queue, 200L)).thenReturn(0L);
+        when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+        when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+        assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "target-topic"))
+                .extracting("matched", "resent", "failed", "outcome")
+                .containsExactly(1, 0, 1, "FAILED");
+
+        
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
+        verify(runtimeAdminClientResolver).executeProducer(eq("instance-a"), 
any());
+        verify(dlqProducer).send(any(Message.class));
         verify(auditService).record(
                 eq("RESEND_DLQ"),
                 eq("DLQ"),
@@ -522,35 +440,24 @@ class RocketMQDLQProviderTest {
         SendResult sendResult = new SendResult();
         sendResult.setSendStatus(SendStatus.SEND_OK);
 
-        try (MockedConstruction<DefaultMQPullConsumer> ignoredConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
-                         when(consumer.searchOffset(queue, 
100L)).thenReturn(0L);
-                         when(consumer.searchOffset(queue, 
200L)).thenReturn(0L);
-                         when(consumer.pull(queue, "*", 0L, 
32)).thenReturn(pullResult);
-                         doNothing().when(consumer).shutdown();
-                     });
-             MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         doNothing().when(producer).start();
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "target-topic"))
-                    .extracting("matched", "resent", "failed", "outcome")
-                    .containsExactly(1, 1, 0, "SUCCESS");
-
-            ArgumentCaptor<Message> messageCaptor = 
ArgumentCaptor.forClass(Message.class);
-            
verify(mockedProducers.constructed().get(0)).send(messageCaptor.capture());
-            Message resent = messageCaptor.getValue();
-            assertThat(resent.getProperties()).containsEntry("traceId", 
"trace-123")
-                    .containsEntry("tenantId", "tenant-a")
-                    .containsEntry("studio_dlq_origin_message_id", 
"msg-with-properties")
-                    .containsEntry("studio_dlq_origin_topic", dlqTopic);
-            assertThat(resent.getProperties())
-                    .doesNotContainEntry(MessageConst.PROPERTY_REAL_TOPIC, 
"system-topic-should-not-copy");
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+        when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+        when(pullConsumer.searchOffset(queue, 200L)).thenReturn(0L);
+        when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+        when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+        assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "target-topic"))
+                .extracting("matched", "resent", "failed", "outcome")
+                .containsExactly(1, 1, 0, "SUCCESS");
+
+        ArgumentCaptor<Message> messageCaptor = 
ArgumentCaptor.forClass(Message.class);
+        verify(dlqProducer).send(messageCaptor.capture());
+        Message resent = messageCaptor.getValue();
+        assertThat(resent.getProperties()).containsEntry("traceId", 
"trace-123")
+                .containsEntry("tenantId", "tenant-a")
+                .containsEntry("studio_dlq_origin_message_id", 
"msg-with-properties")
+                .containsEntry("studio_dlq_origin_topic", dlqTopic);
+        assertThat(resent.getProperties())
+                .doesNotContainEntry(MessageConst.PROPERTY_REAL_TOPIC, 
"system-topic-should-not-copy");
     }
 
     @Test
@@ -571,27 +478,16 @@ class RocketMQDLQProviderTest {
         SendResult sendResult = new SendResult();
         sendResult.setSendStatus(SendStatus.SEND_OK);
 
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
-                         when(consumer.searchOffset(queue, 
100L)).thenReturn(0L);
-                         when(consumer.searchOffset(queue, 
200L)).thenReturn(5001L);
-                         when(consumer.pull(queue, "*", 0L, 
32)).thenReturn(pullResult);
-                         doNothing().when(consumer).shutdown();
-                     });
-             MockedConstruction<DefaultMQProducer> mockedProducers =
-                     mockConstruction(DefaultMQProducer.class, (producer, 
context) -> {
-                         doNothing().when(producer).start();
-                         
when(producer.send(any(Message.class))).thenReturn(sendResult);
-                         doNothing().when(producer).shutdown();
-                     })) {
-            assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "target-topic"))
-                    .extracting("matched", "resent", "failed", "outcome", 
"scanIncomplete", "failedQueueCount")
-                    .containsExactly(5000, 5000, 0, "PARTIAL", true, 0);
-
-            verify(mockedProducers.constructed().get(0), 
times(5000)).send(any(Message.class));
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+        when(pullConsumer.searchOffset(queue, 100L)).thenReturn(0L);
+        when(pullConsumer.searchOffset(queue, 200L)).thenReturn(5001L);
+        when(pullConsumer.pull(queue, "*", 0L, 32)).thenReturn(pullResult);
+        when(dlqProducer.send(any(Message.class))).thenReturn(sendResult);
+        assertThat(provider.resendMessages("instance-a", "group-a", 100L, 
200L, "target-topic"))
+                .extracting("matched", "resent", "failed", "outcome", 
"scanIncomplete", "failedQueueCount")
+                .containsExactly(5000, 5000, 0, "PARTIAL", true, 0);
+
+        verify(dlqProducer, times(5000)).send(any(Message.class));
         verify(auditService).record(
                 eq("RESEND_DLQ"),
                 eq("DLQ"),
@@ -601,15 +497,6 @@ class RocketMQDLQProviderTest {
                 eq("PARTIAL"));
     }
 
-    @Test
-    void createsUniqueProducerGroupsForConcurrentDlqResends() {
-        String first = RocketMQDLQProvider.nextResendProducerGroup();
-        String second = RocketMQDLQProvider.nextResendProducerGroup();
-
-        assertThat(first).startsWith("studio-dlq-resend-");
-        
assertThat(second).startsWith("studio-dlq-resend-").isNotEqualTo(first);
-    }
-
     @Test
     void exportMessagesReturnsMappedDeadLetterMessages() throws Exception {
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
@@ -623,48 +510,36 @@ class RocketMQDLQProviderTest {
         deadLetter.setKeys("key-a,key-b");
         deadLetter.setBody("hello dlq".getBytes(StandardCharsets.UTF_8));
         PullResult pullResult = new PullResult(PullStatus.FOUND, 1L, 0L, 0L, 
List.of(deadLetter));
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
-                         when(consumer.searchOffset(eq(queue), 
anyLong())).thenReturn(0L);
-                         when(consumer.pull(eq(queue), eq("*"), eq(0L), 
eq(32))).thenReturn(pullResult);
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<DLQMessageVO> exported =
-                    provider.exportMessages("instance-a", "group-a", 100L, 
200L, 1000);
-
-            assertThat(exported).hasSize(1);
-            DLQMessageVO vo = exported.get(0);
-            assertThat(vo.getMsgId()).isEqualTo("msg-1");
-            assertThat(vo.getTopic()).isEqualTo(dlqTopic);
-            assertThat(vo.getQueueId()).isEqualTo(0);
-            assertThat(vo.getOffset()).isEqualTo(5L);
-            assertThat(vo.getStoreTime()).isEqualTo(150L);
-            assertThat(vo.getKeys()).isEqualTo("key-a,key-b");
-            assertThat(vo.getBody()).isEqualTo("hello dlq");
-            assertThat(vo.getBodyBase64())
-                    .isEqualTo(Base64.getEncoder().encodeToString("hello 
dlq".getBytes(StandardCharsets.UTF_8)));
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+        when(pullConsumer.searchOffset(eq(queue), anyLong())).thenReturn(0L);
+        when(pullConsumer.pull(eq(queue), eq("*"), eq(0L), 
eq(32))).thenReturn(pullResult);
+        List<DLQMessageVO> exported =
+                provider.exportMessages("instance-a", "group-a", 100L, 200L, 
1000);
+
+        assertThat(exported).hasSize(1);
+        DLQMessageVO vo = exported.get(0);
+        assertThat(vo.getMsgId()).isEqualTo("msg-1");
+        assertThat(vo.getTopic()).isEqualTo(dlqTopic);
+        assertThat(vo.getQueueId()).isEqualTo(0);
+        assertThat(vo.getOffset()).isEqualTo(5L);
+        assertThat(vo.getStoreTime()).isEqualTo(150L);
+        assertThat(vo.getKeys()).isEqualTo("key-a,key-b");
+        assertThat(vo.getBody()).isEqualTo("hello dlq");
+        assertThat(vo.getBodyBase64())
+                .isEqualTo(Base64.getEncoder().encodeToString("hello 
dlq".getBytes(StandardCharsets.UTF_8)));
     }
 
     @Test
     void exportMessagesHonorsMaxCountCap() throws Exception {
         String dlqTopic = MixAll.DLQ_GROUP_TOPIC_PREFIX + "group-a";
         MessageQueue queue = new MessageQueue(dlqTopic, "broker-a", 0);
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
-                         when(consumer.searchOffset(eq(queue), 
anyLong())).thenReturn(0L);
-                         when(consumer.pull(eq(queue), eq("*"), eq(0L), 
eq(32)))
-                                 .thenReturn(new 
PullResult(PullStatus.NO_NEW_MSG, 1L, 0L, 0L, List.of()));
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            // maxCount=0 falls back to the hard cap instead of failing; scan 
still completes.
-            List<DLQMessageVO> exported =
-                    provider.exportMessages("instance-a", "group-a", 100L, 
200L, 0);
-            assertThat(exported).isEmpty();
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues(dlqTopic)).thenReturn(Set.of(queue));
+        when(pullConsumer.searchOffset(eq(queue), anyLong())).thenReturn(0L);
+        when(pullConsumer.pull(eq(queue), eq("*"), eq(0L), eq(32)))
+                .thenReturn(new PullResult(PullStatus.NO_NEW_MSG, 1L, 0L, 0L, 
List.of()));
+        // maxCount=0 falls back to the hard cap instead of failing; scan 
still completes.
+        List<DLQMessageVO> exported =
+                provider.exportMessages("instance-a", "group-a", 100L, 200L, 
0);
+        assertThat(exported).isEmpty();
     }
 }
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
index 2eb00365b..2501eda4b 100644
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
+++ 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/RocketMQMessageProviderTest.java
@@ -27,10 +27,10 @@ import org.apache.rocketmq.common.message.MessageDecoder;
 import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageId;
 import org.apache.rocketmq.common.message.MessageQueue;
-import org.apache.rocketmq.remoting.RPCHook;
 import org.apache.rocketmq.remoting.protocol.body.ClusterInfo;
 import org.apache.rocketmq.remoting.protocol.route.BrokerData;
 import org.apache.rocketmq.studio.cluster.broker.MqAdminExtFactory;
+import org.apache.rocketmq.studio.cluster.broker.MqClientPool;
 import org.apache.rocketmq.studio.cluster.broker.RuntimeAdminClientResolver;
 import org.apache.rocketmq.studio.common.exception.BusinessException;
 import org.apache.rocketmq.studio.common.domain.enums.DeliveryStatus;
@@ -45,7 +45,6 @@ import org.junit.jupiter.api.Timeout;
 import org.junit.jupiter.api.extension.ExtendWith;
 import org.mockito.ArgumentCaptor;
 import org.mockito.Mock;
-import org.mockito.MockedConstruction;
 import org.mockito.junit.jupiter.MockitoExtension;
 
 import java.net.InetSocketAddress;
@@ -67,11 +66,8 @@ import static org.mockito.ArgumentMatchers.anyInt;
 import static org.mockito.ArgumentMatchers.anyLong;
 import static org.mockito.ArgumentMatchers.anyString;
 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.mock;
-import static org.mockito.Mockito.mockConstruction;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
@@ -86,6 +82,9 @@ class RocketMQMessageProviderTest {
     @Mock
     private RuntimeAdminClientResolver runtimeAdminClientResolver;
 
+    @Mock
+    private DefaultMQPullConsumer pullConsumer;
+
     private RocketMQMessageProvider provider;
 
     @BeforeEach
@@ -105,61 +104,36 @@ class RocketMQMessageProviderTest {
         });
         lenient().when(adminExt.examineBrokerClusterInfo())
                 
.thenReturn(clusterInfoWithBrokerAddresses("172.30.10.100:10911"));
+        
lenient().when(runtimeAdminClientResolver.executePullConsumer(anyString(), 
any()))
+                .thenAnswer(invocation -> {
+                    MqClientPool.ClientAction<DefaultMQPullConsumer, Object> 
action =
+                            invocation.getArgument(1);
+                    return action.apply(pullConsumer);
+                });
         provider = new RocketMQMessageProvider(runtimeAdminClientResolver);
     }
 
     @Test
     void queryByTopicReturnsEmptyListWhenQueueSetIsNull() throws Exception {
-        List<List<?>> constructorArguments = new ArrayList<>();
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         constructorArguments.add(context.arguments());
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(null);
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<MessageRecordVO> messages = 
provider.queryMessages("instance-a", "TopicA", null, null, null, 100L, 200L);
-
-            assertThat(messages).isEmpty();
-            assertThat(mockedConsumers.constructed()).hasSize(1);
-            DefaultMQPullConsumer consumer = 
mockedConsumers.constructed().get(0);
-            assertThat(constructorArguments).singleElement();
-            assertThat(constructorArguments.get(0)).hasSize(2);
-            
assertThat(constructorArguments.get(0).get(0)).isEqualTo("studio-msg-query-group");
-            assertThat(constructorArguments.get(0).get(1)).isNull();
-            verify(consumer).setNamesrvAddr("namesrv-a:9876");
-            verify(consumer).start();
-            verify(consumer).fetchSubscribeMessageQueues("TopicA");
-            verify(consumer, never()).pull(any(MessageQueue.class), 
anyString(), anyLong(), anyInt());
-            verify(consumer).shutdown();
-        }
-        verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
-        verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
+        
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(null);
+
+        List<MessageRecordVO> messages = provider.queryMessages("instance-a", 
"TopicA", null, null, null, 100L, 200L);
+
+        assertThat(messages).isEmpty();
+        verify(pullConsumer).fetchSubscribeMessageQueues("TopicA");
+        verify(pullConsumer, never()).pull(any(MessageQueue.class), 
anyString(), anyLong(), anyInt());
+        
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
         verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
     }
 
     @Test
-    void queryByTopicUsesSelectedInstanceCredentialHook() throws Exception {
-        RPCHook credentialHook = mock(RPCHook.class);
-        List<List<?>> constructorArguments = new ArrayList<>();
-        
when(runtimeAdminClientResolver.resolveCredentialHook("instance-a")).thenReturn(credentialHook);
-
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         constructorArguments.add(context.arguments());
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of());
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            provider.queryMessages("instance-a", "TopicA", null, null, null, 
100L, 200L);
-
-            assertThat(mockedConsumers.constructed()).singleElement();
-            assertThat(constructorArguments).singleElement();
-            assertThat(constructorArguments.get(0)).hasSize(2);
-            
assertThat(constructorArguments.get(0).get(0)).isEqualTo("studio-msg-query-group");
-            
assertThat(constructorArguments.get(0).get(1)).isSameAs(credentialHook);
-        }
-        verify(runtimeAdminClientResolver).resolveCredentialHook("instance-a");
+    void queryByTopicReturnsEmptyListWhenNoQueuesExist() throws Exception {
+        
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of());
+
+        List<MessageRecordVO> messages = provider.queryMessages("instance-a", 
"TopicA", null, null, null, 100L, 200L);
+
+        assertThat(messages).isEmpty();
+        
verify(runtimeAdminClientResolver).executePullConsumer(eq("instance-a"), any());
     }
 
     @Test
@@ -170,7 +144,6 @@ class RocketMQMessageProviderTest {
                 .hasMessage("Message query start time must be before end time")
                 .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(400));
 
-        verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
         verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
         verify(adminExt, never()).queryMessage(anyString(), anyString(), 
anyInt(), anyLong(), anyLong());
     }
@@ -183,7 +156,6 @@ class RocketMQMessageProviderTest {
                 .hasMessage("Message query start time must be before end time")
                 .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(400));
 
-        verify(runtimeAdminClientResolver).resolveEndpoint("instance-a");
         verify(runtimeAdminClientResolver).execute(eq("instance-a"), any());
         verify(adminExt, never()).queryMessage(anyString(), anyString(), 
anyInt(), anyLong(), anyLong());
     }
@@ -244,17 +216,14 @@ class RocketMQMessageProviderTest {
 
     @Test
     void queryByTopicSurfacesPullConsumerFailure() throws Exception {
-        try (MockedConstruction<DefaultMQPullConsumer> ignored =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doThrow(new IllegalStateException("nameserver 
unavailable")).when(consumer).start();
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            assertThatThrownBy(() -> provider.queryMessages(
-                    "instance-a", "TopicA", null, null, null, 100L, 200L))
-                    .isInstanceOf(BusinessException.class)
-                    .hasMessage("Failed to query messages by topic: nameserver 
unavailable")
-                    .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(502));
-        }
+        when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+                .thenThrow(new IllegalStateException("nameserver 
unavailable"));
+
+        assertThatThrownBy(() -> provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, 100L, 200L))
+                .isInstanceOf(BusinessException.class)
+                .hasMessage("Failed to query messages by topic: nameserver 
unavailable")
+                .satisfies(error -> assertThat(((BusinessException) 
error).getCode()).isEqualTo(502));
     }
 
     @Test
@@ -262,20 +231,15 @@ class RocketMQMessageProviderTest {
     void queryByTopicStopsWhenPullOffsetDoesNotAdvance() throws Exception {
         MessageQueue queue = new MessageQueue("TopicA", "broker-a", 0);
         PullResult stalledResult = new PullResult(PullStatus.FOUND, 10, 0, 10, 
List.of());
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
-                         mockQueueWindow(consumer, queue, 100L, 200L, 10L, 
10L, 11L, 11L);
-                         when(consumer.pull(eq(queue), eq("*"), eq(10L), 
eq(32))).thenReturn(stalledResult);
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<MessageRecordVO> messages = provider.queryMessages(
-                    "instance-a", "TopicA", null, null, null, 100L, 200L);
-
-            assertThat(messages).isEmpty();
-            verify(mockedConsumers.constructed().get(0), times(1)).pull(queue, 
"*", 10L, 32);
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+        mockQueueWindow(pullConsumer, queue, 100L, 200L, 10L, 10L, 11L, 11L);
+        when(pullConsumer.pull(eq(queue), eq("*"), eq(10L), 
eq(32))).thenReturn(stalledResult);
+
+        List<MessageRecordVO> messages = provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, 100L, 200L);
+
+        assertThat(messages).isEmpty();
+        verify(pullConsumer, times(1)).pull(queue, "*", 10L, 32);
     }
 
     @Test
@@ -286,23 +250,18 @@ class RocketMQMessageProviderTest {
                 List.of(topicMessage("older", 150L)));
         PullResult newerPullResult = new PullResult(PullStatus.FOUND, 11L, 
10L, 10L,
                 List.of(topicMessage("newer", 250L)));
-        try (MockedConstruction<DefaultMQPullConsumer> ignored =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         when(consumer.fetchSubscribeMessageQueues("TopicA"))
-                                 .thenReturn(new 
LinkedHashSet<>(List.of(olderQueue, newerQueue)));
-                         mockQueueWindow(consumer, olderQueue, 100L, 300L, 
10L, 10L, 11L, 11L);
-                         mockQueueWindow(consumer, newerQueue, 100L, 300L, 
10L, 10L, 11L, 11L);
-                         when(consumer.pull(olderQueue, "*", 10L, 
32)).thenReturn(olderPullResult);
-                         when(consumer.pull(newerQueue, "*", 10L, 
32)).thenReturn(newerPullResult);
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<MessageRecordVO> messages = provider.queryMessages(
-                    "instance-a", "TopicA", null, null, null, 100L, 300L);
-
-            assertThat(messages).extracting(MessageRecordVO::getMsgId)
-                    .containsExactly("newer", "older");
-        }
+        when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+                .thenReturn(new LinkedHashSet<>(List.of(olderQueue, 
newerQueue)));
+        mockQueueWindow(pullConsumer, olderQueue, 100L, 300L, 10L, 10L, 11L, 
11L);
+        mockQueueWindow(pullConsumer, newerQueue, 100L, 300L, 10L, 10L, 11L, 
11L);
+        when(pullConsumer.pull(olderQueue, "*", 10L, 
32)).thenReturn(olderPullResult);
+        when(pullConsumer.pull(newerQueue, "*", 10L, 
32)).thenReturn(newerPullResult);
+
+        List<MessageRecordVO> messages = provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+        assertThat(messages).extracting(MessageRecordVO::getMsgId)
+                .containsExactly("newer", "older");
     }
 
     @Test
@@ -316,25 +275,19 @@ class RocketMQMessageProviderTest {
         PullResult illegalOffset = new PullResult(PullStatus.OFFSET_ILLEGAL, 
20L, 0L, 30L, null);
         PullResult foundAfterCorrection = new PullResult(PullStatus.FOUND, 
40L, 20L, 30L, List.of(message));
         PullResult endOfQueue = new PullResult(PullStatus.NO_NEW_MSG, 50L, 
40L, 40L, List.of());
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
-                         mockQueueWindow(consumer, queue, 100L, 200L, 10L, 
10L, 50L, 50L);
-                         when(consumer.pull(eq(queue), eq("*"), eq(10L), 
eq(32))).thenReturn(illegalOffset);
-                         when(consumer.pull(eq(queue), eq("*"), eq(20L), 
eq(32))).thenReturn(foundAfterCorrection);
-                         when(consumer.pull(eq(queue), eq("*"), eq(40L), 
eq(32))).thenReturn(endOfQueue);
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<MessageRecordVO> messages = provider.queryMessages(
-                    "instance-a", "TopicA", null, null, null, 100L, 200L);
-
-            
assertThat(messages).extracting(MessageRecordVO::getMsgId).containsExactly("msg-after-correction");
-            DefaultMQPullConsumer consumer = 
mockedConsumers.constructed().get(0);
-            verify(consumer).pull(queue, "*", 10L, 32);
-            verify(consumer).pull(queue, "*", 20L, 32);
-            verify(consumer).pull(queue, "*", 40L, 32);
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+        mockQueueWindow(pullConsumer, queue, 100L, 200L, 10L, 10L, 50L, 50L);
+        when(pullConsumer.pull(eq(queue), eq("*"), eq(10L), 
eq(32))).thenReturn(illegalOffset);
+        when(pullConsumer.pull(eq(queue), eq("*"), eq(20L), 
eq(32))).thenReturn(foundAfterCorrection);
+        when(pullConsumer.pull(eq(queue), eq("*"), eq(40L), 
eq(32))).thenReturn(endOfQueue);
+
+        List<MessageRecordVO> messages = provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, 100L, 200L);
+
+        
assertThat(messages).extracting(MessageRecordVO::getMsgId).containsExactly("msg-after-correction");
+        verify(pullConsumer).pull(queue, "*", 10L, 32);
+        verify(pullConsumer).pull(queue, "*", 20L, 32);
+        verify(pullConsumer).pull(queue, "*", 40L, 32);
     }
 
     @Test
@@ -347,25 +300,20 @@ class RocketMQMessageProviderTest {
                         .toList());
         PullResult newerPullResult = new PullResult(PullStatus.FOUND, 11L, 
10L, 10L,
                 List.of(topicMessage("newer", 250L)));
-        try (MockedConstruction<DefaultMQPullConsumer> ignored =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         when(consumer.fetchSubscribeMessageQueues("TopicA"))
-                                 .thenReturn(new 
LinkedHashSet<>(List.of(olderQueue, newerQueue)));
-                         mockQueueWindow(consumer, olderQueue, 100L, 300L, 
10L, 10L, 11L, 11L);
-                         mockQueueWindow(consumer, newerQueue, 100L, 300L, 
10L, 10L, 11L, 11L);
-                         when(consumer.pull(olderQueue, "*", 10L, 
32)).thenReturn(olderPullResult);
-                         when(consumer.pull(newerQueue, "*", 10L, 
32)).thenReturn(newerPullResult);
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<MessageRecordVO> messages = provider.queryMessages(
-                    "instance-a", "TopicA", null, null, null, 100L, 300L);
-
-            assertThat(messages).hasSize(200);
-            assertThat(messages.get(0).getMsgId()).isEqualTo("newer");
-            assertThat(messages).extracting(MessageRecordVO::getMsgId)
-                    .contains("newer");
-        }
+        when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+                .thenReturn(new LinkedHashSet<>(List.of(olderQueue, 
newerQueue)));
+        mockQueueWindow(pullConsumer, olderQueue, 100L, 300L, 10L, 10L, 11L, 
11L);
+        mockQueueWindow(pullConsumer, newerQueue, 100L, 300L, 10L, 10L, 11L, 
11L);
+        when(pullConsumer.pull(olderQueue, "*", 10L, 
32)).thenReturn(olderPullResult);
+        when(pullConsumer.pull(newerQueue, "*", 10L, 
32)).thenReturn(newerPullResult);
+
+        List<MessageRecordVO> messages = provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+        assertThat(messages).hasSize(200);
+        assertThat(messages.get(0).getMsgId()).isEqualTo("newer");
+        assertThat(messages).extracting(MessageRecordVO::getMsgId)
+                .contains("newer");
     }
 
     @Test
@@ -376,23 +324,18 @@ class RocketMQMessageProviderTest {
                 List.of(topicMessage("message-a", 250L)));
         PullResult secondPullResult = new PullResult(PullStatus.FOUND, 11L, 
10L, 10L,
                 List.of(topicMessage("message-b", 250L)));
-        try (MockedConstruction<DefaultMQPullConsumer> ignored =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         when(consumer.fetchSubscribeMessageQueues("TopicA"))
-                                 .thenReturn(new 
LinkedHashSet<>(List.of(firstQueue, secondQueue)));
-                         mockQueueWindow(consumer, firstQueue, 100L, 300L, 
10L, 10L, 11L, 11L);
-                         mockQueueWindow(consumer, secondQueue, 100L, 300L, 
10L, 10L, 11L, 11L);
-                         when(consumer.pull(firstQueue, "*", 10L, 
32)).thenReturn(firstPullResult);
-                         when(consumer.pull(secondQueue, "*", 10L, 
32)).thenReturn(secondPullResult);
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<MessageRecordVO> messages = provider.queryMessages(
-                    "instance-a", "TopicA", null, null, null, 100L, 300L);
-
-            assertThat(messages).extracting(MessageRecordVO::getMsgId)
-                    .containsExactly("message-b", "message-a");
-        }
+        when(pullConsumer.fetchSubscribeMessageQueues("TopicA"))
+                .thenReturn(new LinkedHashSet<>(List.of(firstQueue, 
secondQueue)));
+        mockQueueWindow(pullConsumer, firstQueue, 100L, 300L, 10L, 10L, 11L, 
11L);
+        mockQueueWindow(pullConsumer, secondQueue, 100L, 300L, 10L, 10L, 11L, 
11L);
+        when(pullConsumer.pull(firstQueue, "*", 10L, 
32)).thenReturn(firstPullResult);
+        when(pullConsumer.pull(secondQueue, "*", 10L, 
32)).thenReturn(secondPullResult);
+
+        List<MessageRecordVO> messages = provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, 100L, 300L);
+
+        assertThat(messages).extracting(MessageRecordVO::getMsgId)
+                .containsExactly("message-b", "message-a");
     }
 
     @Test
@@ -637,25 +580,19 @@ class RocketMQMessageProviderTest {
         long begin = 10_000L;
         long end = 49_999L;
         long maxOffsetExclusive = 40_000L;
-        try (MockedConstruction<DefaultMQPullConsumer> mockedConsumers =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
-                         mockQueueWindow(consumer, queue, begin, end, 0L, 0L, 
maxOffsetExclusive, maxOffsetExclusive);
-                         when(consumer.pull(eq(queue), eq("*"), anyLong(), 
eq(32)))
-                                 .thenAnswer(invocation -> 
topicPullBatch(invocation.getArgument(2), maxOffsetExclusive, begin));
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<MessageRecordVO> messages = provider.queryMessages(
-                    "instance-a", "TopicA", null, null, null, begin, end);
-
-            assertThat(messages).hasSize(200);
-            assertThat(messages.get(0).getMsgId()).isEqualTo("msg-39999");
-            assertThat(messages.get(199).getMsgId()).isEqualTo("msg-39800");
-            DefaultMQPullConsumer consumer = 
mockedConsumers.constructed().get(0);
-            verify(consumer).pull(queue, "*", 8_000L, 32);
-            verify(consumer, never()).pull(queue, "*", 0L, 32);
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+        mockQueueWindow(pullConsumer, queue, begin, end, 0L, 0L, 
maxOffsetExclusive, maxOffsetExclusive);
+        when(pullConsumer.pull(eq(queue), eq("*"), anyLong(), eq(32)))
+                .thenAnswer(invocation -> 
topicPullBatch(invocation.getArgument(2), maxOffsetExclusive, begin));
+
+        List<MessageRecordVO> messages = provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, begin, end);
+
+        assertThat(messages).hasSize(200);
+        assertThat(messages.get(0).getMsgId()).isEqualTo("msg-39999");
+        assertThat(messages.get(199).getMsgId()).isEqualTo("msg-39800");
+        verify(pullConsumer).pull(queue, "*", 8_000L, 32);
+        verify(pullConsumer, never()).pull(queue, "*", 0L, 32);
     }
 
     @Test
@@ -665,27 +602,22 @@ class RocketMQMessageProviderTest {
         long end = 69_999L;
         long maxOffsetExclusive = 50_000L;
         List<Long> pulledOffsets = new ArrayList<>();
-        try (MockedConstruction<DefaultMQPullConsumer> ignored =
-                     mockConstruction(DefaultMQPullConsumer.class, (consumer, 
context) -> {
-                         doNothing().when(consumer).start();
-                         
when(consumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
-                         mockQueueWindow(consumer, queue, begin, end, 0L, 0L, 
maxOffsetExclusive, maxOffsetExclusive);
-                         when(consumer.pull(eq(queue), eq("*"), anyLong(), 
eq(32)))
-                                 .thenAnswer(invocation -> {
-                                     long offset = invocation.getArgument(2);
-                                     pulledOffsets.add(offset);
-                                     return topicPullBatch(offset, 
maxOffsetExclusive, begin);
-                                 });
-                         doNothing().when(consumer).shutdown();
-                     })) {
-            List<MessageRecordVO> messages = provider.queryMessages(
-                    "instance-a", "TopicA", null, null, null, begin, end);
-
-            assertThat(messages).hasSize(200);
-            assertThat(pulledOffsets).isNotEmpty();
-            assertThat(pulledOffsets.get(0)).isEqualTo(18_000L);
-            assertThat(pulledOffsets).allMatch(offset -> offset >= 18_000L);
-        }
+        
when(pullConsumer.fetchSubscribeMessageQueues("TopicA")).thenReturn(Set.of(queue));
+        mockQueueWindow(pullConsumer, queue, begin, end, 0L, 0L, 
maxOffsetExclusive, maxOffsetExclusive);
+        when(pullConsumer.pull(eq(queue), eq("*"), anyLong(), eq(32)))
+                .thenAnswer(invocation -> {
+                    long offset = invocation.getArgument(2);
+                    pulledOffsets.add(offset);
+                    return topicPullBatch(offset, maxOffsetExclusive, begin);
+                });
+
+        List<MessageRecordVO> messages = provider.queryMessages(
+                "instance-a", "TopicA", null, null, null, begin, end);
+
+        assertThat(messages).hasSize(200);
+        assertThat(pulledOffsets).isNotEmpty();
+        assertThat(pulledOffsets.get(0)).isEqualTo(18_000L);
+        assertThat(pulledOffsets).allMatch(offset -> offset >= 18_000L);
     }
 
     private MQClientAPIImpl mockOffsetLookupClient() {
diff --git 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientNameTest.java
 
b/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientNameTest.java
deleted file mode 100644
index b57a322ce..000000000
--- 
a/server/src/test/java/org/apache/rocketmq/studio/provider/apache/ShortLivedClientNameTest.java
+++ /dev/null
@@ -1,37 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements.  See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License.  You may obtain a copy of the License at
- *
- *     http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-package org.apache.rocketmq.studio.provider.apache;
-
-import static org.assertj.core.api.Assertions.assertThat;
-
-import java.util.Set;
-import java.util.stream.Collectors;
-import java.util.stream.IntStream;
-import org.junit.jupiter.api.Test;
-
-class ShortLivedClientNameTest {
-
-    @Test
-    void generatesUniqueNamesWithoutClockDelays() {
-        Set<String> names = IntStream.range(0, 100)
-                .parallel()
-                .mapToObj(ignored -> 
RocketMQAdminClientImpl.nextMessageSenderGroup())
-                .collect(Collectors.toSet());
-
-        assertThat(names).hasSize(100).allMatch(name -> 
name.startsWith("studio-msg-sender-"));
-    }
-}

Reply via email to