📌系列说明:一个 Java 后端视角的 Spring AI 渐进式实战教程,载体为开源项目「劳小司 · 智能法律助手」。
- 序章:技术栈全景与 AI 学习指南
- 阶段一 · 流式对话内核:篇1 SSE 流式·停止生成·思考可见化(本文)/ 篇2 会话记忆压缩与流式滚动体验
- 阶段二 · 工具调用:篇1 Function Calling 与法律计算器 / 篇2 联网搜索与工具预算
- 阶段三 · RAG 知识库:篇1 起步与底账化 / 篇2 Agentic RAG 与引用可信度 / 篇3 检索质量与体验
- 阶段四 · 多模型路由:篇1 五路级联路由
- 阶段五 · 安全与质量门:篇1 安全层与强制检索 / 篇2 质量门与评估门禁 / 篇3 指代消解与阻塞隔离
- 阶段六 · 产品化与用户体系:篇1 认证·配额·门禁 / 篇2 前端·移动端·身份 / 篇3 劳动法专精与多模态
- 阶段七 · 存储演进与部署:篇1 存储迁移 / 篇2 部署契约
📦本篇涉及:
controller/ChatController.java、service/ChatService.java、前端stores/chat.js。
一、今天要做什么
打开任意大模型产品,回复是一个字一个字蹦出来的,而且生成途中可以随时点"停止"。本篇把这两件事做扎实,再加一件很多教程忽略的——把首 token 之前的黑盒等待变成可见的进度:
- SSE 流式输出:模型每生成一个 token 就实时推送,而不是等全部生成完再返回;
- 真正的停止生成:点停止后,大模型侧的生成也被中断,而不只是前端不显示;
- 思考过程可见化:路由 / 改写 / 检索 / 生成这些前置步骤实时透出,推理模型的
reasoning_content也流式展示。
二、为什么是 SSE,Servlet 栈又怎么流式
| 方案 | 适配度 |
|---|---|
| 轮询 | 延迟高、浪费请求,直接排除 |
| WebSocket | 全双工,对"服务端单向推流"过重,还要处理心跳重连 |
| SSE | 基于 HTTP 的单向推送,天然适配 LLM 逐 token 输出,断线自动重连 |
本项目主体跑在Servlet(Tomcat)栈上(绝大多数 Java 后端更熟悉,也方便复用 Spring Security 等 Servlet 生态),只在流式接口上引入 WebFlux 的Flux<ServerSentEvent>承载 SSE。不整体切 Netty——这正是一个"混合技术栈"的务实取舍:业务用熟悉的 Servlet,流式用响应式,各取所长。
三、一条流上的 8 类事件:哨兵标记设计
流式聊天不只是把 token 推出去,还要在同一条通道里传引用、思考进度、重试信号等。本项目的SSE 事件契约有 8 类:
| 事件 | 含义 |
|---|---|
token | AI 生成的正文片段 |
reasoning | 模型真实推理内容(reasoning_content 透传) |
stage | 首 token 前的阶段进度(路由/改写/检索/生成) |
citation | 结构化引用(法名/条号/相关度 JSON) |
retry | 质量门审校不通过,通知前端清空重接第二轮 |
quota | 配额拒绝(含剩余量与重置时间) |
done | 生成完成,data 为[DONE] |
error | 生成异常 |
难点:Spring AI 的ChatClient.stream()只产一条 token 流,怎么把 citation/stage/retry 这些异构信息也塞进去?
做法:在 token 流里混入哨兵标记(以\u0000空字符为前缀,如\u0000CITATION\u0000{json}),流末尾统一map转成对应的 SSE 事件类型:
.map(token->{if(token.startsWith(CITATION_MARKER)){returnServerSentEvent.<String>builder(token.substring(CITATION_MARKER.length())).event(EVENT_CITATION).build();}if(token.startsWith(STAGE_MARKER)){/* → stage 事件 */}if(token.startsWith(REASONING_MARKER)){/* → reasoning 事件 */}returnServerSentEvent.<String>builder(sseData(token)).event(EVENT_TOKEN).build();})用\u0000前缀是因为正常文本几乎不可能出现空字符,天然不会和模型输出撞车。这样一条 Flux 承载全部语义,前端按addEventListener('citation', ...)分别处理即可。
四、真正的停止生成
这是本篇最值钱的部分。"假停止"只是前端不再渲染,模型那边还在一个字一个字烧你的 token。真停止要让大模型侧的生成也中断。
关键三招:
① 停止信号走 Redis,跨实例可达。stop 请求和 SSE 流可能落在不同实例上,所以用 Redis 键chat:stop:{sessionId}作为信号,而不是进程内的标志位。
②takeUntilOther触发 Reactor cancel 传播。
.takeUntilOther(stopTrigger(sessionId))// stopTrigger 以 400ms 轮询 Redis 停止键takeUntilOther一旦其他源发出信号就终止主流,并向上游传播 cancel——cancel 一路传到发起 HTTP 流式请求的客户端,连接关闭,模型侧生成随之中断。这是 Reactor 响应式流的内置能力,比手动dispose优雅得多。
③ 订阅时清残留、结束时兜底清理。
.doOnSubscribe(s->redissonClient.getBucket(STOP_KEY_PREFIX+sessionId).delete()).doFinally(signal->redissonClient.getBucket(STOP_KEY_PREFIX+sessionId).delete());防止上一轮遗留的停止标志误伤本轮。
前端配合双中断:点停止时既es.close()关闭 EventSource,又POST /chat/stop通知后端——两端都断,体验才干净。
一个工程取舍:stopTrigger的实现是Flux.interval(400ms).filter(stopFlagSet).take(1)——400ms 轮询Redis 标志,而非 Pub/Sub 订阅。停止是低频操作,400ms 的感知延迟完全可接受,却换来实现极简、无需管理订阅连接的生命周期;同时 stop 标志设60s TTL,即便某条流异常退出没清理,信号也会自动过期、不会永久残留。这是一个"够用就好、不为了优雅引入额外复杂度"的典型权衡。
五、思考过程可见化
一次法律问答在首 token 到来前,后端其实做了不少事:读记忆、路由、指代改写、强制检索……全在阻塞等待。若不透出,用户面对的是转圈黑盒。
做法是用一个unicast sink 实时发射阶段事件,与主生成流拼接:
Sinks.Many<String>stageSink=Sinks.many().unicast().onBackpressureBuffer();Mono<ChatContext>prepared=Mono.fromCallable(()->prepare(sessionId,userMessage,stageSink,...)).subscribeOn(Schedulers.boundedElastic())// 阻塞预处理隔离到弹性线程.doFinally(s->stageSink.tryEmitComplete()).cache();// cache 防二次订阅重复执行 prepareprepared.subscribe();// 立即启动 prepare,stage 实时流入 sinkreturnstageSink.asFlux().concatWith(prepared.flatMapMany(this::pipeline))prepare()每走到一步就emitStage(stages, "检索法律知识库…"),前端立刻看到进度条滚动。推理模型的reasoning_content同理以REASONING_MARKER透传,前端渲染成可折叠的"思考面板",并算出"已思考 N 秒"(发送 → 首 token 的耗时)。
这里顺带用到了阻塞隔离(
subscribeOn(boundedElastic)+cache()),让 Tomcat 的服务器线程不被首 token 前的同步 IO 独占——这个主题在阶段五·篇3 会专门展开。
六、踩坑备忘
① GET 把 token 和问题拼进 URL,落进 nginx 访问日志。早期/chat/stream?userMessage=...,法律咨询原文和 JWT 全进了代理日志。改成POST + 请求体后,日志里只剩路径——这是 P0 级安全整改。
②filter不会取消上游,takeUntilOther会。想中断流式,用filter只是丢弃元素,模型那边还在跑;必须用takeUntilOther/take这类会向源传播 cancel 的操作符。
③prepare被订阅两次 = 检索做两遍。stageSink 和 pipeline 分别订阅prepared时,若不cache(),Mono.fromCallable每次订阅都重跑。加.cache()+ 提前subscribe()启动,既让 stage 实时流出,又保证 prepare 只执行一次。
④ SSE data 前导空格被协议吞。SSE 规范会剥掉data:后的一个前导空格,token 恰好以空格开头时就丢了。解决:对 token data 做 JSON 编码(sseData()),前端JSON.parse还原。
七、小结
| 概念 | 一句话 |
|---|---|
| SSE | 单向推送,逐 token,比 WebSocket 轻 |
| 哨兵标记 | 一条 Flux 混装 8 类事件,map 阶段分流 |
| 真停止 | Redis 信号 + takeUntilOther + cancel 传播到模型 |
| 思考可见化 | unicast sink 实时发 stage,reasoning 透传 |
八、下篇预告
流式跑通了,但一个多轮助手如果重启就失忆、聊久了上下文爆 token,体验照样崩。下一篇我们做会话记忆的持久化与两级压缩,并治理前端的流式滚动体验。
🌐源码与体验:Gitee(国内快)https://gitee.com/spaserby/laoxiaosi.git | GitHub https://github.com/spaserby/laoxiaosi.git
🖥 在线演示:https://laoxiaosi.noctisblue.com
本系列全套代码皆开源,觉得这篇有帮助,欢迎顺手点颗 ⭐