使用amqplib实现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

