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

FastAPI多线程共享WebSocket:资源耗尽与SSLTransport错误排查

问题描述

我正在开发一款基于FastAPI的游戏后端服务,所有玩家通过主线程的WebSocket连接并与游戏交互。游戏包含多个分类,每个分类在独立线程中运行,且共享WebSocket作为通用参数。输入消息在主线程接收后,交由对应分类线程处理,再通过WebSocket发送输出消息。

起初运行正常,但服务器日志中持续出现以下错误:

Fatal error on SSL transport
protocol: <asyncio.sslproto.SSLProtocol object at 0x7f79804ba140>
transport: None
Traceback (most recent call last):
  File "/usr/local/lib/python3.10/asyncio/sslproto.py", line 703, in _process_write_backlog
del self._write_backlog[0]
IndexError: deque index out of range

该错误随机出现,甚至在服务器启动10分钟且无玩家连接时也会发生。错误未导致服务崩溃,但我怀疑这是代码流程存在问题的表现。同时,服务器内存持续上涨,数天后会崩溃。

以下是核心代码实现:

api.py

启动时生成所有游戏流程,监听新WebSocket连接/消息并转发至对应线程的游戏:

def create_gameflows(categories: list, websocket_manager) -> Dict[str, Gameflow]:
    return {cat_name: Gameflow(id=i+1, websocket_manager=websocket_manager, category=cat_name) 
            for i, cat_name in enumerate(categories)}

app = FastAPI()
manager = ConnectionManager()
lock = asyncio.Lock()

gameflows = create_gameflows(constants.CATEGORIES, manager)
for gameflow in gameflows.values():
    gameflow.start()

async def handle_websocket(websocket: WebSocket, client_id: int, category: str) -> None:
    """ Handle client incoming messages via websocket

    Parameters:
        websocket (WebSocket):
    """
    async with lock:
        await manager.connect(websocket, client_id, category)
        gf = gameflows[category]
        try:
            while True:
                data = await websocket.receive_json()
                if data['type'] == 'NEW_PLAYER':
                    while gf.game is None:
                        time.sleep(0.1)
                    user = User(**data['data'])
                    user.websocket_id = str(get_cookie(websocket))
                    try:
                        await gf.game.new_player(user, websocket, user.websocket_id, add_to_chat)
                    except :
                        await manager.send_error_message(websocket, user.websocket_id, data['data']['category'] )
                        raise

                elif data['type'] == 'ANSWER':
                    websocket_id = get_cookie(websocket)
                    await gf.game.user_answer(data, websocket_id, websocket)

        except WebSocketDisconnect:
            try:
                websocket_id = get_cookie(websocket)
            except (AttributeError, IndexError):
                websocket_id = str(client_id)
            await manager.close_connection(websocket_id, websocket_id, category)

@app.websocket("/api/ws/{client_id}/{category}")
async def websocket_endpoint(websocket: WebSocket, client_id:int, category: str):
    """ API route that redirect to the websocket endpoint.

    Args:
        websocket (WebSocket):
    """
    await handle_websocket(websocket, client_id, category)

Gameflow.py

Gameflow对象针对指定分类无限创建游戏,每个游戏会获取WebSocket管理器并向连接用户发送消息:

class Gameflow:
    def __init__(self, websocket_manager, category: str):
        self.websocket_manager = websocket_manager
        self.category = category
        self.game = None
        self.waiting_time = None
        self.loop = asyncio.new_event_loop()
        self.executor = ThreadPoolExecutor(max_workers=1)

    async def play_game(self):
        while True:
            try:
                self.game = Game(id=0, websocket_manager=self.websocket_manager, category=self.category)
                await self.game.play_game()
                while not self.game.is_over:
                    await asyncio.sleep(0.1)
                self.waiting_time = time.time()
                await asyncio.sleep(TIME_BETWEEN_GAME)
            except Exception as e:
                print(f"Error in {self.category} game: {e}")
                await asyncio.sleep(10)  # Wait before retrying

    def start(self):
        asyncio.set_event_loop(self.loop)
        self.loop.run_in_executor(self.executor, self.loop.run_forever)
        asyncio.run_coroutine_threadsafe(self.play_game(), self.loop)

ConnectionManager用于管理当前所有游戏的WebSocket连接,如需可提供其代码。

我尝试使用async with lock防止线程并发,但错误仍存在。我想了解该实现是否安全?是什么导致了SSL Transport错误?又是什么导致内存持续上涨?


问题分析与修复建议

1. 当前实现的安全性问题

你的代码存在跨线程操作asyncio对象的严重安全隐患:

  • FastAPI的WebSocket对象绑定在主线程的asyncio事件循环上,直接传递给Gameflow的独立线程循环调用,违反了asyncio的线程安全规则。
  • asyncio.Lock仅在同一事件循环内生效,无法保护跨线程的操作。
  • handle_websocket中使用time.sleep(0.1)会阻塞主线程事件循环,导致所有WebSocket连接无法处理新消息,这是异步代码的典型错误。

2. SSL Transport错误的原因

该错误由跨线程操作WebSocket底层传输对象导致:

  • WebSocket的SSL传输对象属于主线程事件循环,其他线程操作时会引发内部_write_backlog队列的并发冲突,比如空队列执行删除操作,触发IndexError。
  • 即使无玩家连接,Gameflow线程若操作共享的websocket_manager(持有WebSocket引用),也会引发事件循环资源清理的异常。

3. 内存持续上涨的原因

  • WebSocket对象泄漏:跨线程持有WebSocket引用,导致主线程无法回收已断开的连接实例,内存堆积无效对象。
  • Game对象泄漏:Gameflow.play_game循环创建新Game对象,若is_over判断逻辑异常或Game未释放内部资源(如WebSocket引用、任务),旧对象无法被GC回收。
  • 线程与事件循环资源泄漏:每个Gameflow创建独立线程池和事件循环,无正确关闭逻辑,长期运行会积累大量未释放的线程资源。

修复方案

(1)移除独立线程,改用asyncio任务

抛弃线程池,所有游戏逻辑在FastAPI主线程事件循环中运行,用asyncio.create_task启动游戏流程:

# 修改Gameflow.py
class Gameflow:
    def __init__(self, websocket_manager, category: str):
        self.websocket_manager = websocket_manager
        self.category = category
        self.game = None
        self.waiting_time = None
        self.task = None

    async def play_game(self):
        while True:
            try:
                self.game = Game(id=0, websocket_manager=self.websocket_manager, category=self.category)
                await self.game.play_game()
                while not self.game.is_over:
                    await asyncio.sleep(0.1)
                self.waiting_time = time.time()
                await asyncio.sleep(TIME_BETWEEN_GAME)
            except Exception as e:
                print(f"Error in {self.category} game: {e}")
                await asyncio.sleep(10)

    def start(self):
        self.task = asyncio.create_task(self.play_game())

(2)修复WebSocket线程安全问题

  • 绝对不要将WebSocket对象传递给其他线程,所有发送/接收操作必须在主线程事件循环中执行。
  • 通过websocket_manager封装异步发送方法,确保所有WebSocket操作在主线程完成。

(3)替换阻塞睡眠为异步睡眠

把handle_websocket中的time.sleep(0.1)改为await asyncio.sleep(0.1),避免阻塞主线程:

# 修改api.py的NEW_PLAYER分支
if data['type'] == 'NEW_PLAYER':
    while gf.game is None:
        await asyncio.sleep(0.1)  # 异步睡眠不阻塞事件循环
    # ... 其他逻辑

(4)添加资源清理逻辑

  • 在Game类中实现显式的cleanup方法,游戏结束时释放所有持有的资源(如WebSocket引用、任务)。
  • 应用关闭时,取消所有Gameflow的task,避免资源泄漏。

(5)优化WebSocket连接管理

  • 确保ConnectionManager在用户断开时彻底移除WebSocket引用,避免内存泄漏。
  • 添加心跳检测,定期清理管理器中的无效连接。

内容的提问来源于stack exchange,提问作者Jules Civel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 11:57:33