实时日志、AI token、任务进度这类“服务器持续推送、浏览器主要接收”的场景,优先考虑 Server-Sent Events(SSE)。FastAPI 现在提供原生 EventSourceResponse 与 ServerSentEvent;真正的生产难点不在 yield,而在事件身份、恢复语义、代理缓冲、慢客户端、取消传播和可观测性。
什么时候应选 SSE,而不是 WebSocket?
SSE 是基于 HTTP 的单向事件流。浏览器会解析 text/event-stream,连接中断后还可以携带 Last-Event-ID 重新连接。它适合服务端持续推送、客户端偶尔用普通 HTTP 发命令的系统。
| 需求 | SSE | WebSocket | 普通轮询 |
|---|---|---|---|
| 服务端单向推送 | 很合适 | 可以,但更重 | 延迟与请求开销较高 |
| 双向高频交互 | 不合适 | 最合适 | 不合适 |
| 复用 HTTP 鉴权和代理 | 简单 | 需要升级连接配置 | 简单 |
| 原生断线恢复线索 | Last-Event-ID |
需要自定义 | 由请求游标决定 |
| 浏览器 API | EventSource |
WebSocket |
fetch |
如果客户端必须持续上传音频、光标或游戏状态,使用 WebSocket。若只是把生成结果、构建进度或状态变更送到页面,SSE 通常更容易运维。ZoyTown 的原子双语发布系统也采用“状态先明确、传输后验证”的思路:协议选择应服从状态语义,而不是反过来。
最小可恢复 FastAPI SSE 端点怎么写?
FastAPI 官方 SSE 文档展示了 EventSourceResponse、ServerSentEvent 与 Last-Event-ID。下面是工程化骨架;事件存储只是接口示例,不代表已运行的生产实现。
from collections.abc import AsyncIterable
from typing import Annotated
from fastapi import FastAPI, Header, Request
from fastapi.sse import EventSourceResponse, ServerSentEvent
app = FastAPI()
async def read_events_after(cursor: int) -> AsyncIterable[dict]:
# 示例:实际系统可从 append-only log、Redis Streams 或数据库 outbox 读取。
for item in await event_store.list_after(cursor, limit=100):
yield item
@app.get("/runs/{run_id}/events", response_class=EventSourceResponse)
async def run_events(
run_id: str,
request: Request,
last_event_id: Annotated[str | None, Header()] = None,
) -> AsyncIterable[ServerSentEvent]:
cursor = int(last_event_id) if last_event_id else 0
async for item in read_events_after(cursor):
if await request.is_disconnected():
break
yield ServerSentEvent(
data={"runId": run_id, "status": item["status"]},
event="run.status",
id=str(item["sequence"]),
)
事件 id 必须是可比较、可恢复的稳定游标,而不是仅用于展示的随机值。最安全的做法是把业务事件先写入持久化日志或 outbox,再由流端点读取;不要先向客户端发送、再异步落库,否则断线窗口会产生无法重放的“幽灵成功”。
如何设计事件 ID 与断线续传?
浏览器重连时会把最后接收的 id 放进 Last-Event-ID。服务端要明确以下契约:
id在一个流分区内单调递增,或至少能唯一定位顺序。- 查询语义是“严格大于游标”,避免重发最后一条。
- 客户端消费必须幂等,因为网络边界仍可能造成重复。
- 事件保留期大于常见离线窗口。
- 游标过期时返回显式的重建指令,而不是悄悄从最新位置开始。
event: stream.reset
data: {"reason":"cursor_expired","snapshotUrl":"/runs/42"}
事件 ID 不应包含 token、用户邮箱或数据库主键组合等敏感信息。若使用全局序列,要在查询层强制租户与资源范围;拥有一个合法游标不等于有权读取相邻事件。
Nginx 为什么会让流“攒一批才出现”?
反向代理默认可能缓冲上游响应。Nginx proxy_buffering 文档说明,缓冲可由响应头 X-Accel-Buffering 控制。FastAPI 原生 SSE 会设置 X-Accel-Buffering: no、Cache-Control: no-cache,并在空闲时发送心跳;仍应在部署层做显式检查。
location /api/streams/ {
proxy_pass http://api;
proxy_http_version 1.1;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_buffering off;
proxy_cache off;
proxy_read_timeout 75s;
}
不要机械复制超长 proxy_read_timeout。它应该略大于心跳间隔,并与负载均衡器、CDN 和应用层空闲策略一致。生产验收要用计时输出验证首条事件延迟,而不是只看最终 HTTP 200。
心跳、超时和取消应如何配合?
HTML Living Standard 的 SSE 章节定义了事件流格式与重连行为。心跳通常使用注释行,不应伪装成业务事件。FastAPI 的原生实现会在空闲时发送 ping;你的任务仍要响应取消:
- 每次等待外部队列时设置可取消的超时。
- 客户端断开后停止数据库游标、模型生成和下游请求。
- 不要用吞掉
CancelledError的宽泛异常处理。 - 对长任务,把“任务执行”和“客户端订阅”解耦;浏览器断开不应自动取消已确认提交的任务。
这种分离也适用于动态内容页:动态 SSR 双语 SEO 实战强调页面响应与内容状态要各自可验证。SSE 的连接状态同样不能被当作业务状态。
如何处理慢客户端和背压?
无限队列会把一个慢标签页变成内存泄漏。为每个订阅者设置有界缓冲区,并按业务选择:
- 状态类事件:合并中间值,只保留最新状态。
- 审计类事件:不得丢弃,改为让客户端从持久化日志按游标追赶。
- token 流:可以小批量合并,但必须保持顺序。
- 高频指标:按时间窗口采样。
当缓冲区满时,记录一个低基数原因码,终止连接并让客户端重建,而不是继续堆积。不要把完整 prompt、token 或事件正文写进错误日志。
鉴权和跨域有哪些陷阱?
原生 EventSource 不允许像 fetch 那样任意设置请求头。Cookie 同源鉴权最简单,但必须处理 CSRF、SameSite、Origin 校验和短会话。若使用短期订阅票据:
- 票据只授权一个资源和一种动作。
- 有极短 TTL,且最好一次性消费。
- 不放在会被代理、分析或 Referer 记录的长期 URL 中。
- 服务端仍根据当前用户重做资源授权。
跨域时不要使用宽泛的 Access-Control-Allow-Origin: * 搭配凭据。先问是否真的需要跨域;同源反代通常更简单。
生产验收应该检查什么?
curl -Nsv https://example.com/api/streams/42 \
-H 'Accept: text/event-stream'
检查清单:
Content-Type是text/event-stream,且没有被 CDN 转成缓存对象。- 第一条事件在预期时间内抵达,不是连接结束后一次性出现。
- 空闲期心跳早于最短代理超时。
- 杀掉网络后,重连携带
Last-Event-ID并从下一条恢复。 - 重复事件不会造成重复扣款、重复写入或重复通知。
- 客户端关闭后,下游生成或游标按设计释放。
- 并发连接数、连接时长、断开原因和恢复落后量可观测。
- 日志不记录凭据、URL 票据或完整敏感事件。
常见失败模式
HTTP 200,但页面几秒才更新一次
通常是 Nginx、CDN 或应用中间件缓冲。逐层绕过代理对比时间,检查压缩与响应转换,再核对 X-Accel-Buffering 和 proxy_buffering。
重连后漏事件
多半是事件先发送后落库、游标用 >=/> 错误,或保留期不足。用故障注入在“写入、发送、确认”之间断网,证明恢复契约。
内存持续增长
检查每连接队列、生成任务和数据库游标是否有界且可取消。连接数稳定并不代表资源稳定。
浏览器客户端还需要做什么?
EventSource 会自动重连,但产品仍要定义用户可见状态。下面示例只展示状态机,不包含业务鉴权:
const stream = new EventSource("/api/streams/42");
stream.addEventListener("run.status", (event) => {
const payload = JSON.parse(event.data);
renderStatus(payload.status);
});
stream.addEventListener("stream.reset", async (event) => {
stream.close();
const { snapshotUrl } = JSON.parse(event.data);
await reloadSnapshot(snapshotUrl);
});
stream.onerror = () => {
renderConnectionState("reconnecting");
};
不要在每次 onerror 里立即创建新 EventSource;浏览器本身已有重连行为,叠加自定义循环会制造连接风暴。页面进入后台时也不一定要关闭连接,是否暂停应根据事件价值、移动端耗电和服务器连接预算决定。
多 worker 和水平扩容怎么处理?
进程内 asyncio.Queue 只能服务连接所在的 worker。若事件生产者和订阅者可能落在不同进程或主机,需要共享事件源,例如 Redis Streams、消息 broker、数据库 change stream 或 outbox。选择时关注:
- 是否保存可恢复游标,而不只是 pub/sub 瞬时广播。
- 消费者追赶时能否限速和分页。
- 分区键能否保持单个资源的事件顺序。
- broker 故障时是暂停流、返回 503,还是回退数据库。
- 清理策略是否会让合法
Last-Event-ID提前过期。
不要为了 SSE 强行引入 broker。单进程内部工具可以先用有界队列;当可靠恢复、跨 worker 分发或审计成为真实需求时,再升级为持久事件层。
FAQ
SSE 能用 POST 吗?
FastAPI 文档说明 SSE 可以配合 POST 响应,但浏览器原生 EventSource 只直接发起 GET。需要 POST 时可使用 fetch 读取流,或把“创建任务”和“订阅任务”拆成 POST + GET。
可以对 SSE 开 gzip 吗?
压缩中间件可能缓冲小块并增加首事件延迟。不是绝对禁止,但必须用真实链路验证 flush 行为;无法证明时,对事件流路径禁用响应转换更稳妥。
每条 token 都应该是一个事件吗?
不一定。过小事件增加 framing、调度和渲染开销。可以按字符数或很短时间窗口合并,同时保留顺序和取消语义。
最终决策
生产 SSE 的最小可靠闭环是:持久化事件 → 稳定 ID → 有界读取 → 禁止代理缓冲 → 心跳早于超时 → 断线按游标恢复 → 消费幂等 → 资源与敏感数据可审计。先证明这个闭环,再增加多路复用或复杂推送层。更多事实核验与技术表述边界可参考官方来源优先的 GEO/AEO 指南。