使用Node.js+amqplib测试RabbitMQ无法接收消息的排查求助
RabbitMQ消息未被消费者接收问题排查
问题现象
首次使用Node.js结合amqplib测试RabbitMQ,操作流程如下:
- 执行命令
node ./messages/consumer.js,终端输出:- Connected to RabbitMQ
- Channel created
- Waiting for messages...
- 执行命令
node ./messages/producer.js,终端输出:- Connected to RabbitMQ
- Channel created
- Message sent: Hello, world!
- Connection to RabbitMQ closed
通过RabbitMQ管理控制台可见test_exchange、test_queue和test_key,但无任何消息记录,消费者终端始终显示“Waiting for messages...”,未接收到消息。
核心疏漏点与修复方案
1. 生产者提前关闭连接导致消息丢失
producer.js中,调用channel.publish后立即执行rabbitmq.close(),但publish是异步非阻塞操作,此时消息可能还未完全写入RabbitMQ就被强制断开连接,直接导致消息丢失。
修复方法:开启消息确认机制,等待RabbitMQ确认接收消息后再关闭连接。
2. 队列绑定逻辑重复(非致命但不规范)
生产者和消费者都执行了assertQueue和bindQueue操作,虽然不会直接导致问题,但更合理的分工是由消费者负责队列的声明与绑定,生产者仅需声明交换机并发送消息,避免重复操作带来的潜在冲突。
修正后的代码
producer.js
const rabbitmq = require('../lib/rabbitmq'); const config = require('../config/config'); async function produceMessage(message) { try { const channel = await rabbitmq.createChannel(); const exchange = config.rabbitmq.exchange; const key = config.rabbitmq.routingKey; // 仅声明交换机,队列与绑定交由消费者处理 await channel.assertExchange(exchange, 'direct', { durable: true }); // 开启消息确认模式 channel.confirmSelect(); const messageBuffer = Buffer.from(message); channel.publish(exchange, key, messageBuffer); console.log(`Message sent: ${message}`); // 等待RabbitMQ确认消息接收后再关闭连接 await channel.waitForConfirms(); await rabbitmq.close(); } catch (error) { console.error('Error producing message', error); } } produceMessage('Hello, world!');
consumer.js(优化绑定逻辑)
const rabbitmq = require('../lib/rabbitmq'); const config = require('../config/config'); async function consumeMessage() { try { const channel = await rabbitmq.createChannel(); const exchange = config.rabbitmq.exchange; const queue = config.rabbitmq.queue; const key = config.rabbitmq.routingKey; // 消费者负责队列声明与绑定 await channel.assertExchange(exchange, 'direct', { durable: true }); await channel.assertQueue(queue, { durable: true }); await channel.bindQueue(queue, exchange, key); channel.consume(queue, (msg) => { if (msg) { console.log(`Message received: ${msg.content.toString()}`); channel.ack(msg); } }, { noAck: false }); console.log('Waiting for messages...'); } catch (error) { console.error('Error consuming message', error); } } consumeMessage();
额外检查项
- 确认RabbitMQ服务正常运行,
5672端口可正常访问 - 登录RabbitMQ管理控制台,检查队列的
Binding是否正确关联到目标交换机与路由键 - 确保消费者进程在生产者发送消息前已启动并处于等待状态
内容的提问来源于stack exchange,提问作者Alex Aung
相关产品推荐
相关产品推荐

