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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 12:25:57