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

如何在FastAPI中稳定实现SSE并解决多用户队列冲突问题

FastAPI中正确实现SSE并从函数触发事件的方案

问题根源

当前实现的核心问题在于全局共享asyncio.Queue:

  • 所有用户的SSE连接共用同一个队列,事件会被随机客户端接收,无法实现用户专属的状态推送
  • 连接断开后未清理队列资源,长期运行会引发内存泄漏;多用户场景下队列的消费/生产逻辑完全混乱,直接导致崩溃

正确实现方案

核心思路是为每个用户连接维护独立的事件通道,通过用户ID关联专属队列,确保事件精准推送给目标用户。

1. 定义全局订阅存储与锁

import asyncio
from fastapi import FastAPI, Request, StreamingResponse, Depends
from fastapi.security import OAuth2PasswordBearer
import json
from typing import Dict

app = FastAPI()
oauth2_scheme = OAuth2PasswordBearer(tokenUrl="token")

# 存储用户ID到专属事件队列的映射
subscriptions: Dict[str, asyncio.Queue] = {}
# 协程安全锁,保护订阅字典的读写操作
subscriptions_lock = asyncio.Lock()

2. 改造SSE流接口(/stream)

为每个用户创建独立队列,连接断开时自动清理资源:

def parse_user_id_from_token(token: str) -> str:
    # 替换为你的实际Token解析逻辑,从Token中提取用户唯一ID
    return "user_123"

@app.get("/stream")
async def stream(token: str = Depends(oauth2_scheme)):
    user_id = parse_user_id_from_token(token)
    
    async with subscriptions_lock:
        # 为当前用户初始化专属队列,已存在则复用
        if user_id not in subscriptions:
            subscriptions[user_id] = asyncio.Queue()
        user_queue = subscriptions[user_id]
    
    async def event_generator():
        try:
            while True:
                # 等待当前用户的事件
                event = await user_queue.get()
                yield f"data: {json.dumps(event)}\n\n"
        except asyncio.CancelledError:
            # 客户端断开连接时,清理该用户的队列
            async with subscriptions_lock:
                if user_id in subscriptions:
                    del subscriptions[user_id]
            raise

    return StreamingResponse(event_generator(), media_type="text/event-stream")

3. 改造Webhook接口(/webhook)

根据用户ID推送事件到对应专属队列:

@app.post("/webhook")
async def generate_status(request: Request, db: db_dependency):
    payload = await request.json()
    user_id = payload.get("user_id")
    
    if not user_id:
        return {"message": "Webhook payload missing user_id"}
    
    async with subscriptions_lock:
        user_queue = subscriptions.get(user_id)
        if not user_queue:
            return {"message": "User has no active SSE subscription"}
    
    # 推送状态到目标用户的专属队列
    await user_queue.put({"status": "uploaded"})
    return {"message": "Webhook processed and status pushed to user"}

关键注意事项

  • 用户ID的获取:生产环境建议通过OAuth2 Token解析用户ID,避免明文传递;也可通过请求头、会话ID等方式,确保每个SSE连接对应唯一用户
  • 资源清理:必须在客户端断开连接(触发asyncio.CancelledError)时移除用户队列,防止内存泄漏
  • 协程安全:使用asyncio.Lock保护订阅字典的读写,避免多协程并发操作引发的数据混乱
  • 异常处理:Webhook中需处理用户未订阅的情况,避免无意义的队列操作报错

内容的提问来源于stack exchange,提问作者רועי כחלון

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 13:04:55