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
相关产品推荐
相关产品推荐

