跳到主要内容

流式协议与中断

核对日期:2026-08-26。

1. 定义与边界

流式输出(Streaming)是把模型生成过程按增量推到客户端,而不是等完整响应再渲染。对 Agent 而言,流里不只是文字:还包括推理块、工具调用开始、参数增量、审批中断、错误和 run 结束。

本文件只处理传输协议、中断和半成品持久化。工具卡片状态机见 流式Tool-Use与前端状态.md;审批卡片见 HITL交互界面.md;微信 / 支付宝差异见 跨端流式差异.md

不适用:

  • 语音双向实时(应走 WebRTC / Voice 专线,不是文本 SSE)。
  • 把整段聊天历史当唯一状态(状态仍应落在服务端 run store)。

2. 为什么重要

同样的模型和工具,流式实现会决定产品是否可用:

  • 用户在 20–40 秒无反馈时会认为系统卡死。
  • Agent 一步里可能先写字、再调工具、再继续写字;非流式会把这些折叠成一次「转圈」。
  • 用户点停止后,如果服务端还在调模型和写工具,成本和副作用都停不下来。
  • 断线后如果没有 run_id,前端无法续上,只能重跑。

3. 核心机制

推荐把浏览器协议和模型厂商协议分开:

协议职责
模型厂商OpenAI / Anthropic SSEtoken、tool 参数增量
Agent 网关自有 AgentUiEvent与前端契约,屏蔽厂商差异
客户端fetch + ReadableStream 解析 SSEPOST、鉴权、abort

不要让 React / Vue 组件直接消费 Anthropic 的 content_block_delta。厂商事件会变,前端状态机不应绑死。

4. 架构模式

4.1 协议选择

协议适用不适用
POST + SSE(fetch 读流)默认:文本 Agent、工具流、鉴权 header需要服务端主动推送无关事件
EventSource公开、只读、GET 可缓存的演示需要 POST body、Authorization header
WebSocket用户频繁打断、多人协同、语音控制面纯单向 token 流;会增加网关和观测复杂度
轮询 / 长轮询小程序基础库不支持分块时的降级作为 H5 默认方案

EventSource 只能 GET,不能带复杂 JSON body,也难带自定义鉴权头。生产 Agent 几乎都是 POST /runs。因此 H5 默认应是 fetch + 解析 text/event-stream,而不是 new EventSource

4.2 事件契约

type RiskLevel = 'low' | 'medium' | 'high';

type AgentUiEvent =
| {
type: 'run_start';
runId: string;
traceId: string;
sessionId: string;
}
| {
type: 'text_delta';
runId: string;
messageId: string;
delta: string;
}
| {
type: 'tool_call_start';
runId: string;
toolCallId: string;
toolName: string;
riskLevel: RiskLevel;
}
| {
type: 'tool_args_delta';
runId: string;
toolCallId: string;
partialJson: string;
}
| {
type: 'approval_required';
runId: string;
approvalId: string;
toolCallId: string;
}
| {
type: 'run_error';
runId: string;
code: string;
retryable: boolean;
message: string;
}
| {
type: 'run_end';
runId: string;
status: 'completed' | 'aborted' | 'failed' | 'waiting_approval';
};

事件必须带 runId。前端用它做幂等合并:同一 toolCallId 的多次 tool_args_delta 是追加,不是新卡片。

5. 工程实现

5.1 服务端:把 abort 传到模型

Vercel AI SDK 的正确做法是把 HTTP 请求的 AbortSignal 传给 streamText,并在 onAbort 里落半成品,而不是只在浏览器停渲染。

import { streamText } from 'ai';

export async function POST(req: Request): Promise<Response> {
const payload = (await req.json()) as { goal: string; sessionId: string };

const result = streamText({
model: getModel(),
prompt: payload.goal,
abortSignal: req.signal,
onAbort: ({ steps }) => {
persistPartialRun({
sessionId: payload.sessionId,
steps,
status: 'aborted',
});
},
});

return result.toUIMessageStreamResponse();
}

useChat()stop() 只取消浏览器读取。若服务端不转发 req.signal,模型调用和后续工具仍会跑完。

5.2 Vue 3:POST 流 + 可停止

import { onBeforeUnmount, ref } from 'vue';

interface RunStartEvent {
type: 'run_start';
runId: string;
traceId: string;
}

export function useAgentStream() {
const runId = ref<string | undefined>(undefined);
const text = ref('');
const status = ref<'idle' | 'streaming' | 'aborted' | 'ended'>('idle');
let abort: AbortController | undefined;

async function start(goal: string): Promise<void> {
abort?.abort();
abort = new AbortController();
text.value = '';
status.value = 'streaming';

const response = await fetch('/api/agent/runs', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ goal }),
signal: abort.signal,
});

if (!response.ok || response.body === null) {
status.value = 'ended';
throw new Error(`stream_http_${response.status}`);
}

const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = '';

while (true) {
const { done, value } = await reader.read();
if (done) {
break;
}
buffer += decoder.decode(value, { stream: true });
const frames = buffer.split('\n\n');
buffer = frames.pop() ?? '';
for (const frame of frames) {
const dataLine = frame.split('\n').find((line) => line.startsWith('data:'));
if (dataLine === undefined) {
continue;
}
const payload = JSON.parse(dataLine.slice(5).trim()) as RunStartEvent | { type: string; delta?: string; runId?: string };
if (payload.type === 'run_start') {
runId.value = payload.runId;
}
if (payload.type === 'text_delta' && payload.delta !== undefined) {
text.value += payload.delta;
}
}
}

status.value = 'ended';
}

function stop(): void {
abort?.abort();
status.value = 'aborted';
}

onBeforeUnmount(() => abort?.abort());

return { runId, text, status, start, stop };
}

React 侧可用 @ai-sdk/reactuseChat({}).stop;Vue 侧可用 @ai-sdk/vueuseChat。自建网关时不要假设 SDK 事件名稳定,仍应映射到 AgentUiEvent

5.3 断线续读

流中断不等于 run 失败。服务端应:

  1. run_id 持久化当前文本、工具草稿、已执行工具结果。
  2. 提供 GET /runs/:id/events?after=seq 续传。
  3. 前端重连只补事件,不重发 goal。

未完成的工具写操作必须看幂等键,见 ../09-Agent工程化/幂等性设计.md

6. 生产实践

实践说明
先出 run_start用户一点发送就能复制 run / 报障
心跳长推理阶段每 15 秒发 comment / ping,避免代理切断
超时分层网关 120s、模型 60s、单工具 15s,分开告警
停止即审计记录谁停止、停在第几步、是否已产生副作用
半成品可见中断后仍展示已生成文字和未执行的工具草稿
代理兼容关闭 Nginx proxy_buffering;SSE 禁用响应压缩或调好 X-Accel-Buffering

7. 常见反模式

反模式表现后果修正
只用 EventSourceGET 流、token 放 query泄密、无法传工具上下文POST + fetch
前端停、后端不停stop()reader.cancel()继续烧 token、可能写库转发 abortSignal
run_id断线后只能重开对话重复副作用首包下发 run / trace
把厂商 SSE 暴露给客户端前端 switch Anthropic 事件换模型就改 UI网关翻译成 AgentUiEvent
无心跳推理 30 秒无字节负载均衡 idle timeoutcomment 帧或 ping
中断后重放整段 prompt当作用户再发一次重复转账 / 重复建单续传事件,不重跑写工具

8. 评测方法

指标口径
Time to First Event从 POST 到首个 run_starttext_delta
Time to First Tool Card从 POST 到首个 tool_call_start
Abort Propagation点击停止后,模型 span 是否在 SLA 内结束
Partial Persist Rate中断 run 中已落库半成品比例
Reconnect Success断线后用 run_id 续上且不重复写工具
Stalled Stream Rate超过心跳间隔仍无事件的比例

评测不要只看最终文本。一条「答案正确但无法停止」的流,生产上不可接受。

9. 安全与治理

  • 流式响应同样要鉴权;run_id 不可猜测。
  • 不要把工具原始返回整段推到浏览器,只推摘要和允许展示的字段。
  • 中断不能当成安全控制:用户没点停止,高风险工具仍须审批。
  • 错误事件对用户用稳定 code,对内部 trace 才放堆栈。
  • 流里的外部文本按不可信内容处理,渲染规则见 Generative-UI工程.md

10. 权威资料