You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RocketMQ异步发送消息时连接NameServer失败问题求助

RocketMQ异步发送消息抛出RemotingConnectException异常排查与解决

问题现象

异步发送消息时抛出org.apache.rocketmq.remoting.exception.RemotingConnectException,提示连接192.168.2.115:9876失败;仅启动并关闭生产者的同步测试无异常。

依赖配置

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-client</artifactId>
    <version>4.9.3</version>
</dependency>

异常复现代码(异步发送)

public static void main(String[] args) throws Exception {
    // 实例化生产者
    DefaultMQProducer producer = new DefaultMQProducer("async_group");
    producer.setNamesrvAddr("192.168.2.115:9876");
    producer.start();
    
    Message message = new Message("async_basic", "async_tag", ("i am async body,num:").getBytes());
    
    producer.send(message, new SendCallback() {
        @Override
        public void onSuccess(SendResult sendResult) {
            System.out.println("发送状态:" + sendResult);
        }

        @Override
        public void onException(Throwable throwable) {
            throwable.printStackTrace();
        }
    });
    
    // 直接关闭生产者
    producer.shutdown();
}

无异常测试代码(仅启动关闭生产者)

public static void main(String[] args) throws Exception {
    DefaultMQProducer producer = new DefaultMQProducer("async_group");
    producer.setNamesrvAddr("192.168.2.115:9876");
    producer.start();
    producer.shutdown();
}

问题根因

异步发送逻辑是非阻塞的:调用producer.send()后,主线程立即执行后续的producer.shutdown(),此时生产者已销毁与NameServer、Broker的连接资源,但后台线程还在处理异步发送的请求,最终因连接已关闭抛出RemotingConnectException。而同步测试代码未触发实际消息发送,仅完成生产者启停流程,因此无异常。

解决方案

1. 用CountDownLatch等待异步操作完成(推荐)

通过计数器让主线程等待异步发送回调执行完毕,再关闭生产者:

public static void main(String[] args) throws Exception {
    DefaultMQProducer producer = new DefaultMQProducer("async_group");
    producer.setNamesrvAddr("192.168.2.115:9876");
    producer.start();
    
    Message message = new Message("async_basic", "async_tag", ("i am async body,num:").getBytes());
    CountDownLatch latch = new CountDownLatch(1);
    
    producer.send(message, new SendCallback() {
        @Override
        public void onSuccess(SendResult sendResult) {
            System.out.println("发送状态:" + sendResult);
            latch.countDown();
        }

        @Override
        public void onException(Throwable throwable) {
            throwable.printStackTrace();
            latch.countDown();
        }
    });
    
    latch.await(); // 主线程阻塞等待异步回调完成
    producer.shutdown();
}

2. 测试场景下延长主线程存活时间

如果仅用于本地测试,可通过Thread.sleep()让主线程等待足够时长,确保异步发送完成:

// ... 发送消息代码 ...
Thread.sleep(5000); // 等待5秒,根据实际发送耗时调整
producer.shutdown();

3. 生产环境:保持生产者长生命周期

生产环境中,生产者应作为单例组件管理,仅在服务停止阶段(如Spring容器销毁时)调用producer.shutdown(),避免发送消息后立即关闭。

内容的提问来源于stack exchange,提问作者ggyuzhouwudi

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 08:25:56