我们的内部 Agent 平台在推理时需要把 LLM 的 token、思考过程、工具调用、知识库引用、子 Agent 协作状态实时推给前端。本文记录我们在 SSE(Server-Sent Events)之上设计的一套紧凑类型化事件协议,以及围绕"流要真的流起来"踩过的坑与最终方案。
1. 为什么不用裸 SSE data: 帧,而是自定义 JSON 协议
标准 SSE 的 data: + event: 帧模型对人类可读、对浏览器 EventSource 友好,但我们的场景有三个特殊诉求:
- 事件种类多且会持续增长:文本 token、思考 token、工具开始/结果、引用、结构化结果、子 Agent 生命周期……用一个
type字段分发比维护一堆event:名字更直白。 - token 流量大,字节要省:一条对话可能产生上千条 token 事件。我们把字段名压成单字母(
c/i/n/a/r)、事件类型压成 1~2 个字符(t/th/ts/tr),相比{"type":"text","content":"..."}省一半以上体积。 - 前端不走 EventSource:我们用 fetch + ReadableStream 自己读流(需要自定义 Header 鉴权),所以不依赖
data:前缀解析,约定"每条消息一个独立 JSON 对象,消息间以\n\n分隔"即可。
传输层就是普通的 HTTP 长连接流式响应:Content-Type: text/event-stream,服务端单连接推送所有事件,前端按 type 分发。
2. 事件类型设计
2.1 基础聊天事件
| 类型 | 含义 | 关键字段 |
|---|---|---|
t |
文本 token(流式) | c 文本片段 |
th |
思考 token(reasoning 模型) | c 思考片段 |
ts |
工具调用开始 | i 工具 ID / n 工具名 / a 参数 JSON |
tr |
工具调用结果 | i / n / r 结果 JSON |
kr |
知识库引用 | s 引用列表 |
r |
结构化结果帧 | payload 业务 JSON |
d |
主 Agent 流结束 | — |
h |
心跳(每 10s) | — |
e |
错误 | m 错误描述 |
几个设计要点:
文本与思考分流(t vs th)。 DeepSeek 这类 reasoning 模型会同时输出 reasoning_content 和 content 两个通道。我们不下发"开始思考/结束思考"的边界事件,而是让前端按类型各自累积渲染——思考默认折叠。省掉边界事件的代价是前端逻辑略复杂,收益是协议无状态、断线重连后任意一条事件都可独立处理。
{"c":"你好","type":"t"}
{"c":"用户想生成一张小猫图片,我先理解需求...","type":"th"}
工具调用拆成开始/结果两条(ts/tr)。 工具可能执行数秒到数十秒(如生图)。ts 到达时前端立刻渲染"正在调用 XX 工具"的占位卡片,tr 到达后填充结果。参数 a 和结果 r 都是 JSON 字符串而非内嵌对象——它们是"透传载荷",前端只在需要展示时才 parse,避免服务端序列化/前端反序列化两层都理解工具的内部结构。
{"i":"bt_generate_image","n":"generate_image","a":"{\"prompt\":\"一只可爱的小猫\"}","type":"ts"}
{"i":"bt_generate_image","n":"generate_image","r":"{\"type\":\"image\",\"urls\":[\"https://cdn.example.com/gen_abc123_0.jpg\"]}","type":"tr"}
媒体类工具的结果约定含 "type":"image"|"video" + "urls",前端据此渲染媒体组件——这是协议里为数不多的"载荷结构约定",值得用文档钉死。
引用独立成帧(kr)。 知识库检索结果不塞进 tr(工具结果),而是单独发 kr。原因:一次检索的引用列表要驱动前端"来源面板"这种持久 UI,与一次性的工具执行日志是两回事;而且引用可能来自 Agent 自动注入的 KB 上下文,根本没有对应的工具调用。
{"s":[{"docId":"doc_abc123","docName":"2026Q2行业报告.pdf","kbName":"行业数据","score":0.87}],"type":"kr"}
结构化结果帧(r)。 对开启 JSON 输出的 Agent(JSON-contract agent),服务端从流式文本中抽取最终 JSON,在 d 之前一次性下发。t 流仍然照常发(用户能看到生成过程),r 是"提炼后的契约"。这样前端 stepper UI 不用自己从 Markdown 代码块里抠 JSON。
2.2 子 Agent 协作事件
多 Agent 场景下,主 Agent 通过 send_message 工具把任务派给子 Agent(消息总线投递,异步执行)。这部分事件的设计目标是:让用户看见"一个团队在干活",同时协议本身不引入等待语义。
| 类型 | 含义 | 时机 |
|---|---|---|
as |
子 Agent 已加载(懒加载) | 主 Agent 首次向它发任务时 |
am |
Agent 间消息 | 任意 Agent 调 send_message(task / clarify_resp / revision / shutdown) |
ap |
子 Agent 进度汇报 | 执行中 |
ac |
子 Agent 请求澄清 | 对任务有疑问时 |
ar |
子 Agent 交付结果 | 完成时,含 summary/steps/durationMs/media |
ax |
子 Agent 退出 | 完成或超时 |
消息类事件(am/ap/ac/ar)共享一套信封字段:msgId、inReplyTo、from、to、content、ts。inReplyTo 串起对话链,前端可以按任务线渲染线程视图。
{"msgId":"msg_001","from":"agent_main","to":"agent_b","type":"task","content":"画一只可爱的小猫","ts":"2026-07-08T10:00:00Z","type":"am"}
{"msgId":"msg_002","inReplyTo":"msg_001","from":"agent_b","to":"agent_main","type":"progress","content":"已收到任务,正在处理...","ts":"2026-07-08T10:00:01Z","type":"ap"}
{"msgId":"msg_003","inReplyTo":"msg_001","from":"agent_b","to":"agent_main","type":"result","summary":"已生成一只可爱的小猫","steps":1,"durationMs":12345,"media":["https://cdn.example.com/cat.jpg"],"ts":"2026-07-08T10:00:15Z","type":"ar"}
{"agentId":"agent_b","reason":"completed","type":"ax"}
2.3 时序规则(容易写错的部分)
d由主 Agent 的 LLM 流结束触发,不等待子 Agent。 这是"fire-and-forget"协作模型的必然结果:主 Agent 发完任务后继续对用户说话,说完就d;子 Agent 的结果下一轮再带回来。- 协作事件(
am/ap/ar/ac/ax)可以出现在d之后,连接在所有子 Agent 退出后才关闭。前端收到d不能立刻断开。 as先于ap(懒加载语义)。ar至少一条,可能多条(LLM 透传的 send_message 结果 + 子 Agent runner 的最终交付),前端按inReplyTo去重/合并。
一条完整的主+子 Agent 流示例:
{"c":"我","type":"t"}
{"c":"已委托","type":"t"}
{"c":"助手","type":"t"}
{"c":"为您画图","type":"t"}
{"c":"。","type":"t"}
{"i":"bt_send_message","n":"send_message","a":"{\"to\":\"agent_b\",\"type\":\"task\",\"content\":\"画一只可爱的小猫\"}","type":"ts"}
{"msgId":"msg_001","from":"agent_main","to":"agent_b","type":"task","content":"画一只可爱的小猫","ts":"2026-07-08T10:00:00Z","type":"am"}
{"i":"bt_send_message","n":"send_message","r":"{\"sent\":true,\"msgId\":\"msg_001\"}","type":"tr"}
{"agentId":"agent_b","agentName":"图片助手","status":"loaded","subSessionId":"sub_abc123","type":"as"}
{"c":"请稍等片刻","type":"t"}
{"c":"。","type":"t"}
{"type":"d"}
{"msgId":"msg_002","inReplyTo":"msg_001","from":"agent_b","to":"agent_main","type":"progress","content":"已收到任务,正在处理...","ts":"2026-07-08T10:00:01Z","type":"ap"}
{"msgId":"msg_003","from":"agent_b","to":"agent_main","type":"result","summary":"已生成一只可爱的小猫","steps":1,"durationMs":12345,"media":["https://cdn.example.com/cat.jpg"],"ts":"2026-07-08T10:00:15Z","type":"ar"}
{"agentId":"agent_b","reason":"completed","type":"ax"}
3. 让流真的流起来:批处理效应与"三件套"
协议设计完,遇到的第一个工程问题是:上游模型服务逐字推送,但前端收到的事件是"攒一批→瞬间到达",字符一次性弹出一大段。
3.1 逐层排查(排除了所有"缓冲"假设)
| 层 | 组件 | 干预 | 结论 |
|---|---|---|---|
| 客户端读流 | BufferedReader |
8KB → 256B → 1B → 去掉 | 无改善 |
| 字节→字符解码 | InputStreamReader |
换逐字节读 | 无改善 |
| HTTP 连接 | HttpURLConnection 内部缓冲 |
换 java.net.http.HttpClient |
无改善 |
| Reactor 队列 | Flux.merge(prefetch=256) |
改 1 | 无改善 |
| Netty 写入 | ChannelOutboundBuffer |
调 SO_SNDBUF / 水位线 |
无改善 |
| Spring SSE Writer | ServerSentEventHttpMessageWriter |
换 writeAndFlushWith() |
返回类型被框架吞掉 |
最终结论:不是"缓冲"问题,是 Reactor + Netty 的批处理效应。 即使每个事件都调了 sink.next() 和 channel.flush(),当多个事件在同一 event loop 周期内到达时,Netty 会把它们合并进一次 TCP 写;Reactor 的 merge/flatMap 也会在一个 drain 循环里连续推多个事件。应用层唯一能可靠控制的是事件发射的时间间隔。
3.2 三件套
- 输入端小缓冲:读流
BufferedReader给 64B 字符缓冲,避免首次read()预读大量 SSE 行瞬间涌入下游。 - 时间间隔而非缓冲控制:对内容流加 15~30ms 随机间隔的发射节流,让 Netty 在事件之间有空闲独立发送。
- 逐级消除 Reactor 内部攒批:
share()内部 prefetch=256,会把节流后的事件重新攒批——改用publish(1).refCount(1),缓冲限为 1,到达即转发。
var chatFlux = buildChatFlux(...)
.concatMap(ev -> Mono.just(ev).delayElement(
Duration.ofMillis(ThreadLocalRandom.current().nextLong(15, 31)),
Schedulers.boundedElastic()), 32)
.publish(1).refCount(1); // 非 share()
// 延迟和缓冲控制只作用于内容流;心跳与尾随事件保持即时
return Flux.merge(1, chatFlux, heartbeat, tailEvents);
两条设计原则值得强调:
- 应用层解决,不在 Netty/TCP 层强求。调 socket 参数、改水位线都是与实现细节对赌,换一个版本就失效;
delayElement简单、可控、语义明确。 - 纯透传代理不需要这套。如果服务端只是把上游字节原样写回 socket(
SseEmitter.send()→ flush),没有 Reactor 中间层,就没有批处理效应。我们的平台需要解析 chunk、走工具调用链、重新组装事件,路径完全不同,方案不能照抄。
补充一个事后才知道的坑:上面"三件套"里的节流后来被发现会放大另一个潜伏 bug——
Flux.createconsumer 内同步阻塞读导致的背压饿死(request 信号被排队在阻塞任务之后),表现为长回复轮次精确卡在 256 个事件后冻结。教训是:节流可以上,但要确认背压链路本身是健康的,否则"让消费变慢"只会更快耗尽预取额度。详见同目录的事故复盘文章。
4. 错误收尾与容错
流式响应的错误处理有一个刚性约束:响应头一旦 committed,就不可能再用 HTTP 状态码或错误体表达失败,协议内的事件是唯一的通信渠道。我们的规则:
- 错误以
e事件下发,随后正常关闭流。{"m":"推理过程出错: Connection timeout","type":"e"}。服务端不再向已 committed 的 SSE 响应传播异常(传播了前端也收不到,只会让框架刷一条无意义日志),而是把事件收集器和 emitter 都优雅 complete。 - 工具层容错不打断整条流。 模型生成非法 JSON 参数是高频事件(尤其长参数时)。工具回调外包一层容错壳:参数解析失败时把错误文本作为工具结果还给模型,让它自我修正重试,而不是让整个会话崩掉。LLM 工具调用本质上是"会犯错的远端调用",按可重试错误处理,不按异常处理。
- 心跳兜底代理超时。
h事件每 10 秒一条,防止中间代理/CDN 因长时间无数据断开连接——长思考或慢工具时,内容流可能静默几十秒,心跳必须独立于内容流发射(所以上节Flux.merge里 heartbeat 不吃节流)。 d与连接关闭分离。 前端以"连接关闭 + 所有子 Agent 的ax"判定会话真正结束,d只是主输出结束。这把"流结束"从一个单点事件变成一个可推导的状态,对断线重连也更友好。
5. 小结
这套协议的核心权衡可以归纳为三条:
- 紧凑优于自描述:单字母字段名牺牲了一点可读性,换来 token 流量的真实带宽收益;用文档(本文)补足可读性。
- 事件无状态、时序可推导:不发"开始/结束"成对边界事件,每条事件独立可处理;结束状态由
d+ax推导而非单一事件承载。 - 流式体验是应用层责任:批处理效应不靠调 TCP 参数对抗,用发射节流 + 消除中间攒批解决;同时警惕节流掩盖背压链路的健康度问题。