ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java后端SSE流式输出改造:从显式调用到虚拟线程实战

Java后端SSE流式输出改造:从显式调用到虚拟线程实战 先交代一下背景。我最近在一个 AI Agent 项目里负责 Java 后端的流式输出改造核心场景是把大模型的流式回答实时渲染到前端页面并且要支持用户随时中断。整个方案落地过程中我完整经历了从“显式调用 SSE”到“隐式封装 SSE”再到“引入虚拟线程优化性能”这三个阶段。这三个阶段不是简单的 API 使用进阶背后其实是 Java 并发模型、响应式编程思想和 AI 应用架构设计的深度耦合。今天这篇就把整个链路拆开讲透包括我踩过的坑、实测的数据以及最终沉淀下来的可复用方案。1. 从需求说起为什么 AI 对话必须用 SSE1.1 传统 HTTP 请求在 AI 场景下的致命短板如果只是做一个普通的接口请求进来、计算完成、返回结果传统 HTTP 完全够用。但大模型生成回答的耗时通常有几秒到几十秒用户不可能盯着空白页面干等。如果前端用轮询不断问服务器“好了没”不仅浪费大量请求资源体验也谈不上流畅。这里需要一种机制让服务器可以分多次把数据推给客户端客户端拿到一段就渲染一段用户看到的是“一个字一个字蹦出来”的实时效果。SSEServer-Sent Events正好就是干这个的。它是 HTML5 标准里定义的服务器推送技术基于 HTTP 长连接服务端可以持续发送事件流客户端用EventSource接口就能直接接收。和 WebSocket 相比SSE 的优势在于原生支持 HTTP、自动重连、文本场景下协议极简尤其适合大模型这种单向的 Token 流输出。我最早是用 Spring 的SseEmitter做的显式实现后来发现这套东西在 AI Agent 场景里会遇到一个关键问题——一个请求往往要串联多个模型调用和工具调用SSE 的发送逻辑会被散落在各处。1.2 显式调用的痛点SSE 发送逻辑满天飞这里说的“显式调用”就是最原始的方式在 Controller 里创建SseEmitter然后在业务代码的各个地方手动调用emitter.send()推送数据。我第一版代码大致长这样GetMapping(/chat) public SseEmitter chat(RequestParam String question) { SseEmitter emitter new SseEmitter(0L); executor.execute(() - { try { // 调用大模型第一轮 emitter.send(SseEmitter.event().name(message).data(第一轮结果)); // 调用工具 emitter.send(SseEmitter.event().name(tool).data(工具结果)); // 调用大模型第二轮 emitter.send(SseEmitter.event().name(message).data(第二轮结果)); } catch (Exception e) { emitter.completeWithError(e); } }); return emitter; }表面看逻辑清晰但一旦业务复杂起来就失控了。AI Agent 的链路通常是用户提问 → 模型判断是否调用工具 → 调用外部 API → 把工具结果回填给模型 → 模型生成最终回答。整个链路里至少有三四处需要发 SSE 事件思考过程的展示、工具调用状态、最终 Token 流、错误信息。如果全靠手动send业务代码里全是emitter的传参和判空而且一旦漏掉complete()前端连接就会一直挂着最终触发各种莫名其妙的超时。更麻烦的是错误处理。SSE 连接不像普通接口那样一次性返回出错时你得决定是发一个错误事件再结束还是直接completeWithError前端拿到的状态完全不一样。这个阶段我最大的体会是SSE 本身不复杂复杂的是业务链路和推送语义的耦合。1.3 为什么需要“隐式封装”把推送变成事件流既然手动发送太散自然想到封装。所谓“隐式封装”就是把 SSE 的发送动作从业务代码里抽走业务代码只需要“产生消息”由一个统一的组件负责“推送消息”。这其实就是事件驱动思想的落地业务层是一个消息生产者SSE 通道是一个消费者两者之间通过一个简单的抽象解耦。封装之后业务代码里不再出现SseEmitter这个类型取而代之的是一个自定义的消息发布接口。想发什么类型的消息、要不要结束都由统一组件内部处理。这样一来业务方法可以专注于自己的职责调用模型、调用工具、组装结果完全不需要关心“这段数据是怎么到前端的”。这也是我在重构过程中觉得最有价值的一步代码可测试性提升非常明显。2. 隐式封装落地从 SseEmitter 到消息发布抽象2.1 设计一个合理的消息模型在做封装之前我先把 SSE 输出内容梳理了一遍。AI 场景下前端需要区分不同类型的数据比如模型思考的中间过程、工具调用的名称和参数、最终的增量 Token、结束标志和错误信息。如果混在一起发前端就要用正则去猜体验很糟糕。所以我定义了一个SseMessage结构包含三个字段type表示消息类型content是内容traceId用于链路追踪。发送的时候统一转成 JSON 字符串事件名固定用message前端只监听这一个事件名再根据type字段做分支渲染。这种做法的好处是协议足够简单排查问题的时候抓包也清晰。public record SseMessage(String type, Object content, String traceId) { public String toJson() { // 使用 Jackson 序列化 } }消息类型我初期只定义了四种thinking、tool、token、done。后来实际使用中发现错误场景也必须显式区分于是加了error类型。每增加一种类型前端渲染逻辑和测试用例都要同步更新所以类型的定义要克制宁可少不要多。2.2 统一的推送器与订阅注册机制核心封装是一个SsePublisher组件。它的职责有四个建立连接、发消息、处理错误、关闭连接。在这个组件里我使用一个ConcurrentHashMapString, SseEmitter来维护多个会话key 是前端传进来的会话 ID。为什么要支持多会话因为用户可能开多个对话窗口每个窗口独立流式输出。这里有一个细节值得单独说SseEmitter创建之后必须在超时时间内完成首次发送否则连接会被自动关闭。我踩过一次用户问了问题模型却需要先跑一个很慢的工具导致前端连接空等超时后 SSE 断掉等模型结果出来想推送的时候已经晚了。解决方案是连接建立后立即发送一个连接确认事件同时把超时时间调大。这个“保活”思路在后来的虚拟线程方案里也延续了下来。2.3 用回调还是阻塞队列封装时最纠结的一个设计问题是业务层和推送层之间用什么方式通信。我试过两种方案。第一种是回调函数业务层调用publisher.send(callback)推送层在内部触发回调第二种是阻塞队列业务层把消息放进队列推送层从队列里取并发送。实测下来回调方案虽然写起来简单但在长链路场景里容易出现回调地狱而且异常传播路径不清晰。我最终选了阻塞队列方案理由有两个一是队列天然支持缓冲推送层可以批量处理不至于每条消息都触发一次 IO二是队列的消费逻辑可以放到独立线程里和业务线程解耦。这个设计在后来的虚拟线程场景里也发挥了作用因为虚拟线程虽然轻量但也不是无限量队列缓冲可以避免大量连接同时涌入时打爆推送线程池。2.4 隐式封装后的代码长什么样封装完成后业务层的代码非常干净。以 Agent 调用链为例核心方法只需要做三件事调用模型、根据结果决定是否调用工具、把最终结果发布出去。至于推送是走 SSE 还是 WebSocket业务层完全无感。Service public class AgentService { private final SsePublisher publisher; public void chat(String sessionId, String question) { publisher.send(sessionId, new SseMessage(thinking, 正在分析问题, traceId)); String toolResult callTool(question); publisher.send(sessionId, new SseMessage(tool, toolResult, traceId)); String finalAnswer callLlm(toolResult); publisher.send(sessionId, new SseMessage(token, finalAnswer, traceId)); publisher.send(sessionId, new SseMessage(done, 结束, traceId)); } }对比最初的显式版本这段代码几乎没有业务噪音。publisher.send内部会判断当前会话是否有效如果前端已经取消连接就静默丢弃不会抛异常打扰业务逻辑。这也是封装最有价值的地方异常处理被收敛到了边界处。3. 虚拟线程改造SSE 性能瓶颈的解法3.1 为什么原来的线程模型会成为瓶颈SSE 本质是长连接每条连接都会占用一个服务端线程来执行“连接保持 数据推送”的任务。在传统线程模型比如 Tomcat 的线程池下线程数量是有限的通常默认 200。如果同时在线用户超过这个数新的连接只能排队等待流式输出的实时性就无从谈起。我最初的做法是用ExecutorService为每个 SSE 请求分配一个线程这在测试环境几十个连接时完全没问题但一压测就暴露了线程池耗尽、CPU 被上下文切换打满、部分请求响应时间飙升到十几秒。问题不在于服务器性能不够而在于线程资源被不必要地占用了。想想看Agent 链路里大量时间花在等模型响应、等工具 API 返回这些等待属于 IO 阻塞线程阻塞在那里什么都不干纯粹是浪费。3.2 虚拟线程如何改变游戏规则Java 21 正式引入虚拟线程之后这个问题的解法就清晰了。虚拟线程由 JVM 调度不直接映射操作系统线程数量可以开到几十万甚至更多。阻塞操作发生时虚拟线程会自动让出底层载体线程等到数据就绪再恢复执行。这意味着SSE 连接被挂起等待推送时不再需要占住一个昂贵的操作系统线程。我做的改造非常简单。原来用executor.execute(() - { ... })现在直接Thread.ofVirtual().name(sse-push-, 0).start(() - handleSse(sessionId));就这一行连接处理从“线程池 200 上限”变成了“虚拟线程几乎无上限”。我用 1000 个并发 SSE 连接做了压测和原来的线程池方案对比虚拟线程方案的吞吐量大概提升了五倍左右CPU 占用反而更低。原因很直白虚拟线程的创建和切换开销远小于平台线程而且大量阻塞等待的成本被 JVM 吸收掉了。3.3 虚拟线程下的隐式封装依旧成立有人可能会担心虚拟线程这么轻量是不是可以放弃队列缓冲、每条消息直接发我实测下来队列缓冲仍然有存在价值。虚拟线程虽然轻量但也意味着可以开更多连接如果每条推送都直接触发一次网络 IO对 TCP 栈的压力反而更大。队列在这里的作用从“防线程爆炸”变成了“平滑 IO 毛刺”让发送速率更加稳定。另外一个需要注意的点是虚拟线程的局部变量、线程上下文和普通线程没什么区别但使用ThreadLocal时要格外小心。如果你在 SSE 链路里用ThreadLocal传递 traceId 之类的上下文虚拟线程配合池化场景容易出现上下文残留。我后来统一换成了显式参数传递彻底规避了这个问题。3.4 虚拟线程压测数据与调优参数压测环境我用了 8 核 16G 的云主机Java 版本是 21Spring Boot 3.2。测试场景是 1000 个并发 SSE 连接每个连接接收 100 条消息。三组数据对比非常直观方案完成耗时CPU 平均占用线程/连接数Tomcat 默认线程池32s78%200 上限大量等待固定线程池 50021s61%500仍有排队虚拟线程6s43%1000无排队虚拟线程方案在完成耗时上领先明显。刚开始虚拟线程推送时出现乱序问题排查后发现是虚拟线程并行度太高同一个会话的多条消息被多个虚拟线程并发发送。解决办法是按会话加锁或者用单线程的虚拟线程执行器来推送同一个会话保证同一会话内的消息严格按序到达。4. 中断与取消SSE 交互里最容易翻车的地方4.1 abort 场景的前后端配合SSE 的“实时渲染”只是表面功夫真正考验体验的是“用户点击停止按钮”。前端通过EventSource连接 SSE 后如果要中断需要调用abort()方法断开连接。问题在于前端断开后后端怎么感知SseEmitter在客户端断开时会触发完成回调但回调时机不一定即时而且如果你用的是自定义封装这个回调可能没被妥善处理。我在封装里做了这样一件事继承SseEmitter并重写onCompletion和onTimeout把会话状态标记为已失效。这样后续调用publisher.send时组件会先检查状态发现失效就直接返回不会做无意义的 IO。用虚拟线程后断开会话占用的虚拟线程也会很快结束资源释放更干净。4.2 服务端主动取消不只是让前端闭口还有一种中断是服务端主动发起的。比如用户请求触发了某个策略模型继续生成没有意义需要在服务端主动结束推送。这时候单纯complete()是不够的最好先发一个error或者done类型的事件给前端一个明确的终止信号再关闭连接。前端收到终止信号后会显示“已停止”而不是“连接断开”体验差别很大。4.3 超时与闲置断开的兜底SSE 长连接最怕的事情就是“半死不活”连接还在但数据不来了。很多网关或者负载均衡器默认有 60 秒空闲超时一旦没有数据传输就会断开连接。我在压测和线上都遇到过idle timeout问题表现为前端突然收到stream disconnected事件。解决办法是设计一个心跳机制如果 30 秒内没有实际业务数据推送器自动发送一个ping类型消息保持连接活跃。这个心跳消息前端要做忽略处理不要渲染出来。5. 常见问题与排查技巧实录5.1 问题速查表现象可能原因解决方案前端收不到任何消息连接未建立成功或首次保活消息发送失败建立连接后立即发送确认事件流式输出中途断开网关空闲超时增加心跳消息连接建立后线程耗尽使用平台线程池处理长连接切换虚拟线程或按连接分配线程多个会话消息串线会话映射 Key 冲突使用全局唯一会话 ID消息乱序同一会话并发发送同一会话使用独立虚拟线程/加锁业务代码抛异常后连接挂死未调用 complete统一在 finally 中结束连接5.2 调试 SSE 的独家技巧调试 SSE 最直接的工具是curl它能把原始事件流完整打出来。我在排查乱序问题时就是靠curl -N观察消息到达顺序定位到是并发发送导致的。这里有个经验别急着看前端效果先用 curl 确认服务端行为是否正确否则你很难分清是后端发错了还是前端渲染错了。另外日志里一定要记录每个会话的traceId和消息序号排查线上问题时能少走很多弯路。我的推送器里每条消息都会带上会话维度的序号前端也能拿序号做乱序检测前后端联动排查效率非常高。5.3 一个容易被忽视的坑消息吞吐与背压AI 场景下大模型返回是一串 Token 流如果后端转发给前端的频率过高而前端渲染能力跟不上就会出现中间数据积压最终表现为页面卡顿。我在压测时发现每毫秒推送几十条小消息和每 100 毫秒推送一批消息前端的渲染体验差别很大。最终方案是做了简单批处理达到时间窗口或消息条数阈值才推送一次这样网络包的数量大幅下降前端渲染也更平滑。这个“背压”思想在流式系统里很重要。虚拟线程解决了“连接数”的问题但没解决“消息速率”的问题。速率控制属于限流范畴和并发模型是两回事两者要搭配使用才能达到最佳效果。6. 实操复盘与延伸思考6.1 我给项目沉淀的三层结构复盘整个 SSE 改造我最后总结出三层结构每层职责清晰非常适合 Java AI 场景复用接入层Controller SseEmitter 或 WebSocket负责建立和回收连接不关心业务语义。发布层SsePublisher 组件负责消息路由、状态管理、心跳保活、超时处理是唯一需要依赖 SSE 具体 API 的地方。业务层AgentService只关注业务编排通过发布层发送语义化消息不依赖任何 SSE 概念。三层结构之后测试策略也变得简单业务层测试的时候把publisher替换成一个内存版测试替身断言消息类型和顺序是否符合预期即可完全不需要真的建立 HTTP 连接。6.2 非 AI 场景的通用性这套方案虽然是从 AI 场景出发但适用范围远不止大模型对话。比如在 Java 里实现“文件上传进度推送”“长任务执行状态轮询”“服务端主动通知前端刷新”等功能SSE 配合隐式封装的思路都是通用的。我后来在另一个数据同步项目里复用了这套结构只需要替换消息类型业务代码几乎没改。6.3 后续扩展方向如果你继续深挖可以考虑这么几个方向一是将 SSE 升级为 WebSocket 双工通信把用户取消、参数调整等控制信号也走同一条连接二是把消息发布层接上消息中间件比如 Kafka让 SSE 推送能力可以跨节点水平扩展三是在推送层增加限流和降级策略防止大流量时雪崩。我个人在实际操作中的体会是SSE 的技术门槛并不高真正决定项目体验的是边界划分和异常兜底这两件事。把“推送”这个动作从业务代码里抠出来再配合虚拟线程解决资源占用问题AI 应用的流式输出就能做到既优雅又抗压。
返回列表