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

生产环境遇RabbitMQ节点维护导致连接关闭异常求助

问题描述

生产环境中遇到如下异常:

pika.exceptions.ConnectionClosedByBroker: (320, 'CONNECTION_FORCED - Node was put into maintenance mode')

无法在实时环境复现该问题,但Pod会因该异常触发重启,现寻求修复方案。

当前运行代码

consumer.py

class MessageQConsumer:
    """
    消息队列消费者
    """

    def __init__(
            self,
            config: ConfigReader,
            exchange_name: str,
            queue_name: str,
            topic: str
    ):
        self.config = config
        self.logger = LoggerHandler(self.config).logger
        self.publisher = None
        self.topic = topic

        # exchange name
        self.exchange_name = exchange_name
        self.queue_name = queue_name

        self.connection = get_connection()

        # mq channel
        self.channel = self.connection.channel()

        # limit prefetch with a single message, no sense to take more requests from the queue
        self.channel.basic_qos(prefetch_count=1)

        # configuration of the exchange
        self.channel.exchange_declare(
            exchange=self.exchange_name, exchange_type=MessageQExchangesTypes.TOPIC,
            durable='true'
        )

        # declare queue
        temp = self.channel.queue_declare(self.queue_name, durable='true')
        self.queue = temp.method.queue

        # bind routing key in queue
        self.channel.queue_bind(
            exchange=self.exchange_name, queue=self.queue, routing_key=topic
        )

        # list of threads
        self.threads = []

        # retry_count
        self.retry_count = 1

    def run(self) -> None:
        """
        启动消费者
        """
        try:
            on_message_callback = functools.partial(
                self.on_message, args=(self.connection, self.threads))
            self.channel.basic_consume(
                queue=self.queue, on_message_callback=on_message_callback)
            self.channel.start_consuming()
        except KeyboardInterrupt:
            self.channel.stop_consuming()

        # Wait for all to complete.
        for thread in self.threads:
            thread.join()

        # connection close.
        self.connection.close()

    def on_message(self, channel, method, properties, body, args) -> None:
        """
        消息回调
        """
        (connection, threads) = args
        thread = threading.Thread(target=self.do_work, args=(
            connection, method, channel, body, properties))
        thread.start()
        threads.append(thread)

    def do_work(self, connection, method, channel, body, properties) -> None:
        """
        处理消息
        """
        self.message_queue_handle_request(body, properties)
        channel_connection = functools.partial(self.ack_message, channel, method)
        connection.add_callback_threadsafe(channel_connection)

    def ack_message(self, channel, method) -> None:
        """
        确认消息
        """
        if channel.is_open:
            channel.basic_ack(method.delivery_tag)
        else:
            self.logger.error(f"RabbitMq channel is closed !!")
            # Channel is already closed, so we can't ACK this message;

    def message_queue_handle_request(self, body, properties) -> None:
        """
        处理业务逻辑
        """
        data = {}
        self.logger.info(f"Initialize the Message Publisher")
        if not self.publisher:
            self.publisher = MessageQPublisher(
                self.config,
                self.exchange_name
            )

        try:
            # process internal algorithm
            pass
        except Exception as exception:
            pass

publisher.py

class MessageQPublisher:
    """
    消息队列发布者
    """

    def __init__(self, config: ConfigReader, exchange: str):
        self.config = config
        self.logger = LoggerHandler(self.config).logger
        self.exchange = exchange

    def publish(self, routing_key: str, body: str, header_retry_count=None) -> None:
        """
        发布消息
        """
        headers_value = {'retryCount': header_retry_count} if header_retry_count else {}

        # connection
        connection = get_connection(self.config)

        # setup channel
        channel = connection.channel()

        # declare the exchange that should be used
        channel.exchange_declare(exchange=self.exchange,
                                 exchange_type=MessageQExchangesTypes.TOPIC,
                                 durable='true')
        channel.basic_publish(
            self.exchange,
            routing_key,
            body.encode('utf8'),
            properties=pika.BasicProperties(
                # make message persistent
                delivery_mode=2,
                expiration='3600000',
                headers=headers_value
            )
        )

connection.py

def get_connection(config_data=config):
    """
    创建RabbitMQ连接
    """
    if config_data.mb_protocol == MessageQProtocol.AMQPS:
        ssl_context = ssl.SSLContext(ssl.PROTOCOL_TLSv1_2)
        return pika.BlockingConnection(
            pika.ConnectionParameters(
                host=config_data.mb_ip,
                port=config_data.mb_port,
                credentials=pika.PlainCredentials(
                    config_data.mb_user,
                    config_data.mb_password
                ),
                heartbeat=600,
                blocked_connection_timeout=300,
                ssl_options=pika.SSLOptions(context=ssl_context)))
    else:
        return pika.BlockingConnection(
            pika.ConnectionParameters(
                host=config_data.mb_ip,
                port=config_data.mb_port,
                credentials=pika.PlainCredentials(
                    config_data.mb_user,
                    config_data.mb_password),
                heartbeat=600,
                blocked_connection_timeout=300)
        )
修复方案

1. 实现消费者自动重连机制

当前消费者仅在初始化时创建一次连接,断开后无法自动恢复。需要添加重连逻辑,在连接因节点维护等原因断开后,自动重建连接、通道并恢复消费。

修改MessageQConsumer的核心逻辑:

import time
import pika

class MessageQConsumer:
    # 保留原有__init__方法,新增以下方法
    def _setup_connection_and_channel(self):
        """重新建立连接、通道并初始化队列/交换器"""
        self.connection = get_connection()
        self.channel = self.connection.channel()
        self.channel.basic_qos(prefetch_count=1)
        
        # 重新声明交换器、队列、绑定(幂等操作,不影响持久化资源)
        self.channel.exchange_declare(
            exchange=self.exchange_name, exchange_type=MessageQExchangesTypes.TOPIC,
            durable='true'
        )
        temp = self.channel.queue_declare(self.queue_name, durable='true')
        self.queue = temp.method.queue
        self.channel.queue_bind(
            exchange=self.exchange_name, queue=self.queue, routing_key=self.topic
        )

    def run(self) -> None:
        """修改run方法,添加循环重连逻辑"""
        while True:
            try:
                self._setup_connection_and_channel()
                on_message_callback = functools.partial(
                    self.on_message, args=(self.connection, self.threads))
                self.channel.basic_consume(
                    queue=self.queue, on_message_callback=on_message_callback)
                self.logger.info("消息消费者启动成功,开始监听队列")
                self.channel.start_consuming()
            except KeyboardInterrupt:
                self.logger.info("收到中断信号,停止消费")
                self.channel.stop_consuming()
                break
            except pika.exceptions.ConnectionClosedByBroker as e:
                self.logger.error(f"Broker强制关闭连接,5秒后重试: {str(e)}")
                time.sleep(5)
            except pika.exceptions.AMQPConnectionError as e:
                self.logger.error(f"AMQP连接异常,5秒后重试: {str(e)}")
                time.sleep(5)
            finally:
                # 等待所有处理线程结束
                for thread in self.threads:
                    thread.join()
                self.threads.clear()
                # 确保连接关闭
                if self.connection and not self.connection.is_closed:
                    self.connection.close()

2. 优化发布者连接管理与重试

当前发布者每次发送消息都创建新连接,不仅浪费资源,遇到节点维护时也会直接抛出异常。需要复用连接并添加重试逻辑:

class MessageQPublisher:
    def __init__(self, config: ConfigReader, exchange: str):
        self.config = config
        self.logger = LoggerHandler(self.config).logger
        self.exchange = exchange
        self.connection = None
        self.channel = None

    def _get_or_create_channel(self):
        """获取或创建有效通道,复用连接"""
        if not self.connection or self.connection.is_closed:
            self.connection = get_connection(self.config)
        if not self.channel or self.channel.is_closed:
            self.channel = self.connection.channel()
            self.channel.exchange_declare(exchange=self.exchange,
                                     exchange_type=MessageQExchangesTypes.TOPIC,
                                     durable='true')
        return self.channel

    def publish(self, routing_key: str, body: str, header_retry_count=None, retry_times=3) -> None:
        headers_value = {'retryCount': header_retry_count} if header_retry_count else {}
        for attempt in range(retry_times):
            try:
                channel = self._get_or_create_channel()
                channel.basic_publish(
                    self.exchange,
                    routing_key,
                    body.encode('utf8'),
                    properties=pika.BasicProperties(
                        delivery_mode=2,
                        expiration='3600000',
                        headers=headers_value
                    )
                )
                self.logger.info(f"消息发布成功,路由键: {routing_key}")
                return
            except (pika.exceptions.ConnectionClosedByBroker, pika.exceptions.AMQPConnectionError) as e:
                self.logger.error(f"第{attempt+1}次发布失败: {str(e)}")
                # 重置连接,下次重试重新创建
                self.connection = None
                self.channel = None
                time.sleep(2)
        self.logger.error(f"消息发布失败,已重试{retry_times}次,路由键: {routing_key}")

3. 完善异常捕获与消息ACK处理

在ack_message中添加异常捕获,避免因通道状态异常导致未处理崩溃;同时,未ACK的消息会在RabbitMQ连接恢复后自动重新投递,无需手动处理:

def ack_message(self, channel, method) -> None:
    """
    安全确认消息
    """
    try:
        if channel.is_open:
            channel.basic_ack(method.delivery_tag)
            self.logger.debug(f"消息已ACK,delivery_tag: {method.delivery_tag}")
        else:
            self.logger.error(f"通道已关闭,无法ACK消息,delivery_tag: {method.delivery_tag}")
            # RabbitMQ会自动重新投递未ACK的消息
    except Exception as e:
        self.logger.error(f"ACK消息时出错: {str(e)}")

4. 调整心跳参数

当前心跳间隔设置为600秒,过长的心跳会导致客户端无法及时感知连接断开。缩短心跳间隔至60秒,让客户端更快检测到连接状态变化:

修改connection.py中的heartbeat参数:

heartbeat=60,  # 从600改为60
blocked_connection_timeout=300
关键说明
  • CONNECTION_FORCED属于RabbitMQ集群运维场景下的正常断开(如节点维护、重启),客户端必须具备自动恢复能力才能避免Pod崩溃。
  • 重连时重新声明交换器、队列和绑定是幂等操作,不会影响已持久化的资源。
  • 发布者复用连接可大幅减少资源开销,重试逻辑能提升消息发布的可靠性。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 07:02:02