使用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())
方案说明
- 订阅者独立管理:用
subscribers集合维护所有活跃订阅者的Future,每个客户端连接对应一个独立的Future,单个客户端断开只会移除自己的Future,不会影响其他订阅者和发布任务。 - 安全的结果设置:
publish时先检查Future是否未完成,避免给已取消或已完成的Future设置结果。 - 自动清理订阅者:在
subscribe的finally块中确保订阅者无论正常退出还是被取消,都会从集合中移除,防止内存泄漏。
修改后,即使有客户端断开连接,ticker任务仍能持续运行,不会再抛出InvalidStateError。
内容的提问来源于stack exchange,提问作者DurandA
相关产品推荐
相关产品推荐

