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
相关产品推荐
相关产品推荐

