不能直接用 redis.Redis 做SSE消息分发,因其同步 publish() 会阻塞事件循环,导致所有 SSE 客户端卡顿;必须使用 redis.asyncio.Redis 实现异步 publish 和 get_message,并配合全局订阅器与队列分发机制保障实时性与可靠性。

为什么不能直接用 redis.Redis 做SSE消息分发
因为 redis.Redis(同步客户端)在调用 publish() 时会阻塞事件循环,而 SSE 推送必须全程异步——既要响应客户端长连接,又要实时广播消息。一旦某个 publish 调用因网络抖动或 Redis 拥塞变慢,整个 FastAPI 的异步流就会卡住,导致所有 SSE 客户端收不到后续事件。
正确做法是用 redis.asyncio.Redis,它的 publish() 是协程,可被 await,不会拖垮事件循环。同时注意:订阅端(即监听 Redis 消息并转发给 SSE 客户端的逻辑)也必须异步,不能用线程+pubsub.listen() 这种老式阻塞轮询。
如何用 redis.asyncio.Redis + EventSourceResponse 实现广播式 SSE
核心思路是:每个 SSE 连接不直接连 Redis,而是由一个全局的异步订阅器统一接收 Redis 消息,再分发给所有已建立的 SSE 客户端队列。
- 启动时创建一个后台任务(
asyncio.create_task),用redis.asyncio.Redis订阅频道,持续await pubsub.get_message(ignore_subscribe_messages=True) - 维护一个
Set[asyncio.Queue]存储所有活跃 SSE 连接对应的接收队列 - 收到 Redis 消息后,遍历该 Set,对每个
queue调用await queue.put(event_data) - SSE 路由中,每个客户端连接都调用
await sse_manager.subscribe("alarm")获取专属队列,并在响应体中用生成器逐条yield出队列内容
示例关键片段:
async def event_stream():
queue = await sse_manager.subscribe("alarm")
try:
while True:
data = await queue.get()
yield ServerSentEvent(data=json.dumps(data), event="alarm")
except asyncio.CancelledError:
await sse_manager.unsubscribe("alarm", queue)
raise
ServerSentEvent 和纯字典 yield 的区别在哪
直接 yield {"data": "..."} 会触发 FastAPI 自动封装为标准 SSE 格式(data: ...\n\n),但无法控制 event、id、retry 等字段;而显式使用 ServerSentEvent 类可精确控制每条推送的行为:
-
event="update":让前端eventSource.addEventListener("update", ...)可以按类型分流处理 -
id="12345":配合前端Last-Event-ID实现断线重连后的消息去重与续推 -
retry=5000:告诉浏览器重连间隔为 5 秒(默认是 3 秒),避免过频重试压垮服务 -
comment="keep-alive":发送注释行防止连接超时关闭(Nginx 默认 60s 关闭空闲连接)
Redis 订阅器崩溃或重连失败怎么办
异步 pubsub 不像同步版有自动重连,一旦网络中断,get_message() 会抛出 ConnectionError 并终止循环。必须手动兜底:
- 把订阅逻辑包在
while True:循环里,捕获ConnectionError后await asyncio.sleep(1)再重建pubsub实例 - 每次重建前调用旧
pubsub.close(),否则连接句柄泄漏 - 不要依赖
health_check_interval——它只检查连接池内已有连接,不保证pubsub实例本身存活
最容易被忽略的是:Redis 订阅器和 SSE 客户端生命周期完全解耦,订阅器挂了,已连接的客户端不会自动断开,只会“静默失联”。必须在订阅器恢复后主动补发最近几条关键事件(比如用 Redis Stream + XREAD 回溯),否则业务上会出现通知丢失。


















