Ragent流式输出:SSE分事件推送思考、正文与来源,以及跨节点流式取消
【免费下载链接】ragent企业级 Agentic RAG 智能体 - 全链路覆盖文档解析、多路检索、意图识别、问题重写、会话记忆、MCP 工具调用与深度思考。面向真实业务场景,从 0 到 1 完整工程实现。项目地址: https://gitcode.com/gh_mirrors/ragent1/ragent
Ragent 是一个企业级 Agentic RAG 智能体,它的对话界面之所以能做到“边想边说、边查边流”,核心是一套SSE(Server-Sent Events)流式输出协议:思考、正文、检索来源、工具进度被拆成 8 种命名事件分帧推送;配合 Redis 标记 + 广播主题,跨节点流式取消让用户在任意一台服务节点上点“停止”,都能可靠掐断正在运行的 Agent 流。这篇文章带你从协议到取消链路,完整看懂这套机制。
为什么选 SSE 而不是轮询或 WebSocket?
大模型回答动辄几十秒,如果等整个回答生成完再返回,用户只能干等。SSE 基于 HTTP 长连接、单向推送,天然适合“服务端持续产出增量、客户端实时渲染”的场景,且断线重连、事件命名等浏览器原生支持,实现成本远低于 WebSocket。
Ragent 的对话入口就返回一个 SSE 流,超时时间可配:
- 提问入口:AgentChatController.java 中
GET /agent/v1/chat返回SseEmitter - 停止入口:
POST /agent/v1/stop?taskId=...,只需一个任务 ID
上图右侧“SSE 流式通道”和底部“中断与取消、统一收尾”正是本文主角:所有增量帧从 Agent 运行时一路推到接入层,取消信号则从基础设施层反向打断整条链路。
8 种 SSE 事件:思考、正文、来源各走各的通道
Agent 引擎的事件协议定义在 AgentSSEEventType.java,一共产生 8 类事件:
| 事件 | 作用 | 典型内容 |
|---|---|---|
meta | 开场元信息 | conversationId、taskId(停止按钮就靠它) |
message | 增量文本 | type区分think(思考)/response(正文)/error(中断提示) |
block | 文本块封口 | 携带服务端起止时间,用于计算“思考了多久” |
tool | 工具进度 | 工具名、状态(pending/running/end)、结果摘要 |
hint | 运行提示 | 如“已达最大迭代次数,正在生成总结” |
confirm | 写操作确认 | 等待用户批准敏感工具调用 |
finish | 本轮结束 | 消息 ID、状态(NORMAL / INTERRUPTED)、总耗时 |
done | 连接收尾 | 固定[DONE],前端据此关流 |
cancel | 用户取消 | 与 finish 互斥,同样携带落库后的消息状态 |
关键设计是分事件推送而不是一锅烩:
- 思考与正文分离:
think和response都是message事件,但type不同。前端据此把“深度思考”渲染成折叠区,正文渲染成气泡,互不污染。转换逻辑在 AgentStreamEventBridge.java 的onThinkingDelta/onResponseDelta。 - 来源随 finish 交付:RAG 模式下,检索命中的文档级来源(
sources)不是单独一帧,而是随结束载荷一起下发,见 CompletionPayload.java 中的List<SourceRef> sources字段,仅在命中知识库时携带。 - 工具进度独立成帧:
tool事件携带状态机(pending → running → end),前端能画出“正在调用 search_knowledge…”的进度条。
发送侧被封装成线程安全的 SseEmitterSender.java:连接关闭状态用AtomicBoolean守卫,complete()用 CAS 保证只关闭一次——这在“正常完成”“用户取消”“异常失败”三条收尾路径并发竞争时至关重要。
前端解析器 useAgentStream.ts 逐行读取event:/data:行,按事件名分发到各自 handler,think增量实时进入思考面板,finish到达后流式输出正式收尾。
跨节点流式取消:标记、广播、复核三步走
单体应用里“点停止”很简单:同进程内存里打个标记即可。但 Ragent 是多节点部署,用户发请求的节点和实际跑流的节点往往不是同一台——停止请求落在 A 节点,流却在 B 节点上跑。这就是 StreamTaskManager.java 要解决的问题,它由三件东西组成:
- 本地任务注册表(Guava Cache,TTL 30 分钟):记录每个 taskId 的收尾回调与上游中断动作;
- Redis 取消标记(
ragent:stream:cancel:{taskId})+属主标记(ragent:stream:owner:{taskId}):让任何节点都能查到“这条流该不该停、谁有权停”; - Redis 广播主题(
ragent:stream:cancel):停止请求先写标记,再向所有节点发布taskId|发起方消息,持有该流的那台机器收到后执行本地取消。
完整流程如下:
- 流启动时,
register(taskId, userId, 收尾回调)把属主写入 Redis——注意注释里明确说“停止请求可能落在没跑这条流的节点上”,所以属主必须进 Redis 而不是只留本地; - 用户点停止,任意节点收到
/agent/v1/stop,执行 cancelByUser:先比对属主,不匹配就报“任务不存在或已结束”(不区分不存在与非属主,防止停止接口变成他人任务的探测器); - 校验通过后
publishCancel:先落 Redis 标记,再广播,所有节点(含本地)统一走监听器处理,避免重复调用; - 执行端
cancelLocal再次复核发起方(执行端复核兜住“属主还没落地”的抢跑窗口),CAS 保证取消动作只执行一次。
这里的细节很有工程味道:taskId 是雪花 ID,时间有序、可预测,所以“拿到 taskId 就能停别人的流”不是理论风险,必须靠属主比对挡住。
优雅打断:先中断框架,超时才断流
取消信号到达执行节点后,并不是粗暴地掐掉 HTTP 连接。AgentRunHandle.interruptUpstream 的策略是两级:
- 先礼貌中断:调用
agent.interrupt()通知框架走中断分支,等它把本轮 Agent 状态存盘(最多等 2 秒); - 超时再强制断流:礼貌中断失败或超时,直接
dispose()掐断 Reactor 链,并补一次状态存盘,保证工具执行结果不丢。
收尾统一走 finishCancelledStream:把已生成的内容以INTERRUPTED状态落库,补发cancel+done事件——已流出的内容不会丢,刷新页面后看到的和当场看到的一致。
另外还有一个容易被忽略的取消源:用户关页面、断网导致 SSE 断开。bindEmitterLifecycle 在onTimeout/onError/onCompletion三个钩子上都挂了回收动作,以系统身份发布取消——否则 ReAct 循环会空跑到迭代上限,白白烧 token。
小结:这套流式输出协议好在哪
- 分事件推送:思考、正文、工具、来源各走各的通道,前端渲染清晰,协议易于演进;
- 三出口 CAS 互斥:
complete/cancel/fail三条收尾路径只有一个能胜出,杜绝重复发帧; - 跨节点可靠取消:Redis 标记 + 广播 + 执行端复核,任何节点都能停任何节点的流,且属主校验防越权;
- 优雅降级:先等框架存盘再断流,中断的内容照样落库,刷新不丢历史。
想深入阅读,可以从这几个文件入手:
- 事件协议:AgentSSEEventType.java
- 事件桥接与收尾:AgentStreamEventBridge.java
- 取消管理器:StreamTaskManager.java
- 运行句柄与优雅打断:AgentRunHandle.java
- 前端 SSE 解析:useAgentStream.ts
【免费下载链接】ragent企业级 Agentic RAG 智能体 - 全链路覆盖文档解析、多路检索、意图识别、问题重写、会话记忆、MCP 工具调用与深度思考。面向真实业务场景,从 0 到 1 完整工程实现。项目地址: https://gitcode.com/gh_mirrors/ragent1/ragent
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考