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

使用Node.js+amqplib测试RabbitMQ无法接收消息的排查求助

RabbitMQ消息未被消费者接收问题排查

问题现象

首次使用Node.js结合amqplib测试RabbitMQ,操作流程如下:

  1. 执行命令 node ./messages/consumer.js,终端输出:
    • Connected to RabbitMQ
    • Channel created
    • Waiting for messages...
  2. 执行命令 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 11:45:24