最近在搞Java侧对接AI大模型的项目,SSE这个词几乎天天出现在代码和日志里。这个项目的核心,就是用Java后端消费大模型流式对话接口。跑通整个过程后,我完整走过了SSE从显式调用到隐式封装,再到虚拟线程性能优化的三个阶段,每一步都踩了不少坑,也把背后的原理摸了个透。这篇博文就是这段时间实战经验的完整记录,既适合正在做Java+AI集成的后端工程师,也适合准备Java面试时被问到SSE和虚拟线程的同学。我会从最基础的协议细节讲起,把三段演进的动机、代码、踩坑点全部展开,你看完可以直接照着抄。
1. 项目概述与场景拆解:AI接入为什么绕不开SSE
1.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 | 升级为长连接 | 全双工,服务端可推送 | 需自行实现 | 高 | 不错,但大材小用 |
| SSE | HTTP长连接 | 单向服务端推送 | 协议内置 | 低 | 完美匹配流式生成 |
很快锁定了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返回的是Stream<String>,异常处理很别扭,而且对流的生命周期控制不够透明,排查问题不如直接操作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(); HttpResponse<InputStream> 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, Map<String, 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 HttpResponse<InputStream> response; private volatile ExecutorService executor; public SseConnection(SseListener listener) { this.listener = listener; } public void start(String url, Map<String, String> headers, String body) { executor = Executors.newVirtualThreadPerTaskExecutor(); executor.submit(() -> consume(url, headers, body)); } private void consume(String url, Map<String, 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 timeout,20秒静默不会触发超时。真正需要关注的是“流空闲超时”:多长时间没有数据就认为连接死了。这个需要在读取循环里自己实现一个看门狗逻辑——每次读到数据就刷新一个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_8 |
| JSON被截断 | 多行data拼接错误 | 按SSE规范用换行符拼接data行 |
| 连接数持续上涨 | 关闭逻辑没放finally | 在finally统一关闭InputStream和线程池 |
| 取消后仍收到推送 | cancelled检查不完整 | 在循环内外、onDone前都检查标志位 |
这几条每一个都是真金白银换来的。尤其是idle timeout那个问题,如果只看报错信息,很容易误判成是服务端主动断连,然后去查大模型接口配置,绕一大圈才会想到是代理层。
最后再分享一个小技巧。如果你排查SSE问题时想看原始字节流,别一上来就上抓包工具。写一个只打印原始行的临时客户端,把每行都打出来,包括空行和冒号开头的注释行。很多“解析不出来”的谜团,其实是格式和你预期的不一样,眼见为实。等确认协议格式没问题了,再去排查网络层。这个习惯帮我节省了很多时间,也让我发现了一些大模型网关在SSE实现上的非标细节。做Java+AI集成,SSE就是那根管道,管道修扎实了,上面怎么盖楼都不怕。