我们项目最近接了不少大模型平台和 Agent 服务后端是 Java 17 Spring Boot 3.2。刚开始大家对接流式对话的方式很乱。有人直接在 Controller 里手动解析 HTTP 响应流有人用轮询硬磨还有人在 Service 里自己管理连接、重连和心跳。每次联调都有人喊“这里怎么断流了”“前端拿到的 JSON 怎么缺了一半”。后来我把这套逻辑从显式调用逐步收敛成了隐式封装再用虚拟线程把并发能力提了一截。这篇文章就把整个演进过程、核心代码和踩坑记录完整捋一遍适合在 Java 技术栈里做 AI 应用的人参考也适合准备面试时把 SSE、流式对接和虚拟线程串起来聊的同学。1. AI 场景为什么绕不开 SSE1.1 大模型接口和普通 HTTP 接口的根本差异传统 HTTP 接口的逻辑很简单客户端发一个请求服务端算完一次性把结果以 JSON 返回连接关闭。这个过程里客户端能做的事就是等。对于登录、查询订单这类操作几百毫秒的等待完全没有问题。但大模型的对话接口不一样。你问一个问题模型先生成思考过程再逐个 token 输出正文整个过程可能要十几秒甚至更久。如果后端等全部生成完再返回用户体验是灾难级的用户看着空白页面呆坐十几秒然后“唰”地一下冒出整段回答没有打字机效果出错时连个中间反馈都没有。我见过第一次对接的同事用最朴素的方式实现RestTemplate调大模型接口设置一分钟超时然后把整个字符串拼起来返回。这个方案问题很明显首 token 时间被拉长、用户体验差、中间如果流断了根本无法感知只能抛出超时异常。SSE 解决的就是这个问题。它的全称是 Server-Sent Events基于 HTTP 协议服务端可以持续地向客户端推送文本数据。大模型输出的每个 token 到了服务端就立刻通过 SSE 推送出去前端随收随渲染。用户看到的就是一个字一个字往外蹦的对话效果模糊感知明显比干等好很多。1.2 为什么不用 WebSocket 或者轮询聊 SSE 的时候一定会有人问WebSocket 不是更强大吗为什么不用它我实际对比过之后觉得AI 对话流式输出这个场景SSE 是“足够好且最省事”的方案。先看轮询。轮询最简单客户端每隔几百毫秒请求一次“有没有新内容”。但问题也最明显实时性做不到秒级请求频率高了浪费带宽频率低了用户感觉卡顿。大模型输出不均匀思考 5 秒、输出 3 秒、再停顿 1 秒轮询根本没法定准节奏。所以轮询只适合没有流式能力的过渡方案。再看 WebSocket。WebSocket 是全双工通信客户端和服务端可以互相推送功能确实强大。但代价是需要额外处理握手、心跳、断线重连、状态恢复线上代理和网关对 WebSocket 的配置要求高而且 AI 对话这个场景里客户端在等待回答期间基本不说话完全是单向接收。为了一个单向推送引入一套全双工复杂度性价比不高。SSE 的优势在于它还是 HTTP。不需要升级协议不需要处理复杂的握手状态代理网关穿透容易浏览器原生EventSource直接支持自动重连服务端实现也简单就是设置text/event-stream响应头然后持续输出文本块。从协议设计就能看出为什么 SSE 流行于 AI 场景它天生就是“服务端单向向客户端推送序列事件”的模型大模型 token 流恰好就是这个模型。1.3 Java 中 SSE 的基本实现Spring Boot 里实现 SSE 很简单最直接的工具是SseEmitter。先看最朴素的版本RestController public class ChatController { private final ChatService chatService; public ChatController(ChatService chatService) { this.chatService chatService; } GetMapping(value /api/chat/stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream(RequestParam String question) { // 超时时间先给 60 秒实际项目中通常需要更长 SseEmitter emitter new SseEmitter(60_000L); // 这里为了演示先开了一个线程去跑耗时任务 // 正式项目建议用线程池或虚拟线程别裸开单线程 new Thread(() - { try { chatService.streamChat(question, new ChatStreamListener() { Override public void onToken(String token) { emitter.send(SseEmitter.event() .name(message) .data(token, MediaType.TEXT_PLAIN)); } Override public void onDone() { emitter.send(SseEmitter.event() .name(done) .data([DONE])); emitter.complete(); } Override public void onError(Throwable throwable) { emitter.completeWithError(throwable); } }); } catch (Exception e) { emitter.completeWithError(e); } }).start(); return emitter; } }这段代码对应的 SSE 协议格式实际推给前端的网络流长这样event: message data: 你好 event: message data: 请 event: message data: 问有什么可以帮你 event: done data: [DONE]关键点在于协议约定每个事件块之间用空行分隔event:表示事件名称data:表示内容。事件名的作用是让客户端可以根据不同类型做不同处理AI 场景里常用message表示增量 tokendone表示结束error表示错误ping表示心跳。这样看整套链路不复杂。但随着业务复杂化手动实现很快会碰到瓶颈负责拉取大模型流的线程如果处理不完会堆积、重连没有统一处理、各家模型返回格式差异大、日志难以追踪。这也是下面封装要解决的事。2. 从显式调用到隐式封装2.1 裸写的时候我踩过的坑最早我们做的是一个客服知识库问答功能后端需要同时对接两家大模型。第一个版本都是从 Controller 开始直接在大模型 SDK 的 callback 里往SseEmitter塞数据。一开始很爽代码一行一行很直观。但两个迭代之后问题全冒出来了。第一个问题是解析逻辑散落各地。模型厂商返回的 payload 各不相同有的包装在choices[0].delta.content里有的直接返回result字段还有的会先返回 reasoning 内容再返回正式回答。每个接口都要重复写一段“从 JSON 里挖字段”的逻辑稍不留神就挖错层级。第二个问题是断线重连没法统一。大模型流式接口随时可能因为网络抖动、上游超时、代理空闲断开而中断。客户端不知道前端就一直转圈。我们最早的做法是前端发现断流就重新请求结果问题变成用户重复提问模型重复生成。后来才知道这些问题的根源在调用层没有把重连和断点续传做进公共逻辑里。第三个问题是线程和超时管理混乱。每个长连接如果都占一个业务线程连接多了线程池直接被占满。更麻烦的是很多人不知道SseEmitter默认超时时间很短一旦大模型思考超过时限连接就提前断了。我们在测试环境遇到过“前几个字正常后面没下文”的诡异问题查了很久才发现是超时设短了。这三个问题总结下来就一句话SSE 的核心逻辑应该是基础设施不能散落在业务代码里。于是我开始做隐式封装。2.2 一个流式客户端的最小封装封装的第一步是解决“怎么稳定地从远端 AI 接口拿流式数据”。我基于 Java 11 自带的HttpClient写了一个最小的SseClient它负责建立连接、按行解析 SSE 消息、把解析出的事件回调给上层。public class SseClient { private final HttpClient httpClient HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(10)) .build(); public void stream(String url, String apiKey, String payload, ConsumerAiStreamEvent eventConsumer, ConsumerThrowable errorConsumer) { HttpRequest request HttpRequest.newBuilder() .uri(URI.create(url)) .header(Authorization, Bearer apiKey) .header(Accept, text/event-stream) .header(Cache-Control, no-cache) .timeout(Duration.ofSeconds(120)) .POST(BodyPublishers.ofString(payload)) .build(); httpClient.sendAsync(request, BodyHandlers.asLines()) .thenAccept(response - parseLines(response.body(), eventConsumer)) .exceptionally(ex - { errorConsumer.accept(ex); return null; }); } private void parseLines(StreamString lines, ConsumerAiStreamEvent eventConsumer) { StringBuilder dataBuffer new StringBuilder(); String currentEvent message; lines.forEach(line - { if (line.startsWith(event:)) { currentEvent line.substring(6).trim(); } else if (line.startsWith(data:)) { dataBuffer.append(line.substring(5).trim()); } else if (line.isEmpty()) { if (dataBuffer.length() 0) { eventConsumer.accept(new AiStreamEvent(currentEvent, dataBuffer.toString())); dataBuffer.setLength(0); currentEvent message; } } }); } }这段代码有几个细节值得说BodyHandlers.asLines()会把响应流按行切开天然适配 SSE 按行分块的结构。每行要么是event:要么是data:遇到空行就表示一个事件结束。data:开头的行如果内容是 JSON可能包含空格和特殊字符不能简单substring(5).trim()了事要保留 JSON 的原始格式但trim()掉行尾回车和多余空格没问题。这个SseClient是懒人版本没有做重连和心跳。真正能上生产的版本还要在exceptionally里根据异常类型判断是重试还是直接报错并且用指数退避做重连间隔。我后面会展开。2.3 业务层怎么用封装有了SseClient业务层的调用得统一起来。我定义了一个AiChatService接口业务方只需要传入问题然后监听几个回调事件就行。public interface AiChatService { void chat(String userMessage, ChatStreamListener listener); }ChatStreamListener设计成带默认方法的接口业务方只需要关注自己关心的回调public interface ChatStreamListener { default void onToken(String token) { } default void onDone() { } default void onError(Throwable throwable) { } }这样业务代码就能写得很干净。比如在一个控制器里GetMapping(value /api/chat/stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter chat(RequestParam String question) { SseEmitter emitter new SseEmitter(0L); aiChatService.chat(question, new ChatStreamListener() { Override public void onToken(String token) { emitter.send(SseEmitter.event().name(message).data(token)); } Override public void onDone() { emitter.send(SseEmitter.event().name(done).data([DONE])); emitter.complete(); } Override public void onError(Throwable throwable) { emitter.completeWithError(throwable); } }); return emitter; }对比最开始的版本业务代码里不再有“请求上游”“解析协议”“处理重连”这些细节。调用方看到的就是一个chat方法加三个回调。这就是我从显式调用到隐式封装的第一层体会把跨系统的协议细节收进公共组件业务方写的是意图不是流程。2.4 封装时一定要想清楚的细节封装不是“抽个公共类”就完事。有四个细节在封装时一定要想清楚不然组件越封越难用。第一是重连与幂等。SSE 断线后自动重连是基础能力但重连带来一个副作用上游可能已经生成了一部分 token重连后从开头重新生成前端就会出现重复内容。解决思路有两个一是利用 SSE 协议自带的id:字段客户端把它作为游标传给服务端服务端从游标位置续推二是接受重连后从头开始但前端渲染时根据消息 ID 做去重。前者更优雅但对上游接口有要求不是所有大模型平台都支持。考虑成本我们一般先做后者。第二是协议归一化。不同大模型的响应格式不一样有的叫choices有的叫result有的还带工具调用字段。封装层需要定义一个统一事件模型比如TokenEvent、ToolCallEvent、DoneEvent、ErrorEvent然后内部做协议适配。业务方永远只和统一模型打交道换模型商只需要改适配器配置。第三是超时与心跳。SSE 连接什么时候算死不是 TCP 断开才叫死超过一段时间没数据就该判定异常。我在封装里同时设置了整体超时和空闲超时整体超时看模型最大生成时间一般 5 分钟到 10 分钟空闲超时更关键如果 30 秒没有任何data推送主动断开连接并触发重连逻辑避免客户端无限等待。第四是顺序与背压。虽然 SSE 本身按顺序推送但回调里的处理如果异步化可能出现后完成的事件先回调打乱 token 顺序。我的做法是回调默认按接收顺序直接调用不引入额外线程。如果业务方确实需要异步处理由业务方自己保证顺序组件不背这个锅。3. 虚拟线程让 SSE 长连接告别线程池噩梦3.1 长连接场景下传统线程模型的瓶颈SSE 服务端本质上是一个“连接挂起等待推送”的业务。一个用户发起对话后端要同时做两件事维持和用户之间的 SSE 长连接维持和上游大模型之间的 HTTP 长连接。这两条连接在等待期间都是阻塞的。在传统的 Java 线程模型里每个阻塞连接占一个线程。线程是操作系统资源创建和切换成本都不低。默认线程栈大小一般是 1MB 左右千级连接就意味着上千个线程、GB 级内存同时频繁上下文切换让 CPU 持续空转。我拿一个固定 200 线程的线程池做过一次简单压测。模拟 1000 个并发 SSE 连接每个连接保持 60 秒。结果还没到 500 个连接线程池就打满了新建请求全部排队。业务上表现就是用户多了之后点击对话按钮后半天出不来第一个字。传统线程模型的核心矛盾在于线程数量被操作系统限制但 SSE 长连接的数量需求远超线程容量。这也是为什么很多人遇到并发增长第一反应是拥抱 Reactive。3.2 虚拟线程是怎么解决的虚拟线程是 JDK 21 正式推出的轻量级线程核心思路是把“线程”从操作系统线程里解放出来。虚拟线程由 JVM 调度运行在少数平台线程之上。当虚拟线程执行到阻塞操作比如socket.read()、Thread.sleep()时JVM 会把底层平台线程腾出来去调度其他虚拟线程而不是让平台线程傻等。打个比方传统线程是“一个员工盯一个工位堵车了就整个工位停摆”虚拟线程是“一堆工单共享几个工位这个工单在等数据时员工先去处理下一个工单”。在 SSE 场景里这个特性简直是量身定做。每个 SSE 事件的上游生成时间可能间隔几秒甚至几十秒这段空闲时间里传统线程是在干等虚拟线程把等待时间让给了其他任务。所以同样的机器配置虚拟线程能支撑的连接数可以高一个数量级。还有一个容易被忽略的好处虚拟线程让代码风格回到了同步阻塞式。接 AI 流之前我一直纠结要不要把整个链路上 WebFlux后来看到虚拟线程果断放弃了这种复杂度。同步写法调试方便、堆栈可读、不违反直觉。这也是我认为虚拟线程对 SSE 场景最大的价值不用为了并发牺牲代码可读性。3.3 在 Spring Boot 中开启虚拟线程如果你用的是 Spring Boot 3.2 和 JDK 21开启虚拟线程几乎是一行配置的事spring: threads: virtual: enabled: true这个配置会让 Tomcat 的请求处理线程改用虚拟线程执行器。原本每个请求占用一个操作系统线程现在每个请求是一个虚拟线程Tomcat 底层的平台线程按 CPU 核数调度。如果你不想改全局配置也可以用编程式方式在 SpringApplication 构建时开启SpringBootApplication public class AiApplication { public static void main(String[] args) { SpringApplication app new SpringApplication(AiApplication.class); app.setVirtualThreads(true); app.run(args); } }有一个容易混淆的点是Spring Boot 3.2 的spring.threads.virtual.enabledtrue对 WebFlux 并不生效。因为 WebFlux 本身是响应式的不需要也不依赖 Servlet 容器线程。所以如果你项目里用的是 WebFlux虚拟线程的意义更多在于异步任务和阻塞式下游调用而不是 HTTP 出入口。3.4 虚拟线程 SSE 的实测感受开启虚拟线程后我立刻做了一轮简单验证同样的 4 核 8G 机器同一个 SSE 接口500 个并发连接保持 60 秒。传统线程池版本在连接数达到 200 左右时开始出现排队500 并发时大量请求超时。虚拟线程版本同样 500 并发连接全部正常建立响应时间几乎没有劣化CPU 占用也不高。继续加到 2000 并发时才开始看到一些网络层瓶颈但已经和线程池没有关系了。我不建议把这个数据当成绝对标准不同机器、不同网络环境差异很大。我想强调的是虚拟线程把“线程数”这个硬限制解除了之后SSE 长连接服务的扩容空间大得多瓶颈从服务端线程转向了文件描述符、带宽和上游模型服务的并发度。如果你在面试里聊到这个点可以这样串联SSE 是 IO 密集型场景虚拟线程对 IO 密集型任务收益最明显因为阻塞时平台线程可以被让出。3.5 虚拟线程的三个注意点虚拟线程不是银弹实际工程里有三个坑我确实碰到过。第一个是synchronized块内阻塞不会让出平台线程。虚拟线程一旦进入synchronized代码块如果中途发生了阻塞平台线程不会释放因为 Java 的内存监视器monitor不允许虚拟线程在持有锁的状态下被卸载。这意味着你如果在synchronized里调sleep或读网络流虚拟线程的并发收益直接打折扣。解决办法是这类场景改用ReentrantLock。第二个是ThreadLocal的代价变高。虚拟线程数量很大如果每个虚拟线程都往ThreadLocal里塞大对象内存开销会成倍放大。老框架里常见的“用 ThreadLocal 缓存数据库连接”的做法在虚拟线程场景下会变得很危险。建议在虚拟线程跑的代码里谨慎使用继承式ThreadLocal能用参数传递就用参数传递。第三个是 CPU 密集型任务没有收益只有 IO 密集才有效果。如果你把虚拟线程用在纯计算的接口上不仅没有提升反而因为虚拟线程调度器本身有开销性能可能更差。我见过有人把虚拟线程开在所有接口上然后说“还不如以前”本质是场景没选对。4. 常见问题与排查技巧实录4.1 一张速查表SSE 相关的线上问题五花八门但很多其实都踩在同一个坑的不同形态上。我把最常见的几个问题整理成一张表现象可能原因解决思路前端 EventSource 完全收不到数据Nginx 开启了缓冲数据被囤积在代理层关闭proxy_buffering连接几分钟后自动断开代理或客户端有 idle timeout服务端没有心跳包服务端定时发注释行或event: ping流只输出一部分就中断上游大模型长连接空闲超时加大客户端 read timeout超过模型思考时间数据全部挤在一起返回服务端没有 flush或响应头类型不对确认text/event-stream和Cache-Control: no-cache偶发重复内容断线重连后从头生成前端按消息 id 去重或后端做游标续传并发一大就大量超时线程池打满或连接数超过文件描述符上限使用虚拟线程检查ulimit -n报stream disconnected before completion空闲等待超时长时间没有数据心跳加长超时时间检查链路各层 idle timeout表格里的stream disconnected before completion是很多人在对接 AI 流式接口时遇到的报错。字面意思是“流在完成前断开了”但你发现服务端日志里根本没有异常。这时候 90% 的情况是长时间没有事件数据链路某一层认为连接空闲主动掐断了。解决办法很简单服务端定期发心跳数据包比如每 15 秒发一行注释: keep-alive或者一个event: ping空数据。4.2 一次断流的排查链路实录有一次我们线上出现了“用户等待回答时如果在某个环节思考超过 1 分钟前端就会断开”的问题。前端表现是停了 60 秒左右后提示连接失败然后重新发起请求。我当时的排查链路是这样的第一步看服务端日志。服务端没有异常模型生成正常SSE 推送到最后也完整说明生成链路没问题。第二步看接入层。我们前面挂了一层 Nginx检查 access log 发现该请求在 60 秒左右返回了 499 状态码。499 是 Nginx 自定义码表示客户端在服务端处理完之前主动断开了连接。这基本确认问题发生在链路下游。第三步看客户端。前端EventSource默认没有设置超时但浏览器层面的网络栈通常会有自己的超时行为。不过更可疑的是 Nginx 默认的proxy_read_timeout是 60 秒恰好对上症状。第四步改配置。在 Nginx 流式接口的 location 里加上proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s;改完再测60 秒断开的问题消失。这里最值得警惕的是 499 这个状态码很多人看日志只看 500、502忽略了 499 其实隐含了“请求被中途掐断”的信息。4.3 独家避坑技巧我在多个 SSE 项目里摸出几个小习惯不算什么高深理论但能省很多时间。第一个联调时不要用浏览器直接看流式接口用 curl 带-N参数。curl -N http://localhost:8080/api/chat/stream?questionhi会实时显示每个 token一眼就能看出服务端推送是否正常。如果 curl 输出正常而浏览器不正常问题多半在前端或代理层。第二个凡是经 Nginx 代理的 SSE必须检查代理层有没有关闭缓冲。proxy_buffering off是最常见的配置它保证数据到达 Nginx 后立即转发给客户端不会攒一块。漏掉这一步现象就是前端很久没动静然后突然全部内容一起出现。第三个服务端一定要有心跳间隔最好 10 到 30 秒。我习惯用注释行: keep-alive\n\n这种行不会触发客户端事件解析逻辑只是让连接上保持有数据传输避免中间代理把空闲连接回收。第四个SseEmitter超时时间要设置成和业务峰值时长匹配不是越长越好。我们曾经图省事设成Long.MAX_VALUE结果出现异常时连接迟迟不释放导致连接泄漏。后来改成 10 分钟并在finally里确保调用complete()。5. 工程化落地通知、多端和可观测性5.1 通知类 SSE 的实现套路SSE 除了 AI 对话流式输出最常见的场景是服务端主动向指定用户推送通知。两者的核心差异是AI 场景是“请求-流式响应”谁请求谁接收通知场景是“服务端主动推送”用户可能没发起请求但服务端知道该推给谁。实现套路其实不复杂核心是维护一个userId - SseEmitter的注册表。用户登录或进入页面后前端建立 SSE 连接服务端把 emitter 存进 ConcurrentHashMap。需要推送时从 Map 里取出对应用户的 emitter调用send()即可。用户断开时通过onCompletion回调把 emitter 从 Map 移除防止无效推送。Service public class NotificationService { private final MapString, SseEmitter clients new ConcurrentHashMap(); public SseEmitter connect(String userId) { SseEmitter emitter new SseEmitter(0L); clients.put(userId, emitter); emitter.onCompletion(() - clients.remove(userId)); emitter.onTimeout(() - clients.remove(userId)); return emitter; } public void push(String userId, String event, Object data) throws IOException { SseEmitter emitter clients.get(userId); if (emitter ! null) { emitter.send(SseEmitter.event().name(event).data(data)); } } }这里有个细节推送前最好判断emitter是否还在clients里防止用户已经断开但 Map 还没来得及清理导致发送时抛异常。异常处理建议单独捕获不能在推送逻辑里直接让业务失败。5.2 前端与跨语言对接SSE 是跨语言协议前端只要符合协议格式就能接收。浏览器端最简单的是EventSourceconst source new EventSource(/api/chat/stream?questionhello); source.addEventListener(message, (event) { console.log(收到的 token:, event.data); }); source.addEventListener(done, (event) { source.close(); });EventSource只支持 GET 请求如果想在 POST 请求体里携带问题参数可以用fetch配合ReadableStream手动解析const response await fetch(/api/chat/stream, { method: POST, headers: { Content-Type: application/json }, body: JSON.stringify({ question: hello }) }); const reader response.body.getReader(); const decoder new TextDecoder(); while (true) { const { done, value } await reader.read(); if (done) break; console.log(decoder.decode(value)); }Python 端读取 SSE 也很直接requests流式响应逐行解析即可。协议本身没有语言壁垒关键是理解事件边界和data:前缀解析规则。5.3 可观测性设计流式接口的排查难度比普通接口高因为过程是时间和数据双重维度普通接口只要看“成功/失败”就行流式接口还要看“首 token 多快”“token 是否中断”“重连了几次”。我建议至少在封装层埋这几个点位请求进入时生成 traceId贯穿整个链路的日志上游首个 token 到达的时间这是衡量模型服务性能的关键指标token 总耗时和 token 数算平均生成速度断线重连次数超过阈值要告警每次重连都单独打点方便判断是网络问题还是上游问题。这些数据刚开始可以用日志打点线上量大了之后应该接进 Prometheus 或者变色龙这类监控体系按接口、按供应商维度建立报表。我自己的经验是流式接口的监控到位比加压测试更重要因为生产环境和压测环境的网络差异往往比代码性能差异更影响用户体验。结尾从最开始的显式解析到后来封装成公共组件再到用虚拟线程替换传统线程模型整个过程对我来说最大的体会是不是每个场景都需要一开始就上重框架。我先用最朴素的方式跑通链路再根据痛点逐步收敛每一步都有明确收益。虚拟线程出现以后SSE 服务端代码终于可以回归到同步阻塞的直觉写法不需要为了并发硬把代码改成响应式。最后再分享一个小技巧如果你也在对接 AI 流式接口服务端最好每 15 秒发一个注释行做心跳客户端 read timeout 设置得比模型最大思考时间更长这两个动作加起来能省掉一整类“莫名其妙断流”的问题。