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

使用FastAPI+Starlette SSE时,客户端断开引发asyncio.InvalidStateError

FastAPI SSE发布订阅任务终止问题修复

问题原因

当客户端断开连接时,EventSourceResponse会取消对应的异步生成器任务,导致pubsub.subscribe()中的await waiter触发CancelledError,对应的waiter Future进入已取消状态。后续ticker任务调用publish时,尝试给已取消的Future设置结果,就会抛出InvalidStateError,直接终止ticker任务。

修复代码

改用订阅者列表管理的方式,让每个订阅者持有独立的Future,避免单个订阅者断开影响全局发布流程:

import asyncio
from fastapi import FastAPI, Request
from sse_starlette.sse import EventSourceResponse


class PubSub:
    def __init__(self):
        self.subscribers = set()

    def publish(self, value):
        # 遍历当前订阅者副本,避免遍历中集合变更
        for subscriber in list(self.subscribers):
            if not subscriber.done():
                subscriber.set_result((value, asyncio.Future()))

    async def subscribe(self):
        current_future = asyncio.Future()
        self.subscribers.add(current_future)
        try:
            while True:
                value, current_future = await current_future
                yield value
        finally:
            # 无论订阅正常结束还是被取消,都移除订阅者
            self.subscribers.discard(current_future)

    __aiter__ = subscribe

pubsub = PubSub()

async def ticker(pubsub):
    counter = 0
    while True:
        pubsub.publish(counter)
        counter += 1
        await asyncio.sleep(1)

app = FastAPI()

@app.on_event("startup")
async def on_startup():    
    asyncio.create_task(ticker(pubsub), name='my_task')

@app.get('/stream')
async def message_stream(request: Request):
    async def event_publisher():
        try:
            async for event in pubsub:
                yield dict(data=event)
        except asyncio.CancelledError as e:
            print(f"Disconnected from client (via refresh/close) {request.client}")
            raise e
    return EventSourceResponse(event_publisher())

方案说明

  1. 订阅者独立管理:用subscribers集合维护所有活跃订阅者的Future,每个客户端连接对应一个独立的Future,单个客户端断开只会移除自己的Future,不会影响其他订阅者和发布任务。
  2. 安全的结果设置:publish时先检查Future是否未完成,避免给已取消或已完成的Future设置结果。
  3. 自动清理订阅者:在subscribe的finally块中确保订阅者无论正常退出还是被取消,都会从集合中移除,防止内存泄漏。

修改后,即使有客户端断开连接,ticker任务仍能持续运行,不会再抛出InvalidStateError。

内容的提问来源于stack exchange,提问作者DurandA

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 09:28:16