如何在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,提问作者רועי כחלון
相关产品推荐
相关产品推荐

