如何在aio-pika中处理RabbitMQ意外连接关闭:禁用重连抛异常
处理RabbitMQ意外连接关闭:禁用自动重连并触发异常关闭服务
aio-pika的connect_robust是专为自动重连设计的,要实现“连接断开时触发异常并关闭服务”的需求,需要替换成普通连接方法并监听连接关闭事件,具体方案如下:
核心修改点
- 替换
connect_robust为aio_pika.connect:普通连接方法不会自动尝试重连,连接意外断开时会直接抛出异常。 - 注册连接关闭回调:在连接关闭时主动抛出异常,触发服务终止流程。
修改后的完整代码
class AMQPHandler: def __init__(self) -> None: # 注意类型注解切换为普通连接类型 self.connection: AbstractConnection | None = None self.channel: AbstractChannel | None = None async def _on_connection_closed(self, connection, exc): # 连接因异常关闭时抛出错误,触发服务终止 if exc: raise RuntimeError("RabbitMQ连接意外中断") from exc async def init(self) -> None: import settings import asyncio logger.info("Initializing AMQP handler") config = settings.BaseMessageBrokerSettings # 使用普通connect替代自动重连的connect_robust connection = await aio_pika.connect( config.get_dsn(), loop=asyncio.get_event_loop(), timeout=config.CONNECTION_TIMEOUT, ) # 注册连接关闭回调 connection.add_close_callback(self._on_connection_closed) channel = await connection.channel() await channel.set_qos(prefetch_count=10) exchange = await channel.declare_exchange( config.EXCHANGE_NAME, config.EXCHANGE_TYPE, auto_delete=config.EXCHANGE_AUTO_DELETE, durable=True, ) for key in config.BINDING_KEYS: q_name = ( f"{key}.{config.PREFIX_BINDING_KEYS}" if config.PREFIX_BINDING_KEYS else key ) queue = await channel.declare_queue(name=q_name, durable=True) await queue.bind(exchange, q_name) await queue.consume(self.handle_message) logger.info("Queue declared", extra={"queue": q_name}) self.connection = connection self.channel = channel logger.info("AMQP handler initialized")
主程序异常处理示例
在主逻辑中捕获连接异常,执行服务关闭流程:
async def main(): amqp_handler = AMQPHandler() try: await amqp_handler.init() # 保持服务运行的阻塞逻辑(比如等待终止信号) await asyncio.Event().wait() except RuntimeError as e: logger.error(f"RabbitMQ连接异常触发服务终止: {e}") # 这里可以添加其他资源的清理逻辑 exit(1) if __name__ == "__main__": asyncio.run(main())
内容的提问来源于stack exchange,提问作者visiwat303
相关产品推荐
相关产品推荐

