FastAPI WebSocket端点异步任务无法正常退出问题求助
问题分析
你的FastAPI WebSocket服务与Redis Pub/Sub集成运行正常,但在应用退出或重载时,异步任务无法正常终止,导致Uvicorn卡在Waiting for background tasks to complete阶段。核心原因是:
reader函数的无限循环未正确响应异步任务取消信号(CancelledError)- 原代码中
async_timeout.timeout和asyncio.sleep的组合,导致任务取消信号无法及时触发循环退出 - 存在参数名不匹配的潜在问题(
station_idvsid)
解决方案
以下是修复后的完整代码,解决任务取消和应用关闭问题:
1. 修复reader函数,正确响应取消信号
使用aioredis.PubSub.get_message自带的timeout参数替代async_timeout,移除不必要的sleep,并显式处理CancelledError:
async def reader(channel: aioredis.client.PubSub, websocket: WebSocket): try: while True: # 用get_message自带的timeout,避免额外的超时包装 message = await channel.get_message( ignore_subscribe_messages=True, timeout=1 # 每隔1秒检查一次是否有消息或取消信号 ) if message and message['type'] == 'message': await websocket.send_text(message['data'].decode('utf-8')) except asyncio.CancelledError: # 捕获取消信号,正常退出循环并向上传递 raise except Exception as e: print(f"Reader task error: {str(e)}")
2. 优化WebSocket端点,管理异步任务生命周期
使用asyncio.create_task启动reader任务,确保在关闭时能主动取消并等待任务完成,同时修复参数名不匹配问题:
@router.websocket("/api/ws/{station_id}") async def websocket_endpoint(websocket: WebSocket, station_id: str, redis: aioredis.Redis = Depends(get_redis)): await websocket.accept() channel_name = f'updates:{station_id}' # 修复参数名错误:用station_id替代id pubsub = redis.pubsub() await pubsub.subscribe(channel_name) reader_task = None try: # 创建独立任务运行reader,避免阻塞主端点逻辑 reader_task = asyncio.create_task(reader(pubsub, websocket)) # 等待客户端断开连接或任务被取消 await websocket.receive() except WebSocketDisconnect: print("WebSocket disconnected by the client.") except asyncio.CancelledError: print("WebSocket connection cancelled.") except Exception as e: print(f'WebSocket error: {str(e)}') finally: # 主动取消reader任务并等待其完成 if reader_task: reader_task.cancel() try: await reader_task except asyncio.CancelledError: pass # 清理Redis Pub/Sub资源 await pubsub.unsubscribe(channel_name) await pubsub.close() # 确保WebSocket连接关闭 await websocket.close()
3. 额外优化建议
- Redis连接池管理:使用FastAPI的
lifespan事件统一管理Redis连接池,避免每次WebSocket请求创建新连接:@asynccontextmanager async def lifespan(app: FastAPI): # 启动时创建Redis连接池 app.state.redis = aioredis.from_url("redis://localhost") yield # 关闭时清理连接池 await app.state.redis.close() app = FastAPI(lifespan=lifespan) # 修改get_redis依赖 def get_redis(app: FastAPI = Depends(get_app)) -> aioredis.Redis: return app.state.redis - 避免冗余日志:将
print替换为标准日志库(如logging),便于生产环境排查问题。
修复原理
- 移除
asyncio.sleep和async_timeout的组合,让get_message的内置timeout定期释放循环,确保取消信号能及时被捕获 - 使用
asyncio.create_task分离reader任务,主端点逻辑专注于等待连接状态变化,关闭时主动取消reader任务并等待其完成 - 修复参数名错误,确保订阅的Redis频道与预期一致
内容的提问来源于stack exchange,提问作者csimpler
相关产品推荐
相关产品推荐

