我们的内部 Agent 平台在推理时需要把 LLM 的 token、思考过程、工具调用、知识库引用、子 Agent 协作状态实时推给前端。本文记录我们在 SSE(Server-Sent Events)之上设计的一套紧凑类型化事件协议,以及围绕"流要真的流起来"踩过的坑与最终方案。


1. 为什么不用裸 SSE data: 帧,而是自定义 JSON 协议

标准 SSE 的 data: + event: 帧模型对人类可读、对浏览器 EventSource 友好,但我们的场景有三个特殊诉求:

  1. 事件种类多且会持续增长:文本 token、思考 token、工具开始/结果、引用、结构化结果、子 Agent 生命周期……用一个 type 字段分发比维护一堆 event: 名字更直白。
  2. token 流量大,字节要省:一条对话可能产生上千条 token 事件。我们把字段名压成单字母(c/i/n/a/r)、事件类型压成 1~2 个字符(t/th/ts/tr),相比 {"type":"text","content":"..."} 省一半以上体积。
  3. 前端不走 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_contentcontent 两个通道。我们不下发"开始思考/结束思考"的边界事件,而是让前端按类型各自累积渲染——思考默认折叠。省掉边界事件的代价是前端逻辑略复杂,收益是协议无状态、断线重连后任意一条事件都可独立处理。

{"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)共享一套信封字段:msgIdinReplyTofromtocontenttsinReplyTo 串起对话链,前端可以按任务线渲染线程视图。

{"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 时序规则(容易写错的部分)

  1. d 由主 Agent 的 LLM 流结束触发,不等待子 Agent。 这是"fire-and-forget"协作模型的必然结果:主 Agent 发完任务后继续对用户说话,说完就 d;子 Agent 的结果下一轮再带回来。
  2. 协作事件(am/ap/ar/ac/ax)可以出现在 d 之后,连接在所有子 Agent 退出后才关闭。前端收到 d 不能立刻断开。
  3. as 先于 ap(懒加载语义)。
  4. 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 三件套

  1. 输入端小缓冲:读流 BufferedReader 给 64B 字符缓冲,避免首次 read() 预读大量 SSE 行瞬间涌入下游。
  2. 时间间隔而非缓冲控制:对内容流加 15~30ms 随机间隔的发射节流,让 Netty 在事件之间有空闲独立发送。
  3. 逐级消除 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.create consumer 内同步阻塞读导致的背压饿死(request 信号被排队在阻塞任务之后),表现为长回复轮次精确卡在 256 个事件后冻结。教训是:节流可以上,但要确认背压链路本身是健康的,否则"让消费变慢"只会更快耗尽预取额度。详见同目录的事故复盘文章。

4. 错误收尾与容错

流式响应的错误处理有一个刚性约束:响应头一旦 committed,就不可能再用 HTTP 状态码或错误体表达失败,协议内的事件是唯一的通信渠道。我们的规则:

  1. 错误以 e 事件下发,随后正常关闭流。 {"m":"推理过程出错: Connection timeout","type":"e"}。服务端不再向已 committed 的 SSE 响应传播异常(传播了前端也收不到,只会让框架刷一条无意义日志),而是把事件收集器和 emitter 都优雅 complete。
  2. 工具层容错不打断整条流。 模型生成非法 JSON 参数是高频事件(尤其长参数时)。工具回调外包一层容错壳:参数解析失败时把错误文本作为工具结果还给模型,让它自我修正重试,而不是让整个会话崩掉。LLM 工具调用本质上是"会犯错的远端调用",按可重试错误处理,不按异常处理。
  3. 心跳兜底代理超时。 h 事件每 10 秒一条,防止中间代理/CDN 因长时间无数据断开连接——长思考或慢工具时,内容流可能静默几十秒,心跳必须独立于内容流发射(所以上节 Flux.merge 里 heartbeat 不吃节流)。
  4. d 与连接关闭分离。 前端以"连接关闭 + 所有子 Agent 的 ax"判定会话真正结束,d 只是主输出结束。这把"流结束"从一个单点事件变成一个可推导的状态,对断线重连也更友好。

5. 小结

这套协议的核心权衡可以归纳为三条:

  • 紧凑优于自描述:单字母字段名牺牲了一点可读性,换来 token 流量的真实带宽收益;用文档(本文)补足可读性。
  • 事件无状态、时序可推导:不发"开始/结束"成对边界事件,每条事件独立可处理;结束状态由 d + ax 推导而非单一事件承载。
  • 流式体验是应用层责任:批处理效应不靠调 TCP 参数对抗,用发射节流 + 消除中间攒批解决;同时警惕节流掩盖背压链路的健康度问题。