首页
/ FastAPI SSE 流式响应实战指南:从 EventSourceResponse 到 Last-Event-ID 断点续传

FastAPI SSE 流式响应实战指南:从 EventSourceResponse 到 Last-Event-ID 断点续传

2026-09-06 13:31:48作者:裘晴惠Vivianne

Server-Sent Events(SSE)是 FastAPI 用于向浏览器等客户端持续推送数据的内置流式响应方案。本文基于仓库中的官方文档 server-sent-events.md 及其配套示例代码 docs_src/server_sent_events,完整讲解如何用 EventSourceResponse 实现 SSE 流、如何控制 event/id/retry 等 SSE 字段、如何用 raw_data 发送非 JSON 载荷,以及 FastAPI 内置的 keep-alive 心跳、防缓存与防 Nginx 缓冲机制的源码级原理。

什么是 Server-Sent Events(SSE)

SSE 是一个基于 HTTP 的服务器到客户端单向数据推送标准,响应体使用 text/event-stream 格式,浏览器可通过原生的 EventSource API 直接消费,无需 WebSocket 那样的协议升级。

SSE 与 FastAPI 的 JSON Lines 流式响应 类似,但输出格式不同:SSE 中每个事件是一个小的文本块,由若干“字段”组成,包括 dataeventidretry,事件之间以空行分隔。典型报文体如下:

data: {"name": "Portal Gun", "price": 999.99}

data: {"name": "Plumbus", "price": 32.99}

SSE 的典型应用场景包括:

  • AI 聊天的逐 token 流式输出;
  • 实时通知推送;
  • 日志(Logs)与可观测性(Observability)数据的实时上屏;
  • 一切需要服务器主动向客户端推送更新的场景。

提示:如果要流式传输二进制数据(如视频、音频流),SSE 不是合适的选择,应参考高级用户手册中的数据流式传输章节。

SSE 功能由 FastAPI 0.135.0 版本引入。

用 FastAPI 流式输出 SSE:EventSourceResponse

实现 SSE 只需两步:在路径操作函数中用 yield 逐项产出数据,并将路由的 response_class 设置为 EventSourceResponse。从 fastapi.sse 导入该类:

# docs_src/server_sent_events/tutorial001_py310.py(节选)
from collections.abc import AsyncIterable

from fastapi import FastAPI
from fastapi.sse import EventSourceResponse
from pydantic import BaseModel

app = FastAPI()


class Item(BaseModel):
    name: str
    description: str | None


items = [
    Item(name="Plumbus", description="A multi-purpose household device."),
    Item(name="Portal Gun", description="A portal opening device."),
    Item(name="Meeseeks Box", description="A box that summons a Meeseeks."),
]


@app.get("/items/stream", response_class=EventSourceResponse)
async def sse_items() -> AsyncIterable[Item]:
    for item in items:
        yield item

关键行为:

  • 每个 yield 产出的对象会被编码为 JSON,并放入 SSE 事件的 data: 字段。也就是说,上面的端点会依次发送三个 data: {...} 事件块。
  • 声明返回类型 AsyncIterable[Item],FastAPI 会用它对每个流式项做 Pydantic 校验(validate)文档化(document,生成 OpenAPI Schema)序列化(serialize)

从源码结构看,EventSourceResponse 本身在 fastapi/sse.py#L20-L33 中只是一个继承自 Starlette StreamingResponse 的“标记类”,核心职责是设定 media_type = "text/event-stream";真正的 SSE 编码逻辑位于 FastAPI 的路由层 fastapi/routing.py#L520-L551_serialize_sse_item 中:若产出项是 ServerSentEvent,则交给其自身字段决定事件结构;否则先用 Pydantic 流式字段校验并序列化为 JSON 字节,再包裹进 data: 字段。

提示:因为 Pydantic 的序列化在 Rust 侧完成,声明返回类型后能获得比不声明类型明显更高的性能

非 async 的路径操作函数

SSE 端点也可以写成普通的 def 函数(不带 async),同样使用 yield。由于函数不是协程,对应的返回类型应为同步迭代器 Iterable[Item]

# docs_src/server_sent_events/tutorial001_py310.py(节选)
from collections.abc import Iterable

@app.get("/items/stream-no-async", response_class=EventSourceResponse)
def sse_items_no_async() -> Iterable[Item]:
    for item in items:
        yield item

FastAPI 会确保这类同步生成器不会阻塞事件循环。从源码看,fastapi/routing.py#L553-L556 中会先判断端点是否为异步生成器:若是,直接取 gen.__aiter__();若不是,则通过 iterate_in_threadpool(gen) 在线程池中迭代,从而把阻塞式的 __next__ 调用隔离出事件循环。

不声明返回类型

也可以完全省略返回类型注解。此时 FastAPI 退而使用 jsonable_encoder 对每个 yield 的项做转换后再发送:

# docs_src/server_sent_events/tutorial001_py310.py(节选)
@app.get("/items/stream-no-annotation", response_class=EventSourceResponse)
async def sse_items_no_annotation():
    for item in items:
        yield item

这一分支在源码 fastapi/routing.py#L516-L518 中对应 jsonable_encoder(data) + json.dumps 的通用编码路径。tutorial001_py310.py 中还提供了第四个组合 sse_items_no_async_no_annotation(同步函数 + 无注解),四种组合在测试文件 tests/test_sse.py 中都有对应的用例验证(如 test_async_generator_no_annotationtest_sync_generator_no_annotation)。

ServerSentEvent:完整控制 SSE 事件字段

当需要设置 eventidretrycomment 等 SSE 字段时,可以 yield ServerSentEvent 对象而不是裸数据。从 fastapi.sse 导入:

# docs_src/server_sent_events/tutorial002_py310.py
from collections.abc import AsyncIterable

from fastapi import FastAPI
from fastapi.sse import EventSourceResponse, ServerSentEvent
from pydantic import BaseModel

app = FastAPI()


class Item(BaseModel):
    name: str
    price: float


items = [
    Item(name="Plumbus", price=32.99),
    Item(name="Portal Gun", price=999.99),
    Item(name="Meeseeks Box", price=49.99),
]


@app.get("/items/stream", response_class=EventSourceResponse)
async def stream_items() -> AsyncIterable[ServerSentEvent]:
    yield ServerSentEvent(comment="stream of item updates")
    for i, item in enumerate(items):
        yield ServerSentEvent(data=item, event="item_update", id=str(i + 1), retry=5000)

ServerSentEventfastapi/sse.py#L52-L156 中是一个 Pydantic 模型,各字段的含义与校验规则如下:

字段 类型 含义 校验规则
data 任意可 JSON 序列化值 事件载荷,始终被序列化为 JSON(字符串也会加引号,data="hello" 上线为 data: "hello" raw_data 互斥
raw_data str 原样字符串写入 data: 字段,不做 JSON 编码 data 互斥
event str 事件类型名,浏览器端对应 addEventListener(event, ...);缺省时走通用 message 事件 必须单行,不能含 \r/\n
id str 事件 ID,浏览器自动重连时通过 Last-Event-ID 头回传 必须单行且不能含空字符 \0
retry int 断线重连等待时间,单位毫秒 非负整数
comment str : 开头的注释行,客户端忽略,常用于 keep-alive 心跳

data 字段接受任何可 JSON 序列化的值,包括 Pydantic 模型、dict、list、字符串、数字等。在路由层(fastapi/routing.py#L524-L551),如果 data 是 Pydantic 模型(具备 model_dump_json),会直接调用其 Rust 侧序列化;否则回退到 jsonable_encoder + json.dumps。注意:当端点直接产出 ServerSentEvent 时,框架会跳过流式字段的 Pydantic 校验,因为使用者可能有意混合不同类型的事件载荷。

字段的合法性由 Pydantic 校验器强制保证(fastapi/sse.py#L36-L49#L148-L156):event/id 含换行符、id 含空字符、retry 为负数或浮点数、dataraw_data 同时设置,都会直接抛出 ValueErrortests/test_sse.py 中的 test_server_sent_event_null_id_rejectedtest_server_sent_event_single_line_fields_reject_newlinestest_server_sent_event_negative_retry_rejected 等用例逐一验证了这些边界。

发送原始数据:raw_data

如果必须发送不经过 JSON 编码的数据(预格式化文本、日志行,或类似 [DONE] 这样的“哨兵(Sentinel)”标记值),使用 raw_data 而不是 data

# docs_src/server_sent_events/tutorial003_py310.py
from collections.abc import AsyncIterable

from fastapi import FastAPI
from fastapi.sse import EventSourceResponse, ServerSentEvent

app = FastAPI()


@app.get("/logs/stream", response_class=EventSourceResponse)
async def stream_logs() -> AsyncIterable[ServerSentEvent]:
    logs = [
        "2025-01-01 INFO  Application started",
        "2025-01-01 DEBUG Connected to database",
        "2025-01-01 WARN  High memory usage detected",
    ]
    for log_line in logs:
        yield ServerSentEvent(raw_data=log_line)

此时每行日志会原样出现在 data: 字段中,不会被加上 JSON 引号。对应的测试见 tests/test_sse.py 中的 test_raw_data_sent_without_json_encoding

注意:dataraw_data 互斥,每个 ServerSentEvent 只能设置其中之一,否则 Pydantic 模型校验器会直接报错。

Last-Event-ID 恢复流

按 SSE 规范,浏览器在连接中断后自动重连时,会把最近一次成功接收的事件 id 放在请求头 Last-Event-ID 中回传。服务端可以将其作为 Header 参数读取,从客户端断点处继续推送,避免重复或遗漏事件:

# docs_src/server_sent_events/tutorial004_py310.py
from collections.abc import AsyncIterable
from typing import Annotated

from fastapi import FastAPI, Header
from fastapi.sse import EventSourceResponse, ServerSentEvent
from pydantic import BaseModel

app = FastAPI()


class Item(BaseModel):
    name: str
    price: float


items = [
    Item(name="Plumbus", price=32.99),
    Item(name="Portal Gun", price=999.99),
    Item(name="Meeseeks Box", price=49.99),
]


@app.get("/items/stream", response_class=EventSourceResponse)
async def stream_items(
    last_event_id: Annotated[int | None, Header()] = None,
) -> AsyncIterable[ServerSentEvent]:
    start = last_event_id + 1 if last_event_id is not None else 0
    for i, item in enumerate(items):
        if i < start:
            continue
        yield ServerSentEvent(data=item, id=str(i))

实现要点:

  1. 每个事件都带上递增的 id(示例中为 str(i));
  2. 端点通过 Header() 依赖声明 last_event_id 参数(默认 None 表示首次连接);
  3. Last-Event-ID 时从 last_event_id + 1 开始 yield,跳过客户端已收到的事件。

SSE 支持 POST:任何 HTTP 方法都可以

SSE 并不绑定 GETEventSourceResponse 可与任意 HTTP 方法搭配,这对 SSE-over-POST 类协议(如 MCP 模型上下文协议)尤为重要:

# docs_src/server_sent_events/tutorial005_py310.py
from collections.abc import AsyncIterable

from fastapi import FastAPI
from fastapi.sse import EventSourceResponse, ServerSentEvent
from pydantic import BaseModel

app = FastAPI()


class Prompt(BaseModel):
    text: str


@app.post("/chat/stream", response_class=EventSourceResponse)
async def stream_chat(prompt: Prompt) -> AsyncIterable[ServerSentEvent]:
    words = prompt.text.split()
    for word in words:
        yield ServerSentEvent(data=word, event="token")
    yield ServerSentEvent(raw_data="[DONE]", event="done")

这个 AI 聊天式示例展示了两种技巧的组合:逐词以 event="token" 事件推送,最后用 raw_data="[DONE]" 发出非 JSON 的结束哨兵。EventSourceResponse 的类文档中也明确说明其“works with any HTTP method (GET, POST, etc.)”(见 fastapi/sse.py#L20-L31),测试 test_post_method_sse 验证了 POST 场景。

技术细节:内置的 SSE 最佳实践

FastAPI 将若干 SSE 工程惯例内置为开箱即用的行为,无需任何额外配置:

1. Keep-alive 心跳 ping:当生成器在 15 秒内没有产出任何消息时,框架会自动发送一条注释形式的 ping(: ping),防止部分代理/负载均衡器因空闲而关闭长连接——这与 HTML 规范中 SSE 章节的建议一致。

源码层面,心跳常量定义在 fastapi/sse.py#L236-L241

KEEPALIVE_COMMENT = b": ping\n\n"
_PING_INTERVAL: float = 15.0

而调度逻辑在 fastapi/routing.py#L558-L608:FastAPI 用 anyio 内存流把“生成器迭代”与“心跳计时”解耦成两个并发任务——_producer 负责把 yield 出来的项序列化后写入流;_keepalive_inserter 则用 anyio.fail_after(_PING_INTERVAL) 包裹读取,一旦 15 秒超时没有新数据,就向输出流插入 KEEPALIVE_COMMENT。这种设计保证了 keep-alive 计时不会直接取消用户生成器,对线程池中运行的同步生成器同样有效。测试 test_keepalive_ping_asynctest_keepalive_ping_synctest_no_keepalive_when_fast 分别验证了异步/同步生成器的心跳触发,以及数据发送足够快时不会插入多余 ping 的行为。

2. 防缓存:响应固定携带 Cache-Control: no-cache 头,阻止中间层缓存流内容(fastapi/routing.py#L643)。

3. 防 Nginx 缓冲:响应固定携带 X-Accel-Buffering: no 头,让 Nginx 等代理对 SSE 响应禁用缓冲,保证事件实时到达客户端(fastapi/routing.py#L644-L645)。

4. OpenAPI 文档支持fastapi/sse.py#L9-L17 定义了与 OpenAPI 3.2 规范(“Special Considerations for Server-Sent Events”一节)对齐的事件 Schema(data/event/id 为字符串、retry 为不小于 0 的整数)。当端点声明了流式返回类型时,OpenAPI Schema 会据此生成流式事件的描述,tests/test_sse.py 中的 test_sse_router_typed_openapi_schema 验证了 Router 场景下的 Schema 输出。

小结

能力 做法
基础 SSE 流 yield + response_class=EventSourceResponse
带校验/文档/高性能序列化 返回类型标注 AsyncIterable[Item] / Iterable[Item]
控制 event/id/retry/comment yield ServerSentEvent(...)
发送非 JSON 文本(日志、[DONE] ServerSentEvent(raw_data=...)
断线重连续传 读取 Last-Event-ID Header,从断点 yield
非 GET 场景(MCP 等) 任意方法(如 @app.post)+ EventSourceResponse
心跳/防缓存/防代理缓冲 框架内置,零配置

配套的全部示例代码位于 docs_src/server_sent_events(tutorial001 至 tutorial005 五个文件),相关行为测试集中在 tests/test_sse.py,核心实现集中在 fastapi/sse.pyfastapi/routing.py 的 SSE 分支(is_sse_stream 判定见 fastapi/routing.py#L400)。

登录后查看全文
热门项目推荐
相关项目推荐