如何确保FastAPI中aio_pika消费者任务长期稳定运行?
FastAPI中aio_pika消费者长期运行停止问题的分析与解决
可能原因
- 未监控的任务崩溃:启动的消费任务没有异常回调,崩溃后无日志记录,无法定位原因
- AMQP连接重连不彻底:虽然
connect_robust会自动重连,但队列消费在重连后可能未重新初始化,导致消息无法接收 - 发布订阅链失效:
MyManager中的waiterFuture若出现未处理异常,会导致订阅迭代器卡住,无法传递消息 - 未处理的通道/队列异常:消费过程中出现的未捕获异常可能导致消费任务静默终止
解决方案
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
相关产品推荐
相关产品推荐

