生产环境遇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
相关产品推荐
相关产品推荐

