RabbitMQ队列消息堆积问题排查求助
RabbitMQ队列频繁堆积问题排查求助
我使用RabbitMQ已有一段时间,近期队列频繁出现消息堆积问题。部署环境为AWS上的非托管RabbitMQ容器,Node.js后端通过自定义模块,基于amqplib和amqp-connection-manager两个npm包连接RabbitMQ服务。
每隔几周就会出现队列消息堆积且无消费者运行的情况,尽管后端初始化时已配置消费者。怀疑问题与ack逻辑(重复确认dup ack)相关,多次遇到错误码406,但无法准确定位原因。
可以确定的是,堆积消息的redelivered属性均为true,示例消息如下:
[ { "payload_bytes": 309, "redelivered": true, "exchange": "logExchange", "routing_key": "data", "message_count": 191, "properties": { "headers": {} }, "payload": "敏感内容已隐藏", "payload_encoding": "string" }, { "payload_bytes": 309, "redelivered": true, "exchange": "logExchange", "routing_key": "data", "message_count": 190, "properties": { "headers": {} }, "payload": "敏感内容已隐藏", "payload_encoding": "string" }, { "payload_bytes": 309, "redelivered": true, "exchange": "logExchange", "routing_key": "data", "message_count": 189, "properties": { "headers": {} }, "payload": "敏感内容已隐藏", "payload_encoding": "string" }, { "payload_bytes": 309, "redelivered": true, "exchange": "logExchange", "routing_key": "data", "message_count": 188, "properties": { "headers": {} }, "payload": "敏感内容已隐藏", "payload_encoding": "string" }, { "payload_bytes": 309, "redelivered": true, "exchange": "logExchange", "routing_key": "data", "message_count": 187, "properties": { "headers": {} }, "payload": "敏感内容已隐藏", "payload_encoding": "string" }, { "payload_bytes": 309, "redelivered": true, "exchange": "logExchange", "routing_key": "data", "message_count": 186, "properties": { "headers": {} }, "payload": "敏感内容已隐藏", "payload_encoding": "string" }, { "payload_bytes": 309, "redelivered": true, "exchange": "logExchange", "routing_key": "data", "message_count": 185, "properties": { "headers": {} }, "payload": "敏感内容已隐藏", "payload_encoding": "string" }, { "payload_bytes": 309, "redelivered": true, "exchange": "logExchange", "routing_key": "data", "message_count": 184, "properties": { "headers": {} }, "payload": "敏感内容已隐藏", "payload_encoding": "string" } ]
生产者类代码
import { ContactEvents, DataEvents, SocketAgent, CentralEvents, ExposureEvents, } from "./events"; import { exchanges } from "./constants"; import { Logger } from "@private/logger"; import amqp, { Channel, ChannelWrapper, AmqpConnectionManager, } from "amqp-connection-manager"; import { checkEnv } from "@private/shared-utilities"; import { SocketUser } from "./events/SocketUser"; import { IAMServiceEvents } from "./events/IAMServiceEvents"; import { NotificationEvents } from "./events/notification_events"; const logger = new Logger("AMQP Producer"); const exchangeName = exchanges.log; export class Producer { connection: AmqpConnectionManager; channel: ChannelWrapper; contact: ContactEvents; central: CentralEvents; kbrainData: DataEvents; socketAgent: SocketAgent; socketUser: SocketUser; iamService: IAMServiceEvents; exposure: ExposureEvents; notificationEvents: NotificationEvents; constructor() { checkEnv(["RABBITMQ_URL"]); this.connection = amqp.connect(process.env.RABBITMQ_URL); this.connection.on("connect", ({ connection, url }) => { logger.logInfo( `Producer connection established, URL: ${url}, Connection: ${connection.connection.serverProperties}` ); }); this.connection.on("disconnect", (connectionError) => { logger.logError( `Producer connection disconnected: ${JSON.stringify( connectionError.err )}` ); }); this.connection.on("connectFailed", (error) => { logger.logError( `Connection to ${error.url} failed: ${JSON.stringify(error.err)}` ); }); this.connection.on("error", (error, info) => { logger.logError(`Error in producer connection: ${JSON.stringify(error)}`); }); this.connection.on("blocked", ({ reason }) => { logger.logError(`Producer connection blocked: ${JSON.stringify(reason)}`); }); this.connection.on("unblocked", () => { logger.logError(`Producer connection unblocked`); }); this.channel = this.connection.createChannel({ json: true, setup: (channel: Channel) => { return channel.assertExchange(exchangeName, "direct"); }, }); this.channel.on("connect", () => { logger.logInfo("Producer channel established"); }); this.channel.on("error", (error, info) => { logger.logError( `Error in producer channel: ${JSON.stringify( error )}, Info: ${JSON.stringify(info)}` ); }); this.channel.on("close", () => { logger.logInfo("Producer channel closed"); }); this.contact = new ContactEvents(this); this.central = new CentralEvents(this); this.kbrainData = new DataEvents(this); this.socketAgent = new SocketAgent(this); this.socketUser = new SocketUser(this); this.iamService = new IAMServiceEvents(this); this.exposure = new ExposureEvents(this); this.notificationEvents = new NotificationEvents(this); } async publishMessage(routingKey: string, message: Object) { try { await this.channel.publish(exchangeName, routingKey, { logType: routingKey, message: message, dateTime: new Date(), }); logger.logInfo( `The message ${JSON.stringify( message )} is send to exchange ${exchangeName} for queue ${routingKey}` ); } catch (error) { logger.logError( `Error: ${error}, while publishing message: ${JSON.stringify(message)}` ); } } public checkConnection() { return this.connection.isConnected(); } } export const producerInstance = new Producer();
消费者类代码
import { Logger } from "@private/logger"; import { exchanges } from "./constants"; import amqp, { AmqpConnectionManager, Channel, ChannelWrapper, } from "amqp-connection-manager"; import { ConsumeMessage } from "amqplib"; import { checkEnv } from "@private/shared-utilities"; const logger = new Logger("AMQP Consumer"); const exchangeName = exchanges.log; export class Consumer { private static instance: Consumer | null = null; private connection: AmqpConnectionManager; private channel: ChannelWrapper; private handleMessageCB: Function; private constructor(queueName: string, callback: Function) { checkEnv(["RABBITMQ_URL"]); this.connection = amqp.connect(process.env.RABBITMQ_URL); this.setupConnectionListeners(); this.handleMessageCB = callback; this.channel = this.setupChannel(queueName); } public static getInstance(queueName: string, callback: Function): Consumer { if (!Consumer.instance) { Consumer.instance = new Consumer(queueName, callback); } return Consumer.instance; } public isConnectionAlive(): boolean { return this.connection.isConnected(); } private setupConnectionListeners(): void { this.connection.on("connect", ({ connection, url }) => { logger.logInfo( `Consumer connection established, URL: ${url}, Connection: ${connection.connection.serverProperties}` ); }); this.connection.on("disconnect", (connectionError) => { logger.logError( `Consumer connection disconnected: ${connectionError.err}` ); }); this.connection.on("connectFailed", (error) => { logger.logError(`Connection to ${error.url} failed: ${error.err}`); }); this.connection.on("error", (error) => { logger.logError(`Error in consumer connection: ${error}`); }); this.connection.on("blocked", ({ reason }) => { logger.logError(`Consumer connection blocked: ${reason}`); }); this.connection.on("unblocked", () => { logger.logError(`Consumer connection unblocked`); }); } private setupChannel(queueName: string): ChannelWrapper { const channel = this.connection.createChannel({ json: true, setup: async (channel: Channel) => { try { logger.logInfo(`Asserting queue: ${queueName}`); await channel.assertQueue(queueName, { durable: true, exclusive: false, }); logger.logInfo(`Asserting exchange: ${exchangeName}`); await channel.assertExchange(exchangeName, "direct"); logger.logInfo( `Binding queue: ${queueName} to exchange: ${exchangeName}` ); await channel.bindQueue(queueName, exchangeName, queueName); logger.logInfo(`Starting consumer for queue: ${queueName}`); await channel.consume(queueName, this.onMessage); } catch (error) { logger.logError(`Error setting up channel: ${error}`); } }, }); channel.on("connect", () => { logger.logInfo("Consumer channel established"); }); channel.on("error", (error, info) => { logger.logError(`Error in consumer channel: ${error}, Info: ${info}`); }); channel.on("close", () => { logger.logInfo("Consumer channel closed"); }); return channel; } private onMessage = async (data: ConsumeMessage | null) => { logger.logInfo( `Consumer channel received message: ${data ? data.content.toString() : "null"}` ); if (!data) return logger.logError("Empty message"); let message = JSON.parse(data.content.toString()).message; logger.logInfo("Consumer got event: " + message.eventType); try { const message = JSON.parse(data.content.toString()).message; logger.logInfo(`Consumer got event: ${message.eventType}`); if (this.handleMessageCB) { await this.handleMessageCB(message); this.channel?.ack(data); logger.logInfo("Message acknowledged"); } else { logger.logError("Message handler callback is not set"); } } catch (e) { logger.logError( `Error handling message: ${message.eventType} in consumer: ${JSON.stringify(e)}` ); } }; }
后端模块使用方式
import { BadRequestError } from "@private/errors"; import { Consumer, queues } from "@private/event-bus"; import { handleAttackDataEvents } from "./events"; export let consumer: Consumer | null = null; const initAMQP = async () => { try { consumer = Consumer.getInstance(queues.DATA, handleAttackDataEvents); } catch (error) { if (error instanceof Error) { throw new BadRequestError(error.message); } } }; export { initAMQP };
const main = async () => { try { // ... 其他初始化代码 await initAMQP(); } catch (error) { if (error instanceof Error) { return logger.logError(error.message); } } };
已尝试的排查操作:
- 对服务进行压测,未复现问题
- 手动断开AMQP连接/通道,服务能正常重连,未出现堆积问题
恳请帮忙排查队列堆积的根本原因。
内容的提问来源于stack exchange,提问作者lior shein
相关产品推荐
相关产品推荐

