做Java后端接大模型接口,我最头疼的一直不是模型选型,而是流式输出这层网络代码。你在浏览器里看到的“打字机”效果,是AI服务方用SSE(Server-Sent Events)把token一段一段推过来的,但这一层到了Java手里,处理方式就五花八门了:有人直接用OkHttp硬啃字节流,有人封装了通用解析器,还有人在JDK 21之后用虚拟线程把并发直接拉满。我三个阶段的方案都写过,从最原始的显式调用到隐式封装,再到最近用虚拟线程重构,性能差了一个数量级不止。这篇文章就把完整链路拆开讲清楚:SSE协议本身怎么工作、Java端怎么解析、封装层要注意哪些坑、为什么虚拟线程是当前Java+AI流式场景的最优解,以及我踩过的几个实实在在的坑。
1. 为什么AI流式输出绕不开SSE:一个2010年的“老协议”被大模型救了
先说个可能反直觉的事实:SSE不是新东西,它比WebSocket还老,是HTML5规范里自带的能力。早期用它的场景主要是股票行情、服务器日志推送这类单向数据流,这几年AI大模型全面走Token流式输出,SSE才被重新翻出来当主角。原因很简单:AI对话的响应天生是“一次性、单向、持续长时间输出”的,SSE就是为这种场景设计的。
1.1 SSE的协议格式:text/event-stream到底长什么样
SSE走的就是普通HTTP协议,只不过服务端返回的Content-Type是text/event-stream,连接建立之后,服务端可以持续往同一个响应体里写数据。HTTP响应体的流式特性(chunked transfer encoding)天然支持这种模式,所以SSE实现起来非常轻。
一条完整的SSE消息格式是这样的:
id: 1 event: message data: 你好 data: 世界字段之间用换行分隔,每个字段是“字段名: 空格 + 值”的格式。如果一条消息有多个data:行,它们在客户端会被拼接成一个整体,中间补一个换行符。两条消息之间用空行(也就是两次换行\r\n\r\n)隔开。id字段用来支持断线重连时的Last-Event-ID续传,event字段表示事件类型,AI场景下通常用不到那么复杂,最常见的只有data:这一个字段。OpenAI官方协议里,流结束时会发送一条data: [DONE],Java端解析器看到这一行就知道整个流走完了。
1.2 为什么替代方案都在折腾:轮询太蠢,WebSocket又太重
如果不用SSE,传统做法是轮询:前端每隔几秒请求一次后端,后端去问模型状态。这在AI场景下体验极差——首字延迟本来可以到几百毫秒,轮询会把延迟放大到几秒,而且大量无效请求会白白消耗服务端连接资源。
WebSocket是双向全双工的,能实现消息推送,但它的复杂度比SSE高一个量级:需要额外协商升级协议、需要处理二进制帧和掩码、需要维护心跳、前后端都要写一套专门的状态机。而AI对话客户端几乎不需要往服务端推送数据(最多就是发送停止生成的指令),用WebSocket属于“杀鸡用牛刀”。
SSE的优势恰好戳中AI场景的痛点:协议走标准HTTP,防火墙友好,不用额外端口;客户端断线后浏览器或HTTP客户端会基于Last-Event-ID自动重连;服务端可以随时断开;整个技术栈就是InputStream + BufferedReader + while循环,没有状态机,没有帧解析。做Java后端的同学,上手SSE解析的基本门槛就是会读流。
注意:很多人把SSE误以为是长轮询或者WebSocket的另一种叫法,其实是完全不同的东西。SSE本质上是“一个HTTP响应,服务端持续往响应体里写内容”,所以它只能服务端推客户端,客户端中途只能通过关闭连接来中止。
2. 第一阶段:显式调用,手写SSE解析器的日子
最早我接大模型接口时,网络上还找不到很成熟的Java版流式SDK,官方只提供了Python和Node的示例。我当时选的方案是直接用JDK 11自带的HttpClient,配合BodyHandlers.ofLines()把HTTP响应体变成Stream<String>,然后自己解析SSE事件。这也就是典型的“显式调用”——解析逻辑全部裸写在业务代码里。
2.1 用JDK HttpClient手写一个最朴素的SSE读取
核心代码不长,但每一步都有讲究:
HttpClient client = HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(10)) .build(); 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(BodyPublishers.ofString(jsonBody)) .build(); try (Stream<String> lines = client.send(request, BodyHandlers.ofLines()).body()) { Iterator<String> it = lines.iterator(); StringBuilder dataBlock = new StringBuilder(); while (it.hasNext()) { String line = it.next(); if (line.isEmpty()) { // 空行意味着一条SSE事件结束 String payload = dataBlock.toString(); if (payload.startsWith("data:")) { String data = payload.substring(5).trim(); handleData(data); // 交给业务处理 } dataBlock.setLength(0); continue; } dataBlock.append(line).append("\n"); } } catch (Exception e) { log.error("SSE read error", e); }这里关键点有三个:
handleData(data)是真正的业务回调,你需要把它作为Function传入,让上层处理每个token。这一步开始,你的代码里就出现“回调套回调”的味道了——第一个阶段的痛点也是从这里开始的。if (line.isEmpty())判空行,是SSE协议里最容易被忽略的细节。很多AI服务端返回的换行风格不统一,有的用\r\n,有的用\n,BodyHandlers.ofLines()会帮你把平台差异抹掉,但空行判断必须自己处理。- 异常处理要覆盖网络断连、服务端5xx、读取超时三种情况。早期我的做法就是catch Exception然后整个请求重试,但SSE流已经消费了一半,重试时可能重复生成文本,这件事在后面封装阶段才彻底想明白。
2.2 显式调用阶段的三宗罪:样板代码、线程阻塞、可测试性差
手写版本跑通很容易,但用着用着就发现问题了。首先是样板代码泛滥——每次接一个新模型服务(智谱、通义、月之暗面,各家协议细节略不同),就要把上面这坨代码复制一遍,改改URL和JSON字段,然后继续复制。其次是线程阻塞问题:显式调用时,一个请求要占一个线程全程等流读完,一个大模型回答可能持续输出几十秒。用Executors.newFixedThreadPool(10)扛高并发?连接一多线程池直接被打满。第三是没法写单元测试——解析逻辑全嵌在方法和回调里,想模拟一段SSE流来测试,得mock整个HTTP层,痛苦不堪。
当时组里同时接了三个模型供应商,每个供应商的SSE流格式、错误码、重试策略都不一样。意识到继续这样写下去维护成本会失控,我开始进入第二阶段:把SSE读写封装成通用逻辑。
3. 第二阶段:隐式封装,把SSE从业务代码里“拿走”
隐式封装的核心目标,是把“发起HTTP请求 + 解析event-stream + 按事件触发回调 + 错误处理 + 重试”这五件事全部收敛到一层内部组件里,业务方只需要传入请求参数,然后等待一个Publisher<String>或者消费每个token的Consumer。封装之后,业务代码里再也看不到HttpClient、Content-Type、data:这些字样,这就是“隐式”的涵义。
3.1 封装层应该暴露什么API:从回调接口到流式返回
我最终沉淀下来的接口是这么设计的:
public interface StreamClient { CompletableFuture<Void> chat(ChatRequest request, Consumer<String> onToken, Consumer<StreamMeta> onComplete); }这个接口有两个设计要点值得说。
接口里返回的是CompletableFuture<Void>,而不是List<String>或者String。因为流式输出本身是异步的,CompletionStage能天然支持后续编排——你可以.whenComplete()做清理、.exceptionally()做降级。如果返回List,意味着阻塞等待整个流读完才交给你,那就退化成了非流式的普通API。很多封装SDK就是栽在这一步——接口设计成同步返回完整文本,明明底层是SSE,业务层却享受不到流式效果。
Consumer<String> onToken用来接收每一个解析出来的token文本;Consumer<StreamMeta> onComplete用来接收结束元信息。实际使用时,我在内部还封装了一个简单的背压计数:当onToken回调里发生异常时,主动关闭底层连接,停止后续解析。这避免了回调阻塞导致的内存堆积。
3.2 内部实现:用解析器组件处理事件流,而不是裸流
封装内部我拆分成了两层:一是连接层,负责构建HttpClient请求、处理重试、超时;二是解析层,负责把Stream<String>转成标准事件对象。解析层我直接用了okhttp-sse接口,这个库比自研解析器稳定,它对事件流边界、注释行、字段解析都有完善处理:
EventSource.Factory factory = EventSources.createFactory(client); Request request = new Request.Builder() .url(url) .post(RequestBody.create(json, MediaType.get("application/json"))) .header("Accept", "text/event-stream") .build(); EventSource eventSource = factory.newEventSource(request, new EventSourceListener() { @Override public void onEvent(EventSource eventSource, String id, String type, String data) { if ("[DONE]".equals(data.trim())) { onComplete.accept(new StreamMeta(true, null)); eventSource.cancel(); return; } onToken.accept(data); } @Override public void onClosed(EventSource eventSource) { if (dataBlock != null) { // 处理最后一块缓存数据 } } @Override public void onFailure(Throwable t, Response response) { onError.accept(t); } });okhttp-sse内部已经帮你处理好了:
retry:字段自动解析,服务端指定重试时间;- 注释行(以冒号开头)自动跳过;
data:多行合并的逻辑也正确;- 连接意外关闭时会自动重连(前提是你用
EventSource的标准生命周期)。
不过,用okhttp-sse也有它自己的坑:它底层依赖OkHttp连接池,如果同一个域名压到上千并发,OkHttp默认的连接池复用机制会表现出一些诡异行为——具体我在下面踩坑章节详细说。
3.3 隐式封装带来的连带收益:可测试性、可切换性、可观测性
封装完以后,有三个我没想到的收益。
一个是可测试性。现在我可以写这样的单测:构造一段"data: 你好\r\ndata: 世界\r\n\r\n"字符串,喂给解析器,断言onToken按顺序收到两个事件。不需要mock任何网络,SSE解析逻辑的bug在本地就能复现。
第二个是可切换性。业务方依赖的是StreamClient接口,今天用OpenAI兼容接口,明天换成国产大模型的SDK,只要适配层实现同一个接口即可,业务代码完全不动。这种解耦在“多模型接入”场景里非常关键——AI时代模型服务更新迭代太快,接口不稳定是常态。
第三个是可观测性。我在封装层内置了统计点:连接建立时间、首字节延迟、每token间隔、连接时长、异常类型分布。这些数据直接打进监控系统,能很清楚看到每个模型供应商的响应质量。这些在显式调用阶段完全没法做,因为逻辑分散在各业务代码里。
4. 第三阶段:虚拟线程登场,并发轻松冲破“连接数天花板”
前两个阶段,并发能力始终卡在“线程池大小”这件核心症结上。传统线程池方案下,每个SSE连接要占一个平台线程,而平台线程是昂贵资源——默认栈空间1MB起步,切换要内核参与。到大模型场景,一次完整流式响应可能要持续15到60秒,即使只开100个并发,最坏情况下就有100个线程空转等待网络IO。
4.1 虚拟线程为什么是SSE场景的天然搭档
Java 19开始引入虚拟线程预览,Java 21变成正式特性。虚拟线程的引入彻底改变了“一个连接一个线程”的成本模型:虚拟线程在阻塞点(比如读InputStream、读Socket)会自动挂起,让出底层的物理线程(平台线程),IO数据到达之后再自动恢复。这意味着你可以创建成千上万个虚拟线程,真正占用底层资源的只有少量的平台线程(默认是CPU核心数)。
用生活类比来理解:传统线程池就像银行网点,每个柜员(平台线程)服务一个客户直到业务办完,客户多了只能排队;虚拟线程则是取号排队系统,柜员数量固定,但一个柜员可以在客户等待填表时(阻塞IO),立刻去服务下一位客户。
SSE场景正是典型的“阻塞式IO + 长时间等待”模型。读流的时候,大部分时间都花在等服务端推送下一个数据块上,这个等待过程在虚拟线程里几乎零成本。你不需要回调、不需要事件驱动,可以回到最简单的同步阻塞写法——但并发能力却完全不受平台线程限制。
4.2 用虚拟线程重写SSE消费逻辑
改造后的核心代码可以回归到近乎“土的掉渣”的同步风格:
public void processStream() { try (var executor = Executors.newVirtualThreadPerTaskExecutor();) { // 一个虚拟线程处理一个SSE连接 CompletableFuture<?>[] futures = urls.stream() .map(url -> CompletableFuture.runAsync( () -> handleSseConnection(url), executor)) .toArray(CompletableFuture[]::new); CompletableFuture.allOf(futures).join(); } } private void handleSseConnection(String url) { HttpClient client = HttpClient.newBuilder() .executor(Executors.newVirtualThreadPerTaskExecutor()) .build(); // 然后代码可以像普通同步IO一样逐行读流 // 但底层已经不会占用物理线程了 }关键点是HttpClient的executor也要换成虚拟线程池。Java 21的HttpClient默认虚拟线程化,但显式指定更保险。改造之后的压测效果非常直接:在同一台4核8G的机器上,用固定线程池(200个线程)跑SSE,连接数到150左右就开始频繁超时;切换虚拟线程后,压到2000个并发SSE连接,响应延迟基本稳定。
4.3 虚拟线程对SSE封装的“反向优化”
一个有意思的事情是:虚拟线程普及之后,我之前封装的那套“隐式回调式SSE解析”反而显得多余了。回调是为了不阻塞线程而做的妥协,但虚拟线程时代,同步阻塞IO不再需要这种妥协。第二次重构时,我把okhttp-sse的EventSourceListener换成了直接读取:
try (var response = client.send(request, BodyHandlers.ofLines())) { Stream<String> lines = response.body(); Iterator<String> it = lines.iterator(); while (it.hasNext()) { String line = it.next(); if (line.isEmpty()) { handleSseEvent(dataBuffer.toString()); dataBuffer.setLength(0); } else { dataBuffer.append(line).append("\n"); } } }这段代码看起来和我第一阶段的显式调用很像,区别在于:它运行在虚拟线程里,不再占平台线程;解析逻辑沉淀在底层组件里,业务方看到的还是StreamClient接口。所以“显式调用”和“隐式封装”在虚拟线程时代找到了一个最佳平衡点:底层同步阻塞,顶层异步接口,中间虚拟线程兜底。
但是!虚拟线程并不是无脑套上去就完事了。接下来这部分是我认为本文价值最高的部分,全部来自实际排障现场。
5. 三个真实排障现场:idle timeout、虚拟线程锁、HTTP/1.1六连接陷阱
5.1 报错“stream disconnected before completion: idle timeout waiting for sse”的排查链路
有段时间,生产环境频繁出现一条报错:stream disconnected before completion: idle timeout waiting for sse。现象是:用户发一句话,模型输出到一半,流突然断了,用户看到对话在中间某处戛然而止。
排查第一步,我先确认是我们的服务端主动断的,还是上游AI服务端断的。抓包后发现,TCP连接是在服务端等了大概几十秒没有任何数据后才被对端关闭的。再往上游走,发现是云网关层的空闲超时机制——即使SSE连接还活着,如果一段时间内没有新数据包经过负载均衡层,就会被判定为“空闲连接”,直接掐断。
这说明一个关键原则:SSE连接长时间“静默”不等于死连接,但基础设施会把它当成死连接。解决办法也很直接:
- 在封装层加上心跳机制:每隔15到20秒发一个注释行(SSE协议里,以冒号开头的注释行会被客户端忽略,但会作为数据包经过网络链路);
- 在
HttpClient上配置合理的readTimeout,不能太长也不能太短,太短会误杀慢速模型,太长会长时间占着线程不释放。 - 把云网关的idle timeout调大,但这不是解决问题的根本方法,依赖具体部署环境,能调最好,不能调就靠心跳兜底。
这个坑的核心教训是:SSE不是一个“有来有回”的协议,基础设施常会把它当成半死连接。解决方案必须在客户端和服务端之间做双向兜底。
5.2 虚拟线程其实也会“假阻塞”:同步块和锁的坑
这是我在把SSE接入到虚拟线程后,最意外的一个发现。
虚拟线程的设计文档里特别提到一个“固定”(pinning)现象:当一个虚拟线程进入synchronized关键字保护的代码块,并且在块内发生阻塞IO时,这个虚拟线程会被“钉”在底层平台线程上,导致平台线程无法被其他虚拟线程复用。如果并发请求里面有大面积synchronized加锁 + 阻塞IO的组合,虚拟线程的性能优势会严重退化,甚至比固定线程池还差。
我当时封装SSE解析时,为了统计并发量给接口加了一个AtomicInteger计数器,然后用synchronized保护了一段日志打印逻辑。压测时发现虚拟线程模式下QPS反而比固定线程池低,查了JFR火焰图才发现问题出在这里。
修复方案有两个方向:
- 代码层面,把
synchronized替换成java.util.concurrent.locks.ReentrantLock。在虚拟线程的阻塞语义下,ReentrantLock不会导致平台线程固定(JDK 24中synchronized对虚拟线程的处理已经改进,但如果你还在用JDK 21,务必注意这个差异); - 设计层面,尽量减少SSE解析链路里的共享可变状态。能局部缓存就局部缓存,能无锁就无锁。
另外一个小经验:用虚拟线程时,ThreadLocal也要慎重。虚拟线程数量可以轻易上千,但ThreadLocal里的对象是每个线程独立持有的,如果不清理,内存占用会随虚拟线程数量线性膨胀。建议用ScopedValue(JDK 21+ 预览)或者显式清理。
5.3 为什么压测到300并发,服务出端口连接立刻“假死”
还有一个在封装阶段就踩过的坑,虚拟线程之后依然会在连接层出现:HTTP/1.1同域名下,HTTP客户端默认最大连接数是硬顶。
现象是这样的:用OkHttp或JDK HttpClient压测同一个模型的流式接口,并发从100升到300时,突然大量请求报connection pool timeout或直接卡住不返回。排查后发现根本原因在HTTP/1.1规范本身——一个客户端对同一个目标的并发TCP连接数量有限制(浏览器场景是6个左右,Java HttpClient连接池默认也会控制总连接数和每个域名的路由数)。SSE每个请求都是长连接,意味着每个用户请求都会占掉一个连接池槽位,压测几十个并发就会把连接池打爆。
解决方案是把连接池上限调大,或者走HTTP/2多路复用。
JDK HttpClient开启HTTP/2只需要改一行:
HttpClient client = HttpClient.newBuilder() .version(HttpClient.Version.HTTP_2) .executor(Executors.newVirtualThreadPerTaskExecutor()) .build();HTTP/2下,同一个连接可以承载多个SSE流,连接池压力急剧下降。我在生产环境把SSE接入层全面切到HTTP/2之后,同样压测300并发,错误率从8%降到0.2%。如果上游服务不支持HTTP/2(部分老网关不支持),那就只能用调大连接池上限来缓解,但这不是一个优雅的方案。
6. 选型参考与最终建议
写到这里,三个阶段的优缺点其实已经比较清晰了。我用一张表总结一下,方便你根据自己项目情况直接对号入座:
| 阶段方案 | 并发能力 | 代码复杂度 | 维护成本 | 适用场景 |
|---|---|---|---|---|
| 显式调用(手写解析) | 低,受线程池限制 | 高,样板代码多 | 高,每接入一个模型都要复制改 | 单体Demo、一次性脚本 |
| 隐式封装(SDK化+回调) | 中,受线程池/连接池限制 | 中,但抽象良好 | 低,统一升级入口 | 生产环境但JDK版本较低(如17及以下),需要对接多模型 |
| 虚拟线程+同步阻塞封装 | 高,数千并发无压力 | 低,代码回归同步风格 | 低,但依赖JDK21+ | 新项目、线上JDK已升级到21+,SSE类接口专用 |
我的推荐顺序很明确:新项目直接用第三阶段方案,JDK 21 + HTTP/2 + 虚拟线程 + 同步阻塞读取。旧项目在JDK 17上跑的话,可以先把第二阶段封装做好,沉淀解析和重试逻辑,等JDK升级后再把执行层换成虚拟线程——封装层抽象如果做得好,这一步切换成本极低。
还有几个偏个人经验的注意点最后提一下:
一,不要为了“流式”而放弃超时控制。SSE是长连接,但不代表超时时间可以设成无限。建议连接建立超时设10秒,读超时设1到2分钟,空闲心跳设20秒。模型首字延迟如果超过读超时,多半是服务端或网络有问题,及时断掉比干等着好。
二,如果同一份SSE数据需要同时投递给多个下游消费者(比如既要存数据库,又要推给前端,还要喂给审计日志),不要在解析线程里做这些事。解析线程只做解析,下游消费丢给另外的线程池或消息队列,避免一个慢消费者阻塞整个SSE接收。
三,SSE解析器一定要保留“原始事件”的日志能力。我在生产环境排查问题时,最救命的就是能重放某次会话的完整原始SSE字节流。模型侧和网关侧的锅,靠日志一锤定音。
从我个人的实际体验来说,这套方案演进下来,最核心的感悟是:Java + AI的流式接入,最大的瓶颈从来不是模型能力,而是连接和并发这两个工程问题。第一阶段我在显式解析时,被线程池卡得焦头烂额;第二阶段封装完SDK,又遇到连接池和空闲超时;到第三阶段切到虚拟线程后,SSE在Java侧终于不再像“二等公民”了。如果你也正在做类似接大模型流的Java项目,建议直接照着第三阶段的方式搭一套,你会回来谢谢虚拟线程的。