主题
Streaming
要点
- Agent 项目里的流式输出,不只传产物片段,也要传研究、规划、产物、审稿、交付检查等阶段事件。
- SSE 适合服务端到客户端的单向事件流,能覆盖大多数 Agent 任务的进度反馈和增量文本输出。
- 前端消费事件契约,服务端和任务图负责状态推进、错误分类、恢复点和后台写回。
1. 为什么 Agent 任务需要流式输出
Agent 项目通常不是一次模型调用。用户提交目标后,系统可能会经过目标解析、资料检索、来源评估、结构生成、产物生成、事实检查、风格审稿和交付检查。每个节点都可能调用不同 Agent、工具或人工确认点。
如果等所有节点完成后一次性返回,用户只能看到一个很长的等待过程,也难以判断任务停在资料召回、模型生成、审稿规则,还是交付前校验。从体验角度看,关键指标不是总响应时间,而是首 token 到达时间(Time-to-First-Token,TTFT):用户从提交到看到第一段反馈的间隔。OpenAI、Anthropic 和 Google 的聊天接口都把流式输出作为默认方式,正是因为同样的总生成时间,流式呈现会让用户感觉更快,也更容易建立对系统的信任。流式输出要解决的是「过程可见」和「可恢复」:
- 研究 Agent 正在检索哪些来源,哪些来源被接受或剔除。
- 规划 Agent 是否已经给出任务结构,是否等待用户确认。
- 执行 Agent 是否开始输出产物片段,当前片段属于哪个部分。
- 审稿 Agent 发现了哪些阻塞问题,问题对应哪段产物或来源。
- 交付检查是否通过,失败时需要回到哪个节点继续处理。
- 连接中断后,前端能从哪个任务版本继续恢复。
对 Agent 产品来说,流式输出承担的是任务执行日志的一部分。它既服务界面反馈,也服务恢复、追踪和后续审计。
2. 为什么优先选择 SSE
常见方案有三种:
| 方案 | 适用场景 | 代价 |
|---|---|---|
| 普通 HTTP | 短任务、一次性结果 | 用户看不到中间状态 |
| SSE | 服务端持续推送事件 | 单向通信,适合 Agent 任务 |
| WebSocket | 双向实时协作、多人编辑 | 连接管理和状态协调更复杂 |
Agent 任务通常是用户发起任务,服务端持续返回阶段事件、来源信息、产物片段和审稿结果。这个方向更接近 SSE 的单向事件流模型。WebSocket 可以用于多人协作、实时编辑或高频双向交互;只处理任务进度和模型增量输出时,SSE 的连接模型更简单,也更容易接入边缘运行时和普通 HTTP 基础设施。
SSE 还有一个对 Agent 任务很有用的特性:事件可以带 id。根据 MDN 的 EventSource 文档,浏览器收到带 id 的事件后会在断线重连时自动带上 Last-Event-ID 请求头,服务端可以据此判断从哪个事件继续补发。这个机制让前端在断线后恢复已有状态,而不必重新发起整个任务。
需要注意三个工程细节:
- 浏览器原生的
EventSource只支持 GET 请求,且不能设置自定义请求头(如Authorization)。因此带身份认证或需要 POST 请求体的 Agent 接口,通常用fetch+ReadableStream手动解析 SSE 帧,自己在客户端实现重连和Last-Event-ID管理。 - 代理和负载均衡经常会缓冲响应。Nginx、Cloudflare 或 AWS ALB 默认可能把整段流收集后再转发,导致用户看到「最后一次性出现」。服务端需要设置
Content-Type: text/event-stream、Cache-Control: no-cache、Connection: keep-alive和X-Accel-Buffering: no,并在代理层关闭 buffering,必要时用 SSE 注释行(: heartbeat\n\n)每 15–30 秒发送一次心跳,避免空闲超时。 - 慢客户端或高 token 速率下会产生背压(backpressure)。如果服务端一股脑把事件塞进内存,而客户端来不及消费,内存会无界增长。生产环境应通过有界队列或检测
drain事件,让生产速度匹配消费速度。
3. 事件类型设计
只推送 token 会把业务状态藏在产物流里,前端也会被迫从文本变化里猜测任务进度。Agent 任务更适合把事件分层:
- 任务级事件:
stage、checkpoint、done、error。 - 资料级事件:
source、source_rejected、citation_needed。 - 内容级事件:
plan、section_started、token、section_completed。 - 审稿级事件:
review_finding、revision_started、revision_completed。
txt
id: task_123:12
event: stage
data: {"stage":"research","label":"资料检索中"}
id: task_123:13
event: source
data: {"title":"LangGraph Overview","sourceType":"official","reliability":"high"}
id: task_123:18
event: plan
data: {"sections":["为什么需要图编排","State 的作用","人工确认点"]}
id: task_123:26
event: token
data: {"sectionId":"s2","content":"LangGraph 在 Agent 项目里的价值..."}
id: task_123:41
event: review_finding
data: {"severity":"P1","sectionId":"s3","message":"第 3 节缺少来源引用"}
id: task_123:52
event: done
data: {"taskId":"task_123","version":4,"status":"ready_for_delivery_check"}事件类型要和任务状态图对齐。前端收到 stage 就更新进度,收到 source 就展示资料,收到 review_finding 就展示待处理问题,收到 token 就追加产物。业务判断仍由服务端完成,例如「缺少来源引用」是否阻塞交付,不应该由前端根据文案自行判断。
4. 服务端职责
服务端流式接口至少要做几件事:
- 创建或读取任务状态。
- 启动 Agent 编排流程。
- 把节点事件转换成 SSE。
- 在错误时发送结构化错误事件。
- 按任务版本记录事件序号,支持断线恢复。
- 在响应结束后继续执行必要的后台写回。
简化结构如下:
typescript
app.post("/agent-tasks/:id/stream", async (c) => {
const stream = new ReadableStream({
async start(controller) {
const send = (id: string, event: string, data: unknown) => {
controller.enqueue(
new TextEncoder().encode(
`id: ${id}\nevent: ${event}\ndata: ${JSON.stringify(data)}\n\n`,
),
);
};
try {
for await (const update of runAgentGraph(c.req.param("id"))) {
send(update.id, update.type, update.payload);
}
send("final", "done", { status: "completed" });
} catch (error) {
send("error", "error", { code: "STREAM_FAILED" });
} finally {
controller.close();
}
},
});
return new Response(stream, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache, no-transform",
"Connection": "keep-alive",
"X-Accel-Buffering": "no",
},
});
});这段示例为了可读性省略了细节。真实项目还需要处理:
- 鉴权:在建立流之前校验用户是否有权访问该任务。
- 取消与断开:检测客户端是否断开(如 FastAPI 的
request.is_disconnected()),一旦断开就取消上游模型调用,避免继续为已被丢弃的连接付费或占用 GPU 槽位。 - 背压:对慢客户端使用有界队列,让 token 生产速度被消费速度自然限制,防止单连接内存无限增长。
- 心跳:在长工具调用或模型 prefill 阶段没有 token 产出时,每 15–30 秒发送一次 SSE 注释行
: heartbeat\n\n,保持代理和负载均衡的连接温热。 - 重试与恢复:按
taskId、version和事件序号记录事件,前端重连时带上Last-Event-ID,服务端决定补发还是继续。 - Trace 与持久化:把节点事件写入事件日志,用于后续审稿、排障和质量分析。
这里的关键点是:流式接口需要把任务图里的节点事件映射成稳定的前端契约,避免把 Agent 输出原样暴露给浏览器。
5. 前端职责
前端不要推断业务状态。它只消费服务端事件,并把事件落到对应 UI 区域:
stage:更新任务进度。source:展示资料卡片。plan:展示任务结构。token:追加产物。review_finding:展示审稿问题。error:展示可恢复或不可恢复错误。done:结束流并刷新任务详情。
这样前端保持简单,业务判断集中在 Agent 编排和任务状态里。前端仍然需要维护本地渲染状态,例如当前部分、已接收事件 ID、连接状态和重连次数,但这些状态只服务展示和恢复,不参与交付资格判断。
6. 错误和幂等
流式响应比普通请求多几类问题:
- 连接中断:用户只收到部分产物。
- 用户重试:后台任务可能已经继续执行。
- 工具超时:资料召回或审稿节点没有完成。
- 部分成功:产物已经生成,但交付检查失败。
因此任务必须有 taskId、version 和事件序号。重试时前端应该恢复已有状态,而不是盲目重新生成。后台写回也要尽量幂等,例如同一个 sourceId 不重复写入,同一个审稿问题不重复创建,同一个部分版本不会被旧 token 覆盖。
错误事件也要区分类型:
| 类型 | 前端处理 | 服务端处理 |
|---|---|---|
NETWORK_INTERRUPTED | 显示重连状态,保留产物 | 根据最后事件 ID 继续推送 |
TOOL_TIMEOUT | 标记节点可重试 | 记录 trace,允许重跑该节点 |
SOURCE_INSUFFICIENT | 提示补充资料 | 停在人工确认或资料补充节点 |
REVIEW_BLOCKED | 展示审稿阻塞项 | 保留产物,等待改写或确认 |
这些错误不应该都折叠成「生成失败」。Agent 任务里,很多失败是可继续处理的中间状态。
7. 可观测性和审稿联动
Streaming 还会影响后续审稿和排障。只保存最终产物,无法解释某个段落为什么引用了某个来源,也无法复盘审稿 Agent 是在哪一步提出阻塞意见。从可观测性角度,至少应记录以下指标:
- TTFT(首 token 时间):从请求发起到第一个
data:帧刷出的时间,决定用户感知的快慢。 - Inter-token latency:相邻 token 之间的间隔,如果突然升高,通常是模型、队列或网络出现 stall。
- Tokens per second:用户实际接收到的持续解码速率,受网络、背压和前端渲染影响,可能低于模型原生速率。
- Stream completion rate:有多少流正常到达
done/[DONE],有多少被用户放弃或断开。这个指标直接反映取消/断开处理是否生效,也关系到推理成本。
更稳定的做法是同时记录三类信息:
- 事件日志:每个节点发出了什么事件,事件顺序是什么。
- 内容版本:每个部分在执行、审稿、改写后的版本差异。
- 运行 trace:每个 Agent 调用了哪些工具,耗时、错误和重试次数是多少。
这些信息可以用于交付前检查,也可以用于后续质量分析。例如某类产物经常在 citation_needed 阶段阻塞,排查重点应先放到资料召回规则或来源过滤策略,再检查执行 Agent 的引用生成方式。
8. 小结
Agent 项目里的 Streaming,核心是让长任务的状态、产物、资料和审稿结果一起可见。SSE 足够覆盖大多数单向进度场景;前端负责渲染事件和恢复连接,服务端负责状态判断、错误分类、事件持久化和后台写回。
参考资料
- MDN: Using server-sent events — SSE 事件格式、
id、retry和Last-Event-ID的浏览器行为说明。 - MDN: EventSource — 浏览器原生 SSE 客户端 API、连接状态、自动重连限制(GET 与自定义头约束)。
- HTML Living Standard: Server-sent events — SSE 协议的官方规范,包括事件流格式和重连规则。
- OpenAI API: Streaming API responses — OpenAI 基于 SSE 的语义化事件流说明,可作为 LLM/Agent 流式接口设计的参考。
- Vercel AI SDK: Chatbot — 生产级 React 流式聊天 UI 的
useChat+streamText实现参考。