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

RabbitMQ ConfirmChannel发布确认异常:未ACK消息仍返回success

RabbitMQ生产者无法获取消费者确认反馈的问题

我用JavaScript编写RabbitMQ的Hello World测试代码,需求是让生产者获取**消息被消费者确认(ack)**的反馈,但遇到两个问题:

  • 即使队列中存在未被确认(unacked)的消息,channel.waitForConfirms始终打印"success";
  • 使用publish方法的回调时,无论消息是否被消费者确认,err参数始终为null,ok参数始终为undefined。

生产者代码(publisher.js)

var amqp = require('amqplib/callback_api');

amqp.connect('amqp://localhost:5672', function(error0, connection) {
  if (error0) {
    throw error0;
  }
  connection.createConfirmChannel(function(error1, channel) {
    if (error1) {
      throw error1;
    }
    var msg = 'Hello world';
    channel.assertExchange('exchange2', 'fanout', {
        durable: false
      });
    channel.publish('exchange2','', Buffer.from(msg), {},publishCallback);
    channel.waitForConfirms(function(err) {if (err) console.log(err); else console.log('success');})
    console.log(" [x] Sent %s", msg);
  });
});

function publishCallback(err, ok) {
    if (err !== null) {
      console.error('Error:', err);
    } else {
      console.log('Message published successfully:', ok);
    }
  }

消费者代码(receive.js)

var amqp = require('amqplib/callback_api');

amqp.connect('amqp://localhost:5672', function(error0, connection) {
  if (error0) {
    throw error0;
  }
  connection.createChannel(function(error1, channel) {
    if (error1) {
      throw error1;
    }
    var queue = 'hello';

    channel.assertExchange('exchange2', 'fanout', {
        durable: false
      });
  
      channel.assertQueue('', {
        exclusive: true
      }, function(error2, q) {
        if (error2) {
          throw error2;
        }
        console.log(" [*] Waiting for messages in %s. To exit press CTRL+C", q.queue);
        channel.bindQueue(q.queue, 'exchange2', '');
  
        channel.consume(q.queue, function(msg) {
          if(msg.content) {
              console.log(" [x] %s", msg.content.toString());
            //   channel.nack(msg, false, false)
            }
        }, {
          noAck: false
        });
      });
    });
});

问题原因分析

你混淆了RabbitMQ服务器的消息接收确认和消费者的消息处理确认这两个完全独立的机制:

  1. channel.waitForConfirms的作用:这个方法仅负责等待RabbitMQ服务器确认「所有已发布的消息已被服务器接收并完成路由(比如投递到匹配队列)」,和消费者是否确认消息没有任何关系。只要服务器成功接收并路由消息,不管消息是在队列等待消费、还是被消费者获取但未ack,waitForConfirms都会返回成功。
  2. publish回调的作用:在ConfirmChannel中,这个回调是服务器确认接收单条消息时触发的通知,同样只代表服务器已拿到消息,不涉及消费者的处理状态。另外ok参数在amqplib的callback_api设计中本身就不会返回有效值,这是库的特性,不是代码问题。

解决方案:自定义消费者确认反馈机制

RabbitMQ本身不会主动将消费者的ack通知给生产者,要实现需求需要自行构建回调逻辑:

  1. 生产者发送业务消息时,附带唯一消息ID,同时声明并监听一个回调队列;
  2. 消费者处理完业务消息并ack后,向回调队列发送携带原消息ID的确认消息;
  3. 生产者通过回调队列收到确认消息,即可判定对应业务消息已被消费者处理完成。

改进后的生产者代码

var amqp = require('amqplib/callback_api');
const uuid = require('uuid'); // 需先安装:npm install uuid

amqp.connect('amqp://localhost:5672', function(error0, connection) {
  if (error0) throw error0;

  connection.createChannel(function(error1, channel) {
    if (error1) throw error1;

    const exchange = 'exchange2';
    const callbackQueue = 'consumer_confirm_queue';
    const msgId = uuid.v4();
    const msg = 'Hello world';

    // 声明回调队列并监听确认消息
    channel.assertQueue(callbackQueue, { durable: false });
    channel.consume(callbackQueue, function(confirmMsg) {
      if (confirmMsg.content.toString() === msgId) {
        console.log(`[x] 消息 ${msgId} 已被消费者确认`);
        channel.ack(confirmMsg);
      }
    });

    channel.assertExchange(exchange, 'fanout', { durable: false });
    // 发送业务消息时附带唯一ID
    channel.publish(exchange, '', Buffer.from(msg), { messageId: msgId });
    console.log(" [x] Sent %s (ID: %s)", msg, msgId);
  });
});

改进后的消费者代码

var amqp = require('amqplib/callback_api');

amqp.connect('amqp://localhost:5672', function(error0, connection) {
  if (error0) throw error0;

  connection.createChannel(function(error1, channel) {
    if (error1) throw error1;

    const exchange = 'exchange2';
    const callbackQueue = 'consumer_confirm_queue';

    channel.assertExchange(exchange, 'fanout', { durable: false });
    channel.assertQueue('', { exclusive: true }, function(error2, q) {
      if (error2) throw error2;

      console.log(" [*] Waiting for messages in %s. To exit press CTRL+C", q.queue);
      channel.bindQueue(q.queue, exchange, '');

      channel.consume(q.queue, function(msg) {
        if (msg.content) {
          console.log(" [x] Received %s", msg.content.toString());
          // 这里添加业务逻辑处理代码...
          
          // 向回调队列发送确认消息,携带原消息ID
          channel.publish('', callbackQueue, Buffer.from(msg.properties.messageId));
          // 确认业务消息已处理完成
          channel.ack(msg);
        }
      }, { noAck: false });
    });
  });
});

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 01:40:25