如何设计 SSE 流式输出
August 15, 2026
我最近把 Agent 的流式输出,改成了正确的架构,分享下我的经验。
我给自己的个人网站做了一顶 Hat。
它住在我的 Maidang OS 里,会看我正在打开什么、记住一些偏好,也会调用工具。你可以打字问它,也可以按住 V 键直接说话。
最早做流式输出时,我用了一个最自然的写法:
POST /api/os-agent/runs
-> 调用 LLM
-> 一边生成,一边直接返回 SSE
但某次意外刷新后,发现所有状态全没了,我意识到流式输出的设计没那么简单。
把调用 agent 运行和返回 sse 放在一个接口里,是错误的设计。创建任务、执行 Agent run 和 返回 SSE 三个概念混在了一起。
于是我改成这样,把执行 agent run 和订阅事件拆开
POST /api/os-agent/runs
-> PostgreSQL 创建 user message 和 queued run
-> 立即返回 202 { runId, sessionId }
POST /api/os-agent/runs/:runId/events
-> 领取(claim)并执行这个 run
-> 持续返回这个 run 的 SSE 事件
为了验证这个结果,我同时开两个标签页来模拟一个发送消息,另一个页面刷新并且竞争 events 接口
时序图是这样的:

现在系统的主语是 runId。浏览器只是来观看一个已经存在的 Run。刷新页面之后,它可以先从数据库找到当前仍在运行的 runId,如果在运行从,就重新订阅。
这也是我理解 AI 对话的表结构:一个 Agent 产品不能只有 Conversation 和 Message。
至少要把三层东西分开:
| 对象 | 它回答的问题 | Hat 里保存什么 | | ------- | ---------- | ------------------------- | | Session | 这是哪段对话 | owner、摘要、时间 | | Message | 用户最后看见了什么 | user / assistant 的最终内容 | | Run | 系统为了回答做了什么 | 状态、触发方式、优先级、模型、错误、开始和结束时间 |
使用 Redis 解决刷新后 SSE 丢失的问题
有了 Run,只解决了「我正在等哪次任务」。还没有解决「断线期间产生的内容去哪了」。
因为我在某次面试中被人问到,生产环境里,请求会打到不同机器上,如果你执行到一半刷新,你怎么接上这个流?
Hat 在 Agent loop 和 SSE 之间加了一层 Redis Stream:
Agent loop
-> XADD lifecycle / tool / message delta
-> Redis Stream
-> XREAD from cursor
-> SSE
-> Browser
每条 Redis Stream 事件都有自己的 ID。SSE 把这个 ID 放进 id: 字段,前端重连时带回 Last-Event-ID,服务端就从这个游标后面继续读。
于是刷新不再等于重来。
已经完成的对话,从 PostgreSQL 读取最终 Message;仍在生成的回答,从 Redis 重放短期事件。浏览器还会按 Stream ID 去重,所以临时断线、自动重试或多个读取窗口不会把同一个 delta 重复拼上去。
这里有一个很容易说错的地方:Redis 不是事实源。
Hat 的 Redis Stream 只保留一小时,长度大约限制在 10000 条。它负责的是实时事件、游标和短期重放。Session、Message、Run 的状态和最终答案仍然在 PostgreSQL。
换句话说:
> PostgreSQL 回答「最后发生了什么」,Redis 回答「刚才是怎么发生的」。
这也解释了为什么我没有直接用 Redis Pub/Sub。Pub/Sub 适合在线广播,但订阅者断线期间的消息不会回来。LLM 输出需要的是一段短期、可按游标重放的日志,Redis Streams 更合适。
真正拖慢流式输出的,差点不是 LLM
把 Redis 接上之后,我还遇到过一个很反直觉的问题:流能通,但每段输出总要慢五秒左右。
最后用真实 Redis 测试定位到,阻塞读取和写入共用了同一个连接:
```text
XREAD BLOCK 5000
XADD ...
```
同一条连接正在等 XREAD,后面的 XADD 只能排队。实测写入延迟到了约 5708ms。看起来像模型慢,实际是 Redis 客户端内部发生了队头阻塞。
现在的连接拓扑是:
- Serverless 实例共享一个非阻塞 writer,只做
XADD、EXPIRE等写操作 - 每条活跃 SSE 单独
duplicate()一个 blocking reader - Run 结束、浏览器断开或读取异常时,立即释放 reader
这不是微小优化。如果 writer 和 reader 不隔离,LLM 就算每 50ms 吐一次 token,用户仍然只能每五秒看见一坨。
Redis 官方的 LLM Streaming 教程也用了同样的原则:生产者用 XADD 写入,消费者用隔离连接执行阻塞 XREAD,再把 chunk 转发给浏览器。传输层可以是 WebSocket,也可以像 Hat 一样用 SSE;关键是中间有一段可重放的事件日志。
刷新之后,具体发生了什么
现在页面重新打开时,Hat 会先请求当前 Session:
```text
GET /api/os-agent/sessions/current
-> canonical messages
-> activeRun
```
如果没有 active Run,页面只恢复历史消息。
如果有,前端用同一个 runId 重新连接 /runs/:runId/events。服务端从 Redis Stream 读取这个 Run 的事件,已经完成但没有进入当前页面内存的 delta 会被补回来;最终状态则仍然落进 PostgreSQL。
断开连接和停止任务也被分成了两个动作。
刷新、切换页面、Wi-Fi 抖动只是 disconnect,不代表用户想取消生成。只有用户明确按下停止,前端才调用单独的 cancel 接口。Vercel AI SDK 的 resumable stream 文档也专门强调了这个区别:在可恢复流里,client abort 只是断开观看,真正停止底层工作需要独立的 stop endpoint。
这套方案解决了什么,又没有解决什么
它已经解决了几个最影响体验的问题:
- 创建任务不再依赖一条长连接成功返回
- 页面刷新后能找到 active Run 并重新订阅
- 临时断线可以从 Redis 游标后重放
- 完成后的消息以 PostgreSQL 为准,不依赖浏览器内存
- 同一个 Session 只允许一个 Run 真正执行
- 用户消息优先于 Hat 的主动唤醒,可以抢占 proactive Run
- 工具事件、模型调用和最终消息有各自的数据归属
但我不会把它叫作「完全持久化的 Agent」。
当前 /runs/:runId/events 仍然会在这次 Serverless invocation 里领取并执行 Agent loop,route 配置的 maxDuration 是 180 秒。浏览器断开后,代码不会把 disconnect 当成 cancel,但如果 Vercel 最终终止了这个函数,Redis 只能重放已经写进去的事件,不能让死掉的 Agent 从中间继续思考。
> 可恢复的流,不等于可恢复的执行。
这条边界很重要。Redis Streams 是 event transport,不是任务队列,也不是执行引擎。
主流方案其实分三层
查了一圈现在的 AI Chat 和 Agent 基础设施后,我发现大家并不是在争论 SSE 还是 WebSocket,而是在按任务寿命选择不同层级。
第一层是 request-bound streaming。
模型调用仍在一次请求里完成,但服务端会继续消费模型流,并把最终 Message 保存下来。即使客户端断开,用户刷新后至少能从数据库看到完整结果。Vercel AI SDK 的 Message Persistence 指南就是这个思路,同时建议为持久化消息使用稳定 ID。
第二层是 resumable streaming。
数据库保存 activeStreamId,Redis 保存 UI stream,创建和恢复使用两个接口。页面加载后用 chatId 找到 active stream,再继续读取。Vercel AI SDK 当前的官方方案明确要求:消息持久化、active stream 映射、Redis,以及 create/resume endpoints。
Hat 现在大致处在这一层,只是把 chat stream 进一步建模成了 Run event stream。
第三层是 durable execution。
Agent 不再活在 SSE route 里,而是活在 Workflow、独立 worker 或持久化 Agent runtime 里。SSE 只负责订阅,不负责维持任务生命。
Vercel Workflow 的 resumable stream 会返回 workflow runId,恢复端点可以按 startIndex 从上次收到的 chunk 继续;函数超时或页面刷新都不影响 workflow 本身。LangGraph 的思路也类似:用 checkpointer 保存 thread 的图状态,用 store 保存跨 thread 的长期数据,分别解决中断恢复和长期记忆。
所以真正的主流架构不是某一个库,而是同一组边界:
```text
PostgreSQL Redis / stream Workflow / worker
业务事实 实时重放 持久执行
Message / Run delta / cursor retry / checkpoint
```
如果 Hat 的用户量继续上来,我下一步会怎么做
下一步我不会先上 Kafka,也不会立刻拆很多微服务。
对这个网站最合适的升级,是保留现在的 API 和表结构主线,只把执行从 SSE route 里拿出去:
```text
- POST /runs
-> 事务内创建 Message + Run
-> 启动 Vercel Workflow
-> 返回 runId 2. Workflow / worker
-> 执行 Agent loop
-> 更新 Run 状态和 checkpoint
-> 写 Redis Stream
-> 保存最终 Message 3. GET /runs/:runId/events
-> 只做鉴权、重放和实时订阅
-> 不再 claim 或执行 Agent
```
这样浏览器、SSE 函数和 Agent executor 才真正彼此独立。页面可以关,SSE 可以重建,Serverless route 可以超时,Workflow 仍然知道自己跑到了哪一步。
数据模型上,我还会补四件事:
- 给 Message 增加结构化
parts,而不是永远只有一个content字符串。文本、tool call、tool result、approval 和 artifact 应该能独立演进。 - 增加少量持久化
RunEvent,只保存 started、tool completed、approval requested、completed、failed 这类业务事件。高频 token delta 继续留在 Redis,不把 PostgreSQL 写爆。 - 在 Run 上保存
workflowRunId、activeStreamId、幂等键和恢复 checkpoint。客户端重试创建请求时,必须拿回同一个 Run,而不是再生成一份。 - 给事件增加单调递增的
seq。Redis ID 适合传输游标,seq更适合业务去重、审计和跨存储对账。
如果再往上到大量并发用户,我才会继续拆:
- Agent 执行通过 Workflow 或队列控制并发,按 Session 保证顺序
- Redis Stream 从 per-session 调整为 per-run 或分片 stream,减少无关事件过滤
- SSE 从「每个 Serverless 请求一条 Redis reader」迁到常驻 realtime gateway,由网关复用 Redis 连接并向大量客户端 fan-out
- delta 继续按时间和字符数合批,避免每个 token 都产生一次 Redis、网络和 React 更新
- 完成 Run 做 snapshot,旧 stream 按 TTL 清理,恢复时优先加载 snapshot,再补增量
- 用 queue depth、首 token 延迟、重连率、Run 卡死率和每 Run 成本决定是否需要 Kafka、NATS 或更重的事件系统
对于我现在这个个人网站,PostgreSQL + Redis Streams + Vercel Workflow 已经足够走很远。过早上 Kafka,不会让 Hat 更可靠,只会让我多维护一套暂时没有流量证明的基础设施。
最后
我一开始以为流式输出的问题是:怎么让字一个一个出来。
做完这一轮才发现,那个只是最表面的一层。
真正的问题是:当浏览器刷新、网络断开、模型变慢、函数超时的时候,这次任务在系统里还算不算存在?谁保存它的身份,谁保存最终事实,谁保存刚刚产生的事件,又由谁负责把它继续跑完?
SSE 只是最后一公里。
一个 AI 对话从 demo 走向产品,是从「这条连接正在输出」变成「这次 Run 正在发生」。
Hat 现在已经能在刷新后重新接上它的回答。下一步,是让它在没有任何页面看着的时候,也能把该做的事做完。
参考资料
- [Vercel AI SDK:Chatbot Message Persistence](https://ai-sdk.dev/docs/ai-sdk-ui/chatbot-message-persistence)
- [Vercel AI SDK:Chatbot Resume Streams](https://ai-sdk.dev/docs/ai-sdk-ui/chatbot-resume-streams)
- [Vercel Workflow:Resumable Streams](https://useworkflow.dev/docs/ai/resumable-streams)
- [Vercel Workflow:Chat Session Modeling](https://useworkflow.dev/docs/ai/chat-session-modeling)
- [Vercel Functions:Configuring Maximum Duration](https://vercel.com/docs/functions/configuring-functions/duration)
- [Redis:Streaming LLM Output Using Redis Streams](https://redis.io/tutorials/howtos/solutions/streams/streaming-llm-output/)
- [LangGraph:Persistence](https://docs.langchain.com/oss/javascript/langgraph/persistence)
