实时日志、AI token、任务进度这类“服务器持续推送、浏览器主要接收”的场景,优先考虑 Server-Sent Events(SSE)。FastAPI 现在提供原生 EventSourceResponseServerSentEvent;真正的生产难点不在 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 文档展示了 EventSourceResponseServerSentEventLast-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。服务端要明确以下契约:

  1. id 在一个流分区内单调递增,或至少能唯一定位顺序。
  2. 查询语义是“严格大于游标”,避免重发最后一条。
  3. 客户端消费必须幂等,因为网络边界仍可能造成重复。
  4. 事件保留期大于常见离线窗口。
  5. 游标过期时返回显式的重建指令,而不是悄悄从最新位置开始。
event: stream.reset
data: {"reason":"cursor_expired","snapshotUrl":"/runs/42"}

事件 ID 不应包含 token、用户邮箱或数据库主键组合等敏感信息。若使用全局序列,要在查询层强制租户与资源范围;拥有一个合法游标不等于有权读取相邻事件。

Nginx 为什么会让流“攒一批才出现”?

反向代理默认可能缓冲上游响应。Nginx proxy_buffering 文档说明,缓冲可由响应头 X-Accel-Buffering 控制。FastAPI 原生 SSE 会设置 X-Accel-Buffering: noCache-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 校验和短会话。若使用短期订阅票据:

  1. 票据只授权一个资源和一种动作。
  2. 有极短 TTL,且最好一次性消费。
  3. 不放在会被代理、分析或 Referer 记录的长期 URL 中。
  4. 服务端仍根据当前用户重做资源授权。

跨域时不要使用宽泛的 Access-Control-Allow-Origin: * 搭配凭据。先问是否真的需要跨域;同源反代通常更简单。

生产验收应该检查什么?

curl -Nsv https://example.com/api/streams/42 \
  -H 'Accept: text/event-stream'

检查清单:

  • Content-Typetext/event-stream,且没有被 CDN 转成缓存对象。
  • 第一条事件在预期时间内抵达,不是连接结束后一次性出现。
  • 空闲期心跳早于最短代理超时。
  • 杀掉网络后,重连携带 Last-Event-ID 并从下一条恢复。
  • 重复事件不会造成重复扣款、重复写入或重复通知。
  • 客户端关闭后,下游生成或游标按设计释放。
  • 并发连接数、连接时长、断开原因和恢复落后量可观测。
  • 日志不记录凭据、URL 票据或完整敏感事件。

常见失败模式

HTTP 200,但页面几秒才更新一次

通常是 Nginx、CDN 或应用中间件缓冲。逐层绕过代理对比时间,检查压缩与响应转换,再核对 X-Accel-Bufferingproxy_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 指南