FastAPI+SSE实现单一Redis订阅者报错及解决方案咨询
FastAPI + SSE + 单Redis订阅者实现方案
错误原因分析
你用全局变量实现单Redis订阅者时出现ERROR: ASGI callable returned without completing response.,核心问题有两个:
- 全局Redis订阅连接是单例模式,当某个SSE客户端断开时,订阅循环可能抛出异常(如连接中断),导致这个全局连接直接失效,后续请求复用它必然出错。
- 路由里的订阅循环没有正确处理客户端断开的场景,FastAPI期望响应生成器能正常完成,但你的代码在客户端断开后没及时终止循环,引发ASGI响应未完成的错误。
可行的单Redis订阅者实现方案
要实现全局唯一Redis订阅者,同时给所有SSE客户端广播消息,需要用后台任务维护Redis订阅,搭配异步队列中转消息的架构,具体实现如下:
1. 初始化全局队列与后台订阅任务
from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse import asyncio import aioredis app = FastAPI() # 全局异步队列,用于中转Redis订阅到的消息 message_queue = asyncio.Queue(maxsize=100) # Redis订阅连接(由后台任务维护) redis_sub = None async def redis_subscriber(): global redis_sub # 建立Redis连接并订阅目标频道 redis_sub = await aioredis.from_url("redis://localhost") pubsub = redis_sub.pubsub() await pubsub.subscribe("your_target_channel") try: # 持续监听Redis消息,转发到全局队列 async for message in pubsub.listen(): if message["type"] == "message": await message_queue.put(message["data"].decode()) except asyncio.CancelledError: # 服务关闭时,清理订阅与连接 await pubsub.unsubscribe("your_target_channel") await pubsub.close() await redis_sub.close() except Exception as e: print(f"Redis订阅异常: {e}") await redis_sub.close() @app.on_event("startup") async def startup(): # 启动后台Redis订阅任务 asyncio.create_task(redis_subscriber())
2. SSE路由实现
async def sse_message_generator(request: Request): while True: # 检测客户端是否断开连接,及时终止循环 if await request.is_disconnected(): break # 从队列取消息,超时发送心跳保持连接 try: msg = await asyncio.wait_for(message_queue.get(), timeout=5) yield f"data: {msg}\n\n" except asyncio.TimeoutError: yield ": ping\n\n" @app.get("/sse") async def sse_broadcast(request: Request): return StreamingResponse( sse_message_generator(request), media_type="text/event-stream" )
关键设计点说明
- 后台任务独立维护订阅:Redis订阅逻辑在后台单独运行,不受单个SSE客户端断开的影响,保证全局只有一个Redis连接。
- 异步队列中转消息:所有SSE客户端共享同一个队列,天然实现消息广播,避免每个请求重复订阅Redis。
- 客户端断开检测:通过
request.is_disconnected()及时终止生成器,解决ASGI响应未完成的错误。 - 心跳机制:超时发送空消息(ping),防止浏览器因长时间无数据主动断开连接。
内容的提问来源于stack exchange,提问作者Pratik Bhadane
相关产品推荐
相关产品推荐

