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

使用amqplib实现RabbitMQ无限消息消费时出现溢出问题求助

RabbitMQ消费溢出问题排查

问题背景

按照RabbitMQ官方文档,基于amqplib原生实现无限消费消息的方案,最初因setInterval重复调用导致溢出,修复重复调用问题后,运行数天仍出现溢出。

相关代码如下:

const ampq = require('ampqlib');
const JSON = require('JSON');
const winston = require('winston');

let connection = null;

const workerFunction = async (mesage) => {
  // A lot of work not requried for the example
  return result_message;
};

const createConnection = async () => {
  if (connection !== null) {
    // Winston logger setup hidden to simplify example
    logger.info('createConnection: connection already created');
    return connection;
  }

  connection = await amqp.connect({
    hostname: HOSTNAME,
    port: RABBITMQ_PORT,
    username: RABBITMQ_DEFAULT_USER,
    password: RABBITMQ_DEFAULT_PASS,
    vhost: RABBITMQ_DEFAULT_VHOST,
  });

  logger.info('createConnection: connection created');
  return connection;
}

const sendMessage = async (message) => {

  try {
    connection = await createConnection();
    const channel = await connection.createChannel();
    await channel.assertQueue('consumer', { durable: true });

    const msg = JSON.stringify(message, null, 4);
    channel.sendToQueue('consumer', Buffer.from(msg));
    await channel.close();
  } catch (e) {
    logger.error('sendMessage', e);
  }
};

const consumeMessage = async (message) => {
  try{
    logger.info('consumeMessage', message);

    const connection = await createConnection();
    const channel = await connection.createChannel();
    await channel.assertQueue('worker', { durable: true });

    await channel.consume(
      'worker',
      await workerFunction(message),
      { noAck: false }
    );
    logger.info('consumeMessage', message);
  } catch (e) {
    logger.error('consumeMessage', error);
  }
}

const consumeMessages = () => {
  const gson = JSON.parse(message.content.toString());

  setInterval(async () => {
    await consumeMessage();
  }, 200);
};

调用方式:

consumeMessages.catch(logger.error);

可能的溢出原因分析

  • setInterval轮询逻辑完全错误:RabbitMQ的channel.consume本身是事件驱动的持续监听模式,不需要用setInterval轮询。每隔200ms触发一次consumeMessage,如果workerFunction处理耗时超过200ms,会导致任务堆积——前一个任务未完成,新任务已启动,长期运行必然引发内存溢出。

  • Channel资源泄漏:每次调用consumeMessage都会创建新的Channel,且未正确关闭。Channel是RabbitMQ的有限资源,大量未释放的Channel会持续占用内存和连接资源,最终导致资源耗尽。同时,每次调用channel.consume都会新增一个消费者,队列会同时推送消息给多个消费者,加剧资源占用和重复消费问题。

  • channel.consume参数错误:channel.consume的第二个参数应为消息处理回调函数,但代码中传入了await workerFunction(message)的返回值(非函数类型),导致消费者逻辑完全失效,无法正确处理消息,进而引发异常堆积。

  • 消息未确认(ACK):代码设置了noAck: false,但未在消息处理完成后调用channel.ack(message)。RabbitMQ会认为消息未处理完成,持续保留在队列中,消费者会不断收到重复消息,堆积大量未确认消息占用内存。

  • 变量与调用逻辑错误:consumeMessages中使用未定义的message变量,会直接抛出错误;且consumeMessages是普通函数,调用consumeMessages.catch会触发报错,导致异常未被捕获,资源无法正常释放。

  • 异常场景下的资源清理缺失:当连接或Channel出现异常时,没有重连或资源清理逻辑,断开的连接/Channel资源无法释放,长期积累会引发资源溢出。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 12:25:18