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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 16:37:32