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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 18:35:07