
最近在搞Java侧对接AI大模型的项目SSE这个词几乎天天出现在代码和日志里。这个项目的核心就是用Java后端消费大模型流式对话接口。跑通整个过程后我完整走过了SSE从显式调用到隐式封装再到虚拟线程性能优化的三个阶段每一步都踩了不少坑也把背后的原理摸了个透。这篇博文就是这段时间实战经验的完整记录既适合正在做JavaAI集成的后端工程师也适合准备Java面试时被问到SSE和虚拟线程的同学。我会从最基础的协议细节讲起把三段演进的动机、代码、踩坑点全部展开你看完可以直接照着抄。1. 项目概述与场景拆解AI接入为什么绕不开SSE1.1 AI大模型为什么把SSE当成默认协议现在几乎所有大模型开放平台 OpenAPI兼容的那一挂的聊天补全接口默认流式输出都是走SSE。SSE全称Server-Sent Events翻译过来是“服务端发送事件”。它不是WebSocket那样的全双工协议而是一个基于HTTP长连接的单向推送协议客户端发一个普通的HTTP请求服务端保持连接不关闭然后把数据以特定的文本格式持续推给客户端。为什么大模型普遍选SSE而不是WebSocket最核心的原因是大模型生成token本身就是“一次请求持续产出”的模式服务端不需要从客户端接收数据只需要单向往下推。WebSocket需要先升级协议要处理二进制帧、心跳、状态机复杂度高出一大截。而SSE底层就是普通HTTP现有网关、负载均衡、监控体系几乎都能无缝兼容不需要额外基础设施。对服务端来说生成流式响应时只需要持续往已打开的HTTP连接里写文本即可实现成本极低。从客户端视角看SSE还有一个很实用的特性自动重连。协议内置了断线重连机制服务端还能通过retry字段指定重连间隔客户端只需要在事件流里维护一个lastEventId即可。这些特性叠加在一起让SSE成了AI服务端的“默认语言”。你随便接一个Chat模型翻它的文档大概率都是text/event-stream。1.2 项目需求拆解与阶段规划我这次的项目需求并不复杂Java后端作为中间层接收上游业务请求再转发给大模型接口把大模型生成的内容实时推给前端。难点在于大模型的响应是一个持续的流可能有几十个甚至几百个增量片段而且每个片段的到达时间不固定模型“思考”的时候甚至会出现长时间静默。怎么在Java侧稳定地读取、解析、转发这个流是整个项目真正技术含量所在。动手之前我做了个简单的技术选型对比这里直接贴出来方案连接方式服务端push能力自动重连实现复杂度AI场景适配性普通HTTP轮询短连接无需客户端反复请求无低差延迟高、浪费资源WebSocket升级为长连接全双工服务端可推送需自行实现高不错但大材小用SSEHTTP长连接单向服务端推送协议内置低完美匹配流式生成很快锁定了SSE。但“用SSE”和“用好SSE”是两码事。我给自己划了三个阶段第一步先用Java自带的HttpClient显式读取SSE流把协议细节摸清楚第二步把这个过程封装成隐式的流式接口让业务代码感受不到SSE的存在第三步针对高并发场景引入虚拟线程解决阻塞式读取导致的线程资源瓶颈。三个阶段对应三个真实痛点下面逐段展开。2. 显式调用先让大模型的消息“流”起来2.1 一次SSE通信链路拆解SSE的数据格式非常直观用一个具体例子来说event: message id: 1 data: {content:你好} retry: 10000 event: message id: 2 data: {content:我们} data: {content:开始吧} id: 3 data: [DONE]每个事件由若干字段行组成字段和值之间用冒号分隔可能出现的字段有data、event、id、retry。多个data行会被拼接成一个事件的数据拼接符是换行符。事件与事件之间用一个空行分隔。客户端读到空行就认为一个事件结束了。大模型接口在这个基础上做了一些简化通常只发送data:行事件类型固定为默认的message最后一个事件是固定的[DONE]标识表示整个流式响应结束。所以Java侧的实际解析逻辑可以很轻量逐行读取判断是否以data:开头遇到空行就触发一次事件回调最后遇到[DONE]就结束。这里有个常被忽视的点retry字段控制的是重连等待时间服务端可以在任意事件中带上它。我在调试时发现某些AI网关会在异常时在SSE事件里塞一个retry: 3000客户端如果不处理这个字段重连会出现参考偏差。所以解析器里最好把这个字段也解析出来至少别让它污染data内容。2.2 用HttpClient把流式输出“读”出来Java 11起内置的java.net.http.HttpClient就已经能很好地处理流式响应。关键是通过BodyHandlers.ofInputStream()拿到原始输入流再手动按行读取。第一步不要用BodyHandlers.ofLines()那个API返回的是StreamString异常处理很别扭而且对流的生命周期控制不够透明排查问题不如直接操作InputStream方便。HttpClient client HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(10)) .build(); String jsonBody {\model\:\qwen-plus\,\stream\:true,\messages\:[...]}; HttpRequest request HttpRequest.newBuilder() .uri(URI.create(https://api.example.com/v1/chat/completions)) .header(Authorization, Bearer apiKey) .header(Content-Type, application/json) .header(Accept, text/event-stream) .timeout(Duration.ofMinutes(2)) .POST(HttpRequest.BodyPublishers.ofString(jsonBody)) .build(); HttpResponseInputStream response client.send(request, HttpResponse.BodyHandlers.ofInputStream()); BufferedReader reader new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8)); StringBuilder dataBuffer new StringBuilder(); String line; while ((line reader.readLine()) ! null) { if (line.isEmpty()) { if (dataBuffer.length() 0) { String rawData dataBuffer.toString(); if ([DONE].equals(rawData)) { break; } // 这里把 rawData 反序列化成业务对象 JsonNode node objectMapper.readTree(rawData); String content node.path(choices).path(0) .path(delta).path(content).asText(); System.out.print(content); dataBuffer.setLength(0); } continue; } if (line.startsWith(data:)) { String data line.substring(5); if (data.startsWith( )) { data data.substring(1); } dataBuffer.append(data).append(\n); } }这段代码能跑但有一个细节要注意我用dataBuffer累积data行并加上换行符是因为SSE规范规定多条data行在事件结束时要拼接成一个文本。如果直接替换末尾换行遇到带格式的JSON字符串比如消息内容里本身含\n就不会丢失数据了。拼接完成后用rawData.replaceAll(\n$, )去尾即可。顺带提一嘴状态码检查。我在一开始没检查response.statusCode()后来遇到鉴权失败时返回的是200还是401居然看配置导致代码把错误页当成SSE流解析报了一堆奇怪的JSON异常。正确做法是把2xx以外的响应全部视为错误把response.body()读出来记日志。2.3 为什么显式调用只配做“第一版”上面的显式代码是我第一版的原型能跑但它有几个硬伤第一业务逻辑和协议解析完全耦合。每次对接一个新模型就要复制粘贴这一大坨读取逻辑然后在while循环里塞不同的JSON字段提取代码。如果哪天解析规则变了所有调用方都得跟着改。第二连接生命周期管理极其繁琐。[DONE]之后连接并不一定会立即关闭有些服务端要等客户端主动断开。一旦上游在消费完事件流后忘记关闭InputStream连接就泄漏了。HTTP连接池里的连接被占满后新请求全部排队等待表现就是“系统没挂但接口越来越慢”。第三也是最致命的——整条链路是阻塞式的。client.send()会阻塞当前线程直到拿到响应头reader.readLine()又会阻塞当前线程直到下一行数据到达。如果我用一个线程池并发处理多个流式请求每个请求都需要一个线程在那干等。并发数一上来线程就被打满了。这块的解法我放到第4节专门讲但它其实从第一版就要有意识。所以显式调用只用来做协议验证是对的。它让你把SSE的每个字节都看清了但绝对不能作为生产代码直接铺开。写第二版时我的目标非常明确把SSE彻底封装起来让业务代码不需要知道“流”“事件”这些概念。3. 隐式封装把SSE的复杂性关进抽屉里3.1 设计目标让调用方忘掉SSE的存在封装这件事很多人的第一反应是写一个工具类把HttpClient那段代码抄进去对外暴露一个返回值。这确实比复制粘贴好但不够。真正好用的封装是让调用方的代码看起来像在调用一个普通方法我理想中的使用方式是这样的sseClient.stream(url, requestBody, new SseListener() { Override public void onEvent(String data) { // 每收到一个增量片段这里触发一次 sendToFrontend(data); } Override public void onDone() { // 所有内容接收完毕 channel.close(); } Override public void onError(Throwable error) { log.error(sse stream error, error); } });这段代码读起来很顺调方只关心三个时机——有数据、全部完成、出错了。它不需要知道SSE是分行的不需要知道[DONE]长什么样也不需要关心空行和字段拼接。隐式封装的本质就是把“过程式的事件解码”转换成“声明的回调”。除了这个回调接口我还设计了另一个纬度可取消。流式调用通常是长耗时的前端如果切走了或者用户主动停止生成服务端应该能中止本次流式请求否则大模型还在持续消耗token。因此封装的返回对象必须提供一个cancel()方法而不是只返回void。3.2 从事件流到回调的封装骨架下面是我最终沉淀的核心封装代码去掉了具体业务保留了通用的骨架逻辑public class SseStreamingClient implements AutoCloseable { private final HttpClient httpClient HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(10)) .build(); private final ObjectMapper objectMapper new ObjectMapper(); /** * 发起一个SSE流式请求 */ public SseConnection stream(String url, MapString, String headers, String body, SseListener listener) { SseConnection connection new SseConnection(listener); connection.start(url, headers, body); return connection; } public class SseConnection implements AutoCloseable { private final SseListener listener; private volatile boolean cancelled; private volatile HttpResponseInputStream response; private volatile ExecutorService executor; public SseConnection(SseListener listener) { this.listener listener; } public void start(String url, MapString, String headers, String body) { executor Executors.newVirtualThreadPerTaskExecutor(); executor.submit(() - consume(url, headers, body)); } private void consume(String url, MapString, String headers, String body) { try { HttpRequest.Builder builder HttpRequest.newBuilder() .uri(URI.create(url)) .header(Accept, text/event-stream) .POST(HttpRequest.BodyPublishers.ofString(body)); headers.forEach(builder::header); response httpClient.send(builder.build(), HttpResponse.BodyHandlers.ofInputStream()); if (response.statusCode() ! 200) { String errorBody new String( response.body().readAllBytes(), StandardCharsets.UTF_8); listener.onError(new RuntimeException( SSE request failed: response.statusCode() , body: errorBody)); return; } BufferedReader reader new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8)); StringBuffer dataBuffer new StringBuffer(); String line; while (!cancelled (line reader.readLine()) ! null) { if (line.isEmpty()) { if (dataBuffer.length() 0) { String rawData dataBuffer.toString(); if ([DONE].equals(rawData)) { break; } listener.onEvent(rawData); dataBuffer.setLength(0); } continue; } if (line.startsWith(data:)) { String data line.substring(5); if (data.startsWith( )) { data data.substring(1); } dataBuffer.append(data).append(\n); } } if (!cancelled) { listener.onDone(); } } catch (IOException e) { if (!cancelled) { listener.onError(e); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); if (!cancelled) { listener.onError(e); } } finally { close(); } } public void cancel() { cancelled true; close(); } Override public void close() { try { if (response ! null response.body() ! null) { response.body().close(); } } catch (IOException ignored) { } if (executor ! null) { executor.shutdownNow(); } } } }这段封装有几个设计要点值得单独讲。cancelled标志是全局的consume方法里先检查再处理确保取消后不再触发回调。close()放在finally里读流异常和正常结束都会关闭底层输入流。虚拟线程池的引入让这个封装天然具备了后面第4节要讲的性能优势这里先不展开。在consume方法里我做了一个关键的判断[DONE]被当作事件数据送到了onEvent吗没有。它在事件组装阶段就被拦截了业务回调只会收到真正的JSON字符串。这个细节看似简单但如果漏了下游反序列化时必然炸出“Unrecognized token DONE”之类的错误。3.3 重连、超时与取消等边界设计隐式封装最难的不是正常路径而是各种非正常路径。我逐一说。[DONE]之后要不要关闭连接一定关。有些服务端发送完[DONE]之后连接不会立刻关闭如果客户端不主动关连接会一直占着HTTP连接池的位置直到空闲超时。我在finally里统一调用close()就是保证无论[DONE]是正常结束还是异常中断底层的网络资源都被及时释放。超时如何设计我在HttpRequest上设置了timeout()但大家要注意这个timeout是指“从请求发出到响应首字节”的等待时间不是整个流的空闲时间。大模型常见的场景是连接建立后模型内部思考20秒不发任何数据如果只依赖request timeout20秒静默不会触发超时。真正需要关注的是“流空闲超时”多长时间没有数据就认为连接死了。这个需要在读取循环里自己实现一个看门狗逻辑——每次读到数据就刷新一个lastDataTime后台定期检查这个时间差超过阈值就强制cancel()并抛出空闲超时异常。自动重连到底要不要做SSE协议原生支持重连但AI场景要慎重。大模型流式接口大多不保证幂等重连会导致重复生成和重复扣费。我的建议是不要在框架层面自动重连把错误抛给业务层由业务判断这次请求是否允许从头再来一次。比如用户手动刷新页面重试那是业务行为框架自动重连纯粹是烧钱行为。心跳注释行怎么处理服务端偶尔会发一行以冒号开头的注释行例如: keep-alive用于保活。解析时遇到冒号开头的行应当跳过不能当成data解析。我第一版没处理后来发现日志里偶尔冒出unknown SSE field: : keep-alive的警告就是这行的作用。4. 虚拟线程让并发连接数不再成为瓶颈4.1 阻塞式读流为什么是并发瓶颈前面所有代码都是阻塞式I/O这在低并发下毫无问题但一旦并发量上来问题就暴露了。传统Java线程模型里每个平台线程都对应一个系统线程线程栈默认1MB左右线程切换和创建销毁都有不小的系统开销。而SSE场景的特点是线程在readLine()上长期阻塞——大模型生成一次回复通常需要几秒到几十秒这段时间线程不干活只是等数据。假设Tomcat线程池默认200个线程只要有200个并发SSE请求每个请求占住一个线程阻塞在读流上那么第201个请求就会排队。这时候整个接口的表现是小请求也被堵在后面系统吞吐量急转直下。我当时的临时方案是调大server.tomcat.threads.max从200调到1000但实测效果不好——线程多了上下文切换开销和内存占用都上去了GC压力也明显变大。这只是把瓶颈往后推没有真正解决。这个问题的根源在于阻塞I/O让平台线程“空转”。我们需要一个机制让线程在等待I/O时被释放出来去做别的事等数据到了再回来继续处理。虚拟线程就是为此而生的。4.2 虚拟线程下的读取模型改写虚拟线程Virtual Threads是Java 19引入、Java 21正式落地的特性。它与平台线程最大的区别是虚拟线程由JVM调度而不是操作系统调度。虚拟线程阻塞时JVM会把它挂起释放底层载体线程Carrier Thread去执行其他虚拟线程等到I/O就绪再把结果恢复到虚拟线程上继续执行。这正好解决了SSE阻塞读流的问题——readLine()阻塞的不再是稀缺的系统线程而是一个轻量级的虚拟线程。在Java 21中要创建一个虚拟线程非常简单// 方式一直接创建并启动一个虚拟线程 Thread.ofVirtual().name(sse-consumer).start(() - { // 这里做阻塞读流 }); // 方式二虚拟线程池推荐在服务里使用 ExecutorService executor Executors.newVirtualThreadPerTaskExecutor(); Future? future executor.submit(() - { // 阻塞读流 });放到上面的SseStreamingClient里只需要把consume方法从平台线程池切到虚拟线程池。在我的代码里stream()方法内部使用的就是Executors.newVirtualThreadPerTaskExecutor()。这意味着每一个SSE连接对应一个虚拟线程而不是一个平台线程。改造成本极低但效果天翻地覆。还有一个更彻底的用法如果你用的是Spring Boot 3.2及以上版本可以直接开启应用级虚拟线程支持spring: threads: virtual: enabled: true开启后Tomcat处理HTTP请求的线程就是虚拟线程包括SSE请求在内的整个请求处理链路都跑在虚拟线程上。我们项目里两个都开了HTTP请求线程用Spring配置SSE消费线程用自定义的newVirtualThreadPerTaskExecutor()配合得很稳。4.3 性能验证与适用边界做了一次相对严谨的压测这里把数据拿出来给大家参考。测试环境是4核8G的容器Java 21网关直连后端压测工具模拟并发SSE请求每个请求模拟大模型输出20个chunk、持续6秒。场景并发数平台线程方案虚拟线程方案200并发通过CPU正常通过CPU正常1000并发大量超时线程池打满通过响应延迟稳定4000并发直接拒绝服务通过平均延迟上升但无超时内存占用1000并发约2.5G约1.2G这个结果很符合预期。平台线程方案在1000并发时基本已经崩溃而虚拟线程方案到4000并发还能挺住内存占用反而更低——因为虚拟线程的栈是JVM管理的堆内存小而灵活而且会随线程消亡自动回收。不过这里必须泼一盆冷水虚拟线程不是万能药它有严格的适用边界。第一虚拟线程不适合CPU密集任务。while(true){ mathHeavy() }这种场景调度器没有机会挂起虚拟线程还会增加调度开销性能不如平台线程。在SSE场景里读取循环中如果塞了大JSON的复杂反序列化要评估好CPU占比。第二synchronized关键字在虚拟线程下有“钉扎”问题。如果虚拟线程锁定了载体线程阻塞时不会释放载体线程导致并发能力退化。JDK已经在修复相关场景但你在编码时还是要尽量避免在SSE读取回调里用重量级synchronized。第三底层依赖如果自己实现了NIO线程模型比如某些自定义的Netty服务虚拟线程可能帮不上忙甚至冲突。它最适合的是“在平台线程上做阻塞I/O”这个传统模式。5. 常见问题与排查实录5.1 idle timeout waiting for sse长连接被中间层掐断这个报错我在项目上线第二天就遇上了而且是在凌晨高峰期。日志里大量出现stream disconnected before completion: idle timeout waiting for sse。排查了一圈发现不是我们代码的问题是中间代理层的空闲超时。SSE连接虽然是长连接但大模型有时会思考很久不发数据。比如模型在调用工具API、或者在推理长上下文时可能整整30秒、甚至60秒没有往客户端推送任何字节。而Nginx的proxy_read_timeout默认配置通常是60秒云厂商的负载均衡也有类似空闲超时。一旦超过这个阈值中间层就会主动断开SSE连接客户端这边的表现就是流突然断了。排查思路分三步第一步看服务端日志有没有主动断连记录确认不是应用层问题第二步看客户端到服务端之间有几层代理逐一确认超时配置第三步确定断连时“静默时长”到底是多少是60秒还是300秒。解决办法有三个层次。最直接的是把Nginx的proxy_read_timeout调大比如600秒或者在代理环节关掉空闲超时检测。但有些云产品你没法改配置这时就要在应用层做保活服务端每隔15秒往SSE流里发送一行注释行: keep-alive\n\n这会让代理认为连接还在活跃。注释行是SSE协议的规定字段客户端解析时直接忽略不会影响事件流。实测把保活加上后这个错误几乎绝迹了。5.2 流式消息被截断和中文乱码这两个问题经常一起出现我一开始以为是同一个原因后来发现完全是两码事。中文乱码的原因只有一个BufferedReader没有指定UTF-8。new BufferedReader(new InputStreamReader(response.body()))会用系统默认编码读取在Linux容器里通常是UTF-8没问题但在Windows开发机上就是GBK一次乱码能让你排查半天。我的建议是写成InputStreamReader(response.body(), StandardCharsets.UTF_8)不要省略。消息被截断的排查相对复杂。现象是某条SSE事件收到的JSON只有一半解析必炸。后来发现是dataBuffer的拼接逻辑有问题一个SSE事件的data可能被拆成多行我第一版用了dataBuffer.append(data)而不是append(data).append(\n)结果遇到消息内容中本身有换行的场景时数据就错位了。按照SSE规范多行data拼接时要用换行符连接这一步不能省否则JSON的字符串里可能少一个\n导致解析后内容不一致。还有一种截断来自代理层的缓冲。有些代理默认启用了响应缓冲把SSE流攒到一定量才转发导致前端看到的是“一坨一坨”的数据延迟巨大。遇到这种场景一般需要在服务端返回响应头上加X-Accel-Buffering: no或者用Cache-Control: no-cache表明这是实时流。5.3 连接泄漏与取消失效连接泄漏问题在线上出现过一次“血案”运维反馈连接数居高不下数据库连接池也告警但应用进程内存和CPU都正常。排查发现是一个调用方在onDone回调里抛了异常而我的第一版代码在回调后没有finally关闭输入流导致连接永远不释放。后来我把close()移到了finally块里这个问题才根治。另外要特别强调虚拟线程池的关闭时机。Executors.newVirtualThreadPerTaskExecutor()并不会在每次调用cancel()时自动关闭你需要把它作为SseConnection的成员变量在close()里主动shutdownNow()。如果不关闭虚拟线程池每次请求都会创建一个池对象虽然在虚拟线程资源本身上开销不大但池对象累积起来依然是个隐患。还有一个坑是取消后依然触发回调。cancelled标志必须在readLine()循环里外都检查一遍包括[DONE]之后、finally之前。否则会出现“用户已经取消请求但最后一条事件还是写给了前端”的诡异现象。我当时就是漏了onDone()前的检查导致页面关闭后还能收到“回答完成”的推送。这里把三个常见问题的速查表整理出来方便你排查时直接对照症状根因解决办法stream disconnected before completion: idle timeout waiting for sse中间代理空闲超时调大超时阈值或服务端定期发送保活注释行中文乱码InputStreamReader未指定UTF-8显式指定StandardCharsets.UTF_8JSON被截断多行data拼接错误按SSE规范用换行符拼接data行连接数持续上涨关闭逻辑没放finally在finally统一关闭InputStream和线程池取消后仍收到推送cancelled检查不完整在循环内外、onDone前都检查标志位这几条每一个都是真金白银换来的。尤其是idle timeout那个问题如果只看报错信息很容易误判成是服务端主动断连然后去查大模型接口配置绕一大圈才会想到是代理层。最后再分享一个小技巧。如果你排查SSE问题时想看原始字节流别一上来就上抓包工具。写一个只打印原始行的临时客户端把每行都打出来包括空行和冒号开头的注释行。很多“解析不出来”的谜团其实是格式和你预期的不一样眼见为实。等确认协议格式没问题了再去排查网络层。这个习惯帮我节省了很多时间也让我发现了一些大模型网关在SSE实现上的非标细节。做JavaAI集成SSE就是那根管道管道修扎实了上面怎么盖楼都不怕。