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

FastAPI WebSocket端点异步任务无法正常退出问题求助

问题分析

你的FastAPI WebSocket服务与Redis Pub/Sub集成运行正常,但在应用退出或重载时,异步任务无法正常终止,导致Uvicorn卡在Waiting for background tasks to complete阶段。核心原因是:

  • reader函数的无限循环未正确响应异步任务取消信号(CancelledError)
  • 原代码中async_timeout.timeout和asyncio.sleep的组合,导致任务取消信号无法及时触发循环退出
  • 存在参数名不匹配的潜在问题(station_id vs id)
解决方案

以下是修复后的完整代码,解决任务取消和应用关闭问题:

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 08:31:09