使用FastAPI+asyncio单例管理WebSocket时遇Future跨事件循环错误
问题:FastAPI+asyncio单例连接管理器用Lock触发"Future pending attached to a different loop"错误
错误日志
cb=[WebSocketProtocol.on_task_complete()] got Future
attached to a different loop
代码结构
- endpoint.py:实现WebSocket连接管理的
ConnectionManager单例类,通过asyncio.Lock()保护active_conversations字典 - main.py:FastAPI应用初始化及路由挂载
- test.js:k6编写的WebSocket连接负载测试脚本
代码片段
endpoint.py
import logging import asyncio from asyncio import Lock from dataclasses import dataclass import json from fastapi import APIRouter, HTTPException, WebSocket, WebSocketDisconnect from fastapi.param_functions import Query router = APIRouter() logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] - %(message)s") logger = logging.getLogger(__name__) @dataclass class Connection: websocket: WebSocket class ConnectionManager: _instance = None def __new__(cls, *args, **kwargs): if not cls._instance: cls._instance = super(ConnectionManager, cls).__new__(cls) return cls._instance def __init__(self) -> None: if not hasattr(self, "initialized"): self.active_conversations: dict[str, Connection] = {} self.lock = Lock() self.initialized = True async def connect(self, websocket: WebSocket, session_id: str): async with self.lock: await websocket.accept() connection = Connection(websocket=websocket) self.active_conversations[session_id] = connection async def disconnect(self, session_id: str): async with self.lock: connection = self.active_conversations.pop(session_id, None) await asyncio.sleep(1) if not connection: return try: await connection.websocket.close() logger.info(f"WebSocket closed for session {session_id}") except Exception as e: logger.error(f"Error closing WebSocket for session {session_id}: {e}") manager = ConnectionManager() @router.websocket("/conversation/{session_id}") async def conversation(websocket: WebSocket, session_id: str, token: str = Query(...)): if not token: raise HTTPException(status_code=400, detail="Token is required!") try: await manager.connect(websocket=websocket, session_id=session_id) while True: data = await websocket.receive_text() message = json.loads(data) if message.get("type") == "ping": await websocket.send_text(json.dumps({"type": "pong"})) except Exception as e: logger.error(f"Error in WebSocket connection: {e}") finally: await manager.disconnect(session_id=session_id)
main.py
import os from fastapi import FastAPI from fastapi.middleware.cors import CORSMiddleware from jac.api.endpoints import router as endpoint_router app = FastAPI(title="Deepdive LLM") origins = ["*"] app.add_middleware( CORSMiddleware, allow_origins=origins, allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) app.include_router(endpoint_router) if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8080)
test.js
import ws from 'k6/ws'; import { check } from 'k6'; export const options = { scenarios: { warmup: { executor: 'per-vu-iterations', vus: 5, iterations: 1, maxDuration: '2m', }, }, }; export default function () { const url = 'ws://localhost:8080/conversation/session_' + __VU; const token = 'your_token'; const connection = `${url}?token=${token}`; const response = ws.connect(connection, {}, function (socket) { socket.on('open', function open() { console.log('connected'); socket.send(JSON.stringify({ type: 'ping', session_id: 'session_' + __VU, token: token })); socket.on('message', function (msg) { console.log('Message received: ', msg); check(msg, { 'is pong': (msg) => msg.includes('pong') }); socket.close(); }); socket.on('close', () => console.log('disconnected')); }); socket.on('error', function (e) { if (e.error() != 'websocket: close sent') { console.error('An unexpected error occured: ', e.error()); } }); }); }
问题原因
核心问题是单例初始化时机与事件循环不匹配:
模块导入时就创建了ConnectionManager实例(manager = ConnectionManager()),此时asyncio.Lock()会绑定到当前的临时事件循环。但FastAPI启动uvicorn时会创建全新的事件循环处理请求,导致锁对象和请求使用的事件循环不属于同一个,触发跨循环的Future错误。
解决方案
方案1:延迟初始化Lock(推荐)
不在类初始化时创建Lock,而是在第一次使用锁的时候初始化,确保Lock绑定到当前运行的事件循环:
修改ConnectionManager类:
class ConnectionManager: _instance = None def __new__(cls, *args, **kwargs): if not cls._instance: cls._instance = super(ConnectionManager, cls).__new__(cls) return cls._instance def __init__(self) -> None: if not hasattr(self, "initialized"): self.active_conversations: dict[str, Connection] = {} self.lock = None # 先不初始化Lock self.initialized = True async def connect(self, websocket: WebSocket, session_id: str): # 第一次调用时初始化Lock,绑定到当前事件循环 if self.lock is None: self.lock = asyncio.Lock() async with self.lock: await websocket.accept() connection = Connection(websocket=websocket) self.active_conversations[session_id] = connection async def disconnect(self, session_id: str): # 确保Lock已初始化 if self.lock is None: self.lock = asyncio.Lock() connection = None async with self.lock: connection = self.active_conversations.pop(session_id, None) # 把sleep移出锁块,避免占用锁资源 if connection: await asyncio.sleep(1) try: await connection.websocket.close() logger.info(f"WebSocket closed for session {session_id}") except Exception as e: logger.error(f"Error closing WebSocket for session {session_id}: {e}")
方案2:利用FastAPI启动事件初始化Lock
通过FastAPI的startup事件,在应用启动完成(事件循环已创建)后再初始化Lock:
- 修改
ConnectionManager的__init__:
def __init__(self) -> None: if not hasattr(self, "initialized"): self.active_conversations: dict[str, Connection] = {} self.lock = None self.initialized = True
- 在
endpoint.py中添加启动事件:
@router.on_event("startup") async def init_lock(): manager.lock = asyncio.Lock()
额外优化建议
- 不要在锁块内执行耗时操作:原代码中
disconnect的锁块里有await asyncio.sleep(1),这会导致锁被长时间占用,其他连接请求会被阻塞。建议把sleep和WebSocket关闭操作移到锁块外面,只在锁内处理字典的修改。 - 简化单例实现:可以用
functools.lru_cache实现更简洁的单例:
from functools import lru_cache class ConnectionManager: def __init__(self) -> None: self.active_conversations: dict[str, Connection] = {} self.lock = None @classmethod @lru_cache(maxsize=1) def get_instance(cls): return cls() # 使用时 manager = ConnectionManager.get_instance()
内容的提问来源于stack exchange,提问作者ibiscp
相关产品推荐
相关产品推荐

