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

使用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:

  1. 修改ConnectionManager的__init__:
def __init__(self) -> None:
    if not hasattr(self, "initialized"):
        self.active_conversations: dict[str, Connection] = {}
        self.lock = None
        self.initialized = True
  1. 在endpoint.py中添加启动事件:
@router.on_event("startup")
async def init_lock():
    manager.lock = asyncio.Lock()

额外优化建议

  1. 不要在锁块内执行耗时操作:原代码中disconnect的锁块里有await asyncio.sleep(1),这会导致锁被长时间占用,其他连接请求会被阻塞。建议把sleep和WebSocket关闭操作移到锁块外面,只在锁内处理字典的修改。
  2. 简化单例实现:可以用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 16:52:34