Skip to content

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 请求头,服务端可以据此判断从哪个事件继续补发。这个机制让前端在断线后恢复已有状态,而不必重新发起整个任务。

需要注意三个工程细节:

  1. 浏览器原生的 EventSource 只支持 GET 请求,且不能设置自定义请求头(如 Authorization)。因此带身份认证或需要 POST 请求体的 Agent 接口,通常用 fetch + ReadableStream 手动解析 SSE 帧,自己在客户端实现重连和 Last-Event-ID 管理。
  2. 代理和负载均衡经常会缓冲响应。Nginx、Cloudflare 或 AWS ALB 默认可能把整段流收集后再转发,导致用户看到「最后一次性出现」。服务端需要设置 Content-Type: text/event-streamCache-Control: no-cacheConnection: keep-aliveX-Accel-Buffering: no,并在代理层关闭 buffering,必要时用 SSE 注释行(: heartbeat\n\n)每 15–30 秒发送一次心跳,避免空闲超时。
  3. 慢客户端或高 token 速率下会产生背压(backpressure)。如果服务端一股脑把事件塞进内存,而客户端来不及消费,内存会无界增长。生产环境应通过有界队列或检测 drain 事件,让生产速度匹配消费速度。

3. 事件类型设计

只推送 token 会把业务状态藏在产物流里,前端也会被迫从文本变化里猜测任务进度。Agent 任务更适合把事件分层:

  • 任务级事件:stagecheckpointdoneerror
  • 资料级事件:sourcesource_rejectedcitation_needed
  • 内容级事件:plansection_startedtokensection_completed
  • 审稿级事件:review_findingrevision_startedrevision_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. 服务端职责

服务端流式接口至少要做几件事:

  1. 创建或读取任务状态。
  2. 启动 Agent 编排流程。
  3. 把节点事件转换成 SSE。
  4. 在错误时发送结构化错误事件。
  5. 按任务版本记录事件序号,支持断线恢复。
  6. 在响应结束后继续执行必要的后台写回。

简化结构如下:

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,保持代理和负载均衡的连接温热。
  • 重试与恢复:按 taskIdversion 和事件序号记录事件,前端重连时带上 Last-Event-ID,服务端决定补发还是继续。
  • Trace 与持久化:把节点事件写入事件日志,用于后续审稿、排障和质量分析。

这里的关键点是:流式接口需要把任务图里的节点事件映射成稳定的前端契约,避免把 Agent 输出原样暴露给浏览器。

5. 前端职责

前端不要推断业务状态。它只消费服务端事件,并把事件落到对应 UI 区域:

  • stage:更新任务进度。
  • source:展示资料卡片。
  • plan:展示任务结构。
  • token:追加产物。
  • review_finding:展示审稿问题。
  • error:展示可恢复或不可恢复错误。
  • done:结束流并刷新任务详情。

这样前端保持简单,业务判断集中在 Agent 编排和任务状态里。前端仍然需要维护本地渲染状态,例如当前部分、已接收事件 ID、连接状态和重连次数,但这些状态只服务展示和恢复,不参与交付资格判断。

6. 错误和幂等

流式响应比普通请求多几类问题:

  • 连接中断:用户只收到部分产物。
  • 用户重试:后台任务可能已经继续执行。
  • 工具超时:资料召回或审稿节点没有完成。
  • 部分成功:产物已经生成,但交付检查失败。

因此任务必须有 taskIdversion 和事件序号。重试时前端应该恢复已有状态,而不是盲目重新生成。后台写回也要尽量幂等,例如同一个 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],有多少被用户放弃或断开。这个指标直接反映取消/断开处理是否生效,也关系到推理成本。

更稳定的做法是同时记录三类信息:

  1. 事件日志:每个节点发出了什么事件,事件顺序是什么。
  2. 内容版本:每个部分在执行、审稿、改写后的版本差异。
  3. 运行 trace:每个 Agent 调用了哪些工具,耗时、错误和重试次数是多少。

这些信息可以用于交付前检查,也可以用于后续质量分析。例如某类产物经常在 citation_needed 阶段阻塞,排查重点应先放到资料召回规则或来源过滤策略,再检查执行 Agent 的引用生成方式。

8. 小结

Agent 项目里的 Streaming,核心是让长任务的状态、产物、资料和审稿结果一起可见。SSE 足够覆盖大多数单向进度场景;前端负责渲染事件和恢复连接,服务端负责状态判断、错误分类、事件持久化和后台写回。

参考资料

基于 MIT 协议开源