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

FastApi+Aio-Pika重启Rabbit Broker后无法绑定队列消费问题

FastAPI + AioPika 重连RabbitMQ集群后无法绑定队列消费

我使用FastAPI(0.78)结合AioPika(9.0.5)开发消息消费服务,部署在AWS K8s环境中,对接3节点RabbitMQ集群实现负载均衡。AWS每1-2周会通过逐个下线旧实例更新后重新加入集群的方式维护Broker,当最后一个旧实例下线后,应用能成功重连新Broker实例,但无法绑定队列进行消费。

启动代码(main.py)

@app.on_event('startup')
async def startup() -> None:
    """Execute messages consumption on the application startup."""
    loop = asyncio.get_event_loop()
    task = loop.create_task(consume(loop))

    await task

消费逻辑代码(tasks.py)

async def consume(loop: BaseEventLoop) -> Connection:
    """Consume messages from RabbitMQ.

    Args:
        loop: Event loop.

    Returns:
        Aio-pika connection.
    """
    logger.info("Starting messages consumption...")

    conn: Connection = await connect_robust(settings.rabbitmq_url, loop=loop)
    channel = await conn.channel()
    exchange = await channel.get_exchange(settings.rabbitmq_exchange)
    queue_bridge = await channel.declare_queue('queue_name')

    await queue_bridge.bind(exchange, 'queue_name')
    await queue_bridge.consume(on_message)

    return conn


async def on_message(message: IncomingMessage) -> None:
    """Process incoming message with report.

    Args:
        message: Message received from the RabbitMQ.
    """
    async with message.process(ignore_processed=True):
        message_body = json.loads(message.body.decode('utf-8'))

        logger.info("Received message on %s", message.routing_key, extra={
            'body': message_body,
        })

        if message.routing_key == 'routing.key.name':
            # Do something with message payload
            pass

应用日志

Unexpected connection close from remote "amqps:SOMETHING", Connection.Close(reply_code=320, reply_text="CONNECTION_FORCED - broker forced connection closure with reason 'shutdown'")
NoneType: None
Unexpected connection close from remote "amqps:SOMETHING", Connection.Close(reply_code=320, reply_text="CONNECTION_FORCED - broker forced connection closure with reason 'shutdown'")
NoneType: None
Connection attempt to "amqps:SOMETHING" failed: [Errno 111] Connect call failed (SOME IP AND PORT). Reconnecting after 5 seconds.
Connection attempt to "amqps:SOMETHING" failed: [Errno 111] Connect call failed (SOME IP AND PORT). Reconnecting after 5 seconds.
Connection attempt to "amqps:SOMETHING" failed: [Errno 111] Connect call failed (SOME IP AND PORT). Reconnecting after 5 seconds.
Connection attempt to "amqps:SOMETHING" failed: [Errno 111] Connect call failed (SOME IP AND PORT). Reconnecting after 5 seconds.
Connection attempt to "amqps:SOMETHING" failed: [Errno 111] Connect call failed (SOME IP AND PORT). Reconnecting after 5 seconds.
Connection attempt to "amqps:SOMETHING" failed: [Errno 111] Connect call failed (SOME IP AND PORT). Reconnecting after 5 seconds.

按照Aio-Pika文档说明,connect_robust应自动重建连接、交换器和队列,但当前重连成功后无消费者运行,查阅相关修复记录仍未解决问题,寻求帮助。


问题分析与解决方案

核心问题

connect_robust仅负责自动重建连接和通道,但队列声明、绑定规则、消费者注册这些业务逻辑不会自动重复执行。同时原启动代码用await task阻塞FastAPI启动流程,消费任务中断后无法自动重启。

修复步骤

1. 重构消费逻辑,支持重连后自动恢复

将队列、绑定、消费者注册逻辑封装为独立函数,利用Aio-Pika连接的on_reconnect事件触发重连后的恢复操作:
修改tasks.py:

async def setup_consumption(channel: Channel):
    """Setup queue, binding and consumer on a channel."""
    exchange = await channel.get_exchange(settings.rabbitmq_exchange)
    queue_bridge = await channel.declare_queue('queue_name')
    await queue_bridge.bind(exchange, 'queue_name')
    await queue_bridge.consume(on_message)
    logger.info("Consumption setup completed successfully")

async def consume(loop: BaseEventLoop) -> None:
    logger.info("Starting messages consumption...")
    conn: Connection = await connect_robust(settings.rabbitmq_url, loop=loop)

    # 首次启动初始化消费
    channel = await conn.channel()
    await setup_consumption(channel)

    # 重连时自动恢复消费配置
    @conn.on_reconnect
    async def on_reconnect(connection: Connection):
        logger.info("Reconnected to RabbitMQ, restoring consumption...")
        new_channel = await connection.channel()
        await setup_consumption(new_channel)

    # 保持任务持续运行,避免退出
    await conn.close_event.wait()

2. 修复FastAPI启动阻塞问题

启动时将消费任务作为后台任务运行,避免阻塞FastAPI服务启动:
修改main.py:

@app.on_event('startup')
async def startup() -> None:
    """Execute messages consumption on the application startup."""
    loop = asyncio.get_event_loop()
    # 创建后台任务,不阻塞服务启动
    loop.create_task(consume(loop))

3. 额外优化建议

  • 确保队列/交换器持久化:如果队列或交换器是非持久化的,Broker重启后会丢失,重连时需重新声明。修改声明参数:
    # 声明持久化队列
    queue_bridge = await channel.declare_queue('queue_name', durable=True)
    # 获取持久化交换器(需确保交换器创建时已设置durable=True)
    exchange = await channel.get_exchange(settings.rabbitmq_exchange, durable=True)
    
  • 检查DNS解析:AWS维护后新实例IP可能变更,确保K8s环境能正确解析RabbitMQ集群域名,避免连接旧IP。
  • 升级AioPika版本:9.0.5为较旧版本,新版本修复了更多重连相关bug,建议升级到兼容的最新稳定版。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 06:47:14