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

如何确保FastAPI中aio_pika消费者任务长期稳定运行?

FastAPI中aio_pika消费者长期运行停止问题的分析与解决

可能原因

  • 未监控的任务崩溃:启动的消费任务没有异常回调,崩溃后无日志记录,无法定位原因
  • AMQP连接重连不彻底:虽然connect_robust会自动重连,但队列消费在重连后可能未重新初始化,导致消息无法接收
  • 发布订阅链失效:MyManager中的waiter Future若出现未处理异常,会导致订阅迭代器卡住,无法传递消息
  • 未处理的通道/队列异常:消费过程中出现的未捕获异常可能导致消费任务静默终止

解决方案

1. 监控消费任务异常

为启动的消费任务添加异常回调,捕获并记录任务崩溃的详细信息:

my_manager = asyncio.Future()

@app.on_event("startup")
async def on_startup():
    manager = MyManager()
    my_manager.set_result(manager)
    consume_task = asyncio.create_task(manager.run())
    
    def handle_task_exception(task):
        try:
            task.result()
        except Exception as e:
            logger.critical("消费任务意外崩溃", exc_info=True)
    
    consume_task.add_done_callback(handle_task_exception)

2. 实现可靠的重连与消费逻辑

修改run方法,确保连接断开后自动重连并重新初始化队列消费,全程捕获异常并记录:

async def run(self):
    while True:
        try:
            connection = await connect_robust(
                settings.amqp_url,
                loop=asyncio.get_running_loop()
            )
            connection.add_close_callback(lambda _: logger.warning("AMQP连接已关闭"))
            
            channel = await connection.channel()
            # 声明持久化队列(根据业务需求调整参数)
            my_queue = await channel.declare_queue('my-queue', durable=True)
            
            logger.info("成功启动队列消费: my-queue")
            await my_queue.consume(self.on_message)
            
            # 等待连接关闭事件,而非无限阻塞
            await connection.wait_closed()
            logger.warning("AMQP连接关闭,即将尝试重连")
        except Exception as e:
            logger.error("AMQP连接或消费过程出错", exc_info=True)
            # 重连前等待5秒,避免频繁重试消耗资源
            await asyncio.sleep(5)

3. 修复发布订阅逻辑的健壮性

增强subscribe方法的异常处理,避免迭代器因Future异常卡住;同时优化publish方法,确保Future操作的有效性:

async def subscribe(self):
    waiter = self.waiter
    while True:
        try:
            value, waiter = await waiter
            yield value
        except Exception as e:
            logger.error("订阅迭代器处理出错", exc_info=True)
            # 重置waiter,恢复订阅链
            self.waiter = asyncio.Future()
            waiter = self.waiter

def publish(self, value):
    # 确保当前waiter未完成时才进行替换操作
    if not self.waiter.done():
        current_waiter, self.waiter = self.waiter, asyncio.Future()
        current_waiter.set_result((value, self.waiter))
    else:
        # 若当前waiter已完成,直接初始化新的Future并设置结果
        self.waiter = asyncio.Future()
        self.waiter.set_result((value, self.waiter))

4. 优化SSE路由的稳定性

简化SSE事件发布逻辑,添加客户端断开检测和异常处理:

@router.get('/items')
async def items_stream(request: Request):
    async def event_publisher():
        manager = await my_manager
        try:
            async for event in manager:
                # 提前检测客户端是否断开,避免无效消息发送
                if await request.is_disconnected():
                    break
                yield dict(data=event)
        except asyncio.CancelledError:
            logger.info(f"客户端断开连接: {request.client}")
            raise
        except Exception as e:
            logger.error("SSE流处理出错", exc_info=True)
            yield dict(data="服务端异常", event="error")
    
    return EventSourceResponse(event_publisher())

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.11 01:53:15