☰
Java+AI 流式输出:SSE 从手写解析到虚拟线程优化全攻略
2026/10/3 15:36:52 网站建设 项目流程

做 AI 应用接入大模型流式输出,只要你碰的是 Java 后端,就一定绕不开 SSE(Server-Sent Events)。我第一次接这类需求时,以为这只是个“普通接口多等几秒而已”,结果项目从手写 HTTP 流式解析,到后面接 Spring 的封装,再到现在用虚拟线程撑起大批量长连接,踩的坑比预想中多得多。这篇文章就把这条路径完整复盘一遍,聊聊 Java+AI 场景下 SSE 从显式调用到隐式封装、再到虚拟线程性能优化的演进过程。

这篇文章适合正在做 AI 应用开发的 Java 工程师,也适合后端技术负责人评估流式接口改造方案。你会看到协议层细节、代码示例、封装思路、压测心得,以及一些常规文档里不会写的问题排查经验。

1. 为什么 AI 流式输出场景偏偏选了 SSE

1.1 SSE 是一段不会结束的 HTTP 响应

很多第一次接触 SSE 的人会下意识把它和“轮询”混在一起,但 SSE 的本质是一次完整的 HTTP 请求,只是服务端不立刻返回结果,而是把连接一直开着,持续向客户端推送文本数据。

服务端响应的 Content-Type 必须是text/event-stream,数据格式非常朴素,每一行是一个字段,字段和字段之间用空行分隔成一个事件块。一个典型的事件块长这样:

data: {"token":"你"} data: {"token":"好"}

客户端收到的是普通 HTTP 响应体,服务端通过不断追加这个体来持续传输数据。这种设计让人很容易上手:不需要像 WebSocket 那样先走一通握手协议、升级连接,SSE 只是“一个响应迟迟没读完”的普通请求。

SSE 协议里定义了五种字段:data表示数据内容,event表示自定义事件名,id表示事件编号,retry表示断线重连的间隔时间,以及用冒号开头的注释行。注释行看起来没用,但它是服务端用来“刷心跳”的常用手段——发一行: ping就能让连接保持活跃,又不会污染业务数据。

1.2 和 WebSocket 对比后的选型逻辑

很多人会问,流式推送为什么不用 WebSocket?WebSocket 确实更全能,但 AI 对话这个场景,实际需要的只是一条“服务器单向推到客户端”的下行通道。

对比维度SSEWebSocket
连接模型普通 HTTP 长响应TCP 专用长连接
协议复杂度低,基于文本行较高,需要帧解析和掩码
握手方式无需额外握手需要 101 升级握手
数据流向单向下行双向
自动重连内置,基于 Last-Event-ID需要自己实现
调试成本浏览器控制台直接看 Network 流需要专门的调试面板或工具
AI 流式场景完全匹配功能过剩

从工程视角看,AI 对话基本是“客户端发一个问题,服务端返回一段流式文本”,几乎没有实时上行需求。用 WebSocket 服务端要维护会话状态、处理心跳帧、考虑消息分片,而 SSE 只需要把响应写好就行。

还有一个现实原因:大模型厂商的服务端接口大多直接输出 SSE 格式,Java 后端要做的不是发明协议,而是接住上游的流,再原样或加工后推给下游。这个单向链路用 SSE 最顺手。

1.3 AI 应用里的完整 SSE 调用链路

在一个典型的 Java+AI 项目中,SSE 会贯穿两层链路。

用户在前端点“发送”,前端用EventSource发起请求,Java 后端收到这个请求后,再去调用大模型接口。大模型接口通常也是 SSE 流式返回,Java 后端一边读上游 token,一边通过自己的 SSE 连接把 token 推给前端。前端每收到一块新数据,就把它渲染到对话框里,形成“打字机”效果。

这中间要处理三件事:上游模型的流式数据解析、后端到前端的协议转换、以及两段连接的异常处理。很多人只关注“调用 API”,忽略了传输链路上的断连和超时,后面会详细讲。

2. 显式调用:手写 SSE 客户端的那段日子

2.1 早期实现:HttpClient + BufferedReader 手动读流

最早我写 SSE 客户端,完全是“硬读”。Java 自带java.net.http.HttpClient支持把响应体作为InputStream接收,然后就能用BufferedReader一行行往下读。代码看起来很简单:

HttpClient client = HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(10)) .build(); HttpRequest request = HttpRequest.newBuilder() .uri(URI.create("http://localhost:8080/v1/chat/completions")) .header("Content-Type", "application/json") .header("Accept", "text/event-stream") .POST(BodyPublishers.ofString(""" {"model":"demo","prompt":"你好","stream":true} """)) .build(); HttpResponse<InputStream> response = client.send(request, HttpResponse.BodyHandlers.ofInputStream()); if (response.statusCode() != 200) { // 这里要先把错误体完整读取出来记录,再抛异常 } try (BufferedReader reader = new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; while ((line = reader.readLine()) != null) { if (line.startsWith("data:")) { String payload = line.substring("data:".length()).trim(); if ("[DONE]".equals(payload)) { break; } // 用 Jackson 或 Gson 把 payload 解析成 JSON // 从 choices[0].delta.content 里取出增量 token } } }

这段代码放到真实环境里能跑通 demo,但它真的太“显式”了——所有协议细节都暴露在业务代码里,所有异常都要业务侧自己兜底。

2.2 显式调用最容易忽略的协议细节

SSE 的事件块解析并不是“读一行就完事”这么简单。按协议规范,一个事件块里可能出现多行data,多行内容会被拼接成一段完整数据,中间用换行符连接。也就是说下面这种格式是合法的:

data: {"delta":"你"} data: {"delta":"好"}

这段内容会拼成{"delta":"你"}\n{"delta":"好"},解析时不能只取第一行。

event、id、retry字段也各有用途。id字段特别重要,客户端断线重连时会带Last-Event-ID头,服务端可以根据这个 ID 决定从哪条事件之后开始重推。retry字段则是服务端建议客户端下次重连的等待毫秒数。

我后来写了一个通用解析方法,负责把一个事件块的多行内容转换成结构化对象:

static SseEvent parseEventBlock(List<String> lines) { String data = ""; String event = "message"; String id = null; Integer retry = null; for (String raw : lines) { String line = raw.endsWith("\r") ? raw.substring(0, raw.length() - 1) : raw; if (line.isEmpty() || line.startsWith(":")) { continue; } int sep = line.indexOf(':'); String field = sep < 0 ? line : line.substring(0, sep); String value = sep < 0 ? "" : line.substring(sep + 1); if (value.startsWith(" ")) { value = value.substring(1); } switch (field) { case "data" -> data = data.isEmpty() ? value : data + "\n" + value; case "event" -> event = value; case "id" -> id = value; case "retry" -> retry = Integer.valueOf(value); } } return new SseEvent(id, event, data, retry); }

这套解析逻辑看着繁琐,但一旦接入多家大模型供应商,就会发现只处理data字段的代码完全不够用,有的是事件名不同,有的是断点续传需要id。

2.3 显式调用阶段踩过的坑

这一阶段最大的体会是:SSE 客户端最难的从来不是“怎么读”,而是“读的过程中连接挂了怎么办”。

请求超时不能瞎设。很多人在HttpRequest上加.timeout(Duration.ofSeconds(30)),然后发现流式响应一到 30 秒就被客户端主动断掉。因为 HttpClient 的 timeout 是“整个请求从开始到结束”的总超时,而 SSE 恰恰是一个“持续时间很长”的请求。正确做法是把 connectTimeout 单独设置好,不在请求级别设置总超时,超时控制交给专门的空闲读超时逻辑。

必须处理非 200 状态。模型服务如果鉴权失败、限流或参数错误,会直接返回 4xx,响应体里是一个普通 JSON 错误信息,根本不是 event-stream 格式。流式解析代码遇到这种混入的响应直接懵。所以进入读流之前,先检查状态码,把错误体完整读出并记录。

断线重连和心跳是刚需。早期我没有做重连逻辑,结果一次网络抖动就让用户看到“生成中断”,而且没有恢复机制。后来才理解 SSE 协议自带 Last-Event-ID 重连机制就是给这个场景用的,客户端维护一个最后处理的事件 ID,重连时带上去,服务端决定要不要补发。

背压很容易被忽略。如果上游生成速度快,下游消费者处理不过来,BufferedReader.readLine()会持续读到新数据,内存里的消息队列就会膨胀。尤其在做 Agent 场景时,中间的日志、状态更新、工具调用结果都要处理,这一段消费链路要设计好缓冲和丢弃策略,不能无脑 while 循环。

2.4 显式调用留下的实际价值

虽然显式调用代码最繁琐,但我不建议新手一上来就套封装框架。因为后续排查线上问题,比如“为什么读到一半断了”“为什么这个事件没触发”,最终都要回到协议的原始格式去看。

亲自动手实现一遍 SSE 解析,之后再看各种封装库的源码,就会清楚它内部到底在做什么。

3. 隐式封装:SSE 从手写到开箱即用

3.1 服务端发送:SseEmitter 和 WebFlux 怎么选

Java 后端给前端推 SSE,如果是 Spring MVC 项目,最直接的方式是返回SseEmitter。

@PostMapping("/chat") public SseEmitter chat(@RequestBody ChatRequest request) { SseEmitter emitter = new SseEmitter(0L); // 不设超时,由业务控制何时完成 modelCallExecutor.execute(() -> { try { // 模拟模型逐步返回 emitter.send(SseEmitter.event() .id("1") .name("token") .data("你好")); emitter.send(SseEmitter.event() .id("2") .name("token") .data(",我是 AI")); emitter.complete(); } catch (Exception ex) { emitter.completeWithError(ex); } }); return emitter; }

SseEmitter.event()返回一个事件构建器,id、name、data正好对应 SSE 协议里的字段,底层会帮你拼成text/event-stream响应。这里emit线程的选择值得注意:模型调用如果是阻塞式 HttpClient 请求,就不要占用 Tomcat 的请求工作线程,否则长连接一多,容器线程很快被占满。这也是后面虚拟线程切入的入口。

如果你已经用了 Spring WebFlux 的响应式栈,可以直接返回Flux<ServerSentEvent>:

@GetMapping(value = "/chat", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<ServerSentEvent<String>> chat() { return Flux.interval(Duration.ofSeconds(1)) .map(i -> ServerSentEvent.builder("token " + i) .id(String.valueOf(i)) .build()); }

两种方式的选择建议很简单:项目原本就是 Spring MVC,用 SseEmitter,不要为了 SSE 把整个技术栈换成 WebFlux;项目本来就是全链路响应式,用 Flux 更自然。强行混用会让团队同时维护两套编程模型,代价远超收益。

3.2 客户端接收:WebClient 和 OkHttp 的封装差异

Java 后端调用大模型接口时,推荐优先用 Spring WebClient 的bodyToFlux。它对 SSE 做了内置支持,几行代码就能把上游数据流变成一个响应式的 Flux。

WebClient client = WebClient.builder() .baseUrl("http://localhost:8080") .build(); Flux<String> stream = client.post() .uri("/chat") .bodyValue(request) .accept(MediaType.TEXT_EVENT_STREAM) .retrieve() .bodyToFlux(String.class);

默认的bodyToFlux(String.class)拿到的是data:字段的内容,如果你还想拿到事件的id、event、retry元信息,可以用bodyToFlux(ServerSentEvent.class):

Flux<ServerSentEvent<String>> events = client.post() .uri("/chat") .bodyValue(request) .accept(MediaType.TEXT_EVENT_STREAM) .retrieve() .bodyToFlux(new ParameterizedTypeReference<ServerSentEvent<String>>() {});

如果你的项目已经大量使用 OkHttp,不考虑引入 WebClient,可以用okhttp-sse扩展包:

OkHttpClient okHttpClient = new OkHttpClient.Builder().build(); EventSources.createFactory(okHttpClient).newEventSource(request, new EventSourceListener() { @Override public void onEvent(EventSource eventSource, String id, String type, String data) { // 每收到一个事件回调一次 } @Override public void onClosed(EventSource eventSource) { // 连接正常关闭 } @Override public void onFailure(EventSource eventSource, Throwable t, Response response) { // 连接意外失败 } });

OkHttp 的封装把断线重连逻辑做了一部分,但它只负责客户端侧,服务端实现要自己写。WebClient 的优势在于和 Spring 生态的响应式链路融合得更紧密,适合“上游到下游全程流式”的架构。

3.3 把流式接口统一抽象成 StreamResponse 模式

工程里接入多个模型供应商后会发现,每家模型的流式协议大体相同但细节各有差异。有的返回[DONE],有的返回data: [DONE];有的 token 在choices[0].delta.content里,有的在data.choices[0].message.content里。

这个时候最该做的不是到处写 if-else,而是抽象一个统一回调接口:

public interface StreamCallback { default void onStart(StreamContext context) {} default void onToken(StreamContext context, String token) {} default void onFinish(StreamContext context) {} default void onError(StreamContext context, Throwable throwable) {} }

每个供应商实现一个适配模块,内部负责把各自的 SSE 事件解析成统一的onToken回调。业务层只需要关心“我收到了哪个字符串”,不需要关心它是 OpenAI 的格式还是国产模型的自定义格式。

这种封装的本质,是把我前面手写的那套协议解析、断线重连、超时控制、背压处理逻辑沉淀成公共组件。它不只是省代码量,更关键的是让所有接入方都拿到一致性的可靠性保障。

4. 虚拟线程:SSE 并发量的真实拐点

4.1 平台线程池和 SSE 长连接是一个天然矛盾

传统 Tomcat 默认请求线程池是 200 个平台线程,SSE 却是“一个连接长时间占着一个线程直到断开”的模型。如果同时有 200 个用户挂着流式对话,请求线程池被占满,其他普通接口全部排队。

这个问题的本质是:等待大模型返回期间,线程完全阻塞在 IO 上,不干任何事,但仍然占用内存和调度资源。一个平台线程的默认栈内存大约 1MB,启动 5000 个线程就近乎 5GB 的内存开销,这还没算操作系统上下文切换的成本。

前面提到的方式,无论是 NIO 还是响应式 WebClient,都是非阻塞方案的变体。但这条路有两个代价:一是代码风格从同步改造成链式回调,整个团队要重新学习;二是与现有 Spring MVC、MyBatis、事务管理等阻塞式技术栈耦合时非常别扭。

4.2 虚拟线程怎么解决“阻塞等待”问题

虚拟线程是 JDK 21 正式提供的能力,它和平台线程最大的区别是:平台线程是操作系统调度的,虚拟线程是 JVM 内部调度的。

JVM 会创建少数平台线程作为载体线程池,虚拟线程运行在载体线程上。当虚拟线程执行到阻塞点(比如读取网络流、Thread.sleep、LockSupport.park),JVM 会自动把这个虚拟线程从载体线程上摘下来,让载体线程去执行另一个虚拟线程。整个过程由 JVM 调度器完成,无需业务代码介入。

生活化理解:平台线程像固定工位,一个人坐在工位上等快递,工位就浪费了;虚拟线程像临时工,快递没到就先干别的活,快递到了再回来接着干。SSE 场景里大量连接都在“等 token”,正好是虚拟线程最擅长消化的一类负载。

4.3 虚拟线程 + SSE 的落地改造

在 Spring Boot 3.2 及以后版本,最简单的开启方式是在配置文件里设一个开关:

spring: threads: virtual: enabled: true

开启后,容器接收请求时会用newVirtualThreadPerTaskExecutor()默认创建的虚拟线程执行器处理请求。对 SSE 接口来说,意味着每个 SSE 长连接都可以占一个虚拟线程,而不是平台线程。

如果你的项目不打算全局开启,也可以只在流式接口的“模型调用”环节用虚拟线程执行器:

ExecutorService modelCallExecutor = Executors.newVirtualThreadPerTaskExecutor(); @PostMapping("/chat") public SseEmitter chat(@RequestBody ChatRequest request) { SseEmitter emitter = new SseEmitter(0L); modelCallExecutor.execute(() -> { try { // 这里内部是阻塞式 HttpClient 调用大模型 modelService.streamChat(request, token -> { try { emitter.send(SseEmitter.event().name("token").data(token)); } catch (IOException e) { throw new UncheckedIOException(e); } }); emitter.complete(); } catch (Exception ex) { emitter.completeWithError(ex); } }); return emitter; }

虚拟线程的特点是“创建成本低、用完即丢”,所以不要用线程池去池化它。每次调用就execute,执行完自动销毁即可。

我做过一个简单对比测试:一台 8 核 16G 的开发机器,模拟 300 个 SSE 客户端连接,每个连接保持 60 秒,服务端每 2 秒向下推一条消息。平台线程池模式下,Tomcat 默认线程很快被打满,普通接口响应开始出现几秒延迟;切到虚拟线程后,CPU 占用平稳,普通接口响应时间基本没受流式连接影响。当然这个数据只是中午跑的小实验,权当参考,但虚拟线程在长连接场景的收益方向是明确的。

4.4 虚拟线程使用边界和几个隐藏问题

虚拟线程不是银弹,有几类场景反而要格外小心。

synchronized会让虚拟线程 pins 到载体线程上。如果在虚拟线程里锁竞争激烈,虚拟线程无法被摘除,会直接占用一个平台线程,并发效果立刻退化。能用ReentrantLock的场景尽量替换。

ThreadLocal 在虚拟线程里虽然能用,但开启虚拟线程的 ThreadLocal 复制成本更高。不要在虚拟线程里传重量级上下文对象,尤其跨线程池传数据时要重新设计。

虚拟线程适合 IO 密集型,不适合 CPU 密集型。如果你在流式任务里做大量 JSON 大字段解析、递归计算、正则回溯,这些计算不能靠虚拟线程换并发,反而会因为线程库增加而带来额外调度开销。

另外,不要在代码里为了“看起来用到了虚拟线程”而手动创建大量虚拟线程去做无限循环任务。SSE 长连接的本质是 IO 等待,但如果连接建立后业务代码本身不阻塞、只拼 CPU 死等,虚拟线程帮不上忙。

5. 常见问题与排查技巧实录

5.1 SSE 问题速查表

现象可能原因排查方向快速处理
前端 EventSource 不触发 open响应 Content-Type 不是 text/event-stream看 Network 响应头后端显式指定produces = TEXT_EVENT_STREAM_VALUE
数据不实时刷新,一次性返回反向代理层开启了响应缓冲检查网关/Web 代理配置关闭代理缓冲,添加X-Accel-Buffering: no响应头
连接在固定时间后断开网关或客户端空闲超时看服务端日志、网关日志服务端定期推送: ping注释行心跳
收到乱码字符编码不一致检查响应头 charset 字段统一使用 UTF-8,显式指定编码
断线重连后内容重复没有维护 Last-Event-ID检查客户端是否传重连 ID服务端按 ID 实现断点续推
开启虚拟线程后性能下降synchronized pin 或 CPU 密集计算看线程 dump、锁竞争替换锁、拆分CPU密集任务

5.2 遇到 stream disconnected before completion 怎么查

这个报错本质是 SSE 流没有读完就断开了,常见于客户端、网关、服务端三者的超时策略不一致。

第一步看客户端。如果用的是 WebClient,检查是否设置了全局超时或者连接空闲超时,响应式框架有多个超时维度,很容易误伤长连接。第二步看网关。Nginx 默认proxy_read_timeout是 60 秒,SSE 连接超过这个时间没有新数据,网关直接断开。配合服务端心跳注释行,并调大proxy_read_timeout,或者干脆在特定路径上关闭缓冲。

第三步看服务端。服务端如果正常complete()了,但客户端以为还该继续收数据,也会出现这种报错。排查时要区分是“符合预期的正常关闭”还是“异常中断”,日志里记录的堆栈是关键。

5.3 断线续传和业务幂等怎么设计

最常见的场景是:用户提问后,模型已经生成了 30 个 token,第 31 个 token 发送时断网。重连后如果把整个请求重发一遍,模型要重新生成全部 token,体验很差,还可能造成费用重复。

SSE 自带的id字段就是干这个的。服务端持续为每个事件生成递增 ID,客户端重连时带上Last-Event-ID,服务端从该 ID 之后的事件继续推送。但这里有个复杂点:如果模型服务侧没有缓存中间生成结果,服务端拿到Last-Event-ID也无从恢复。所以更通用的做法是服务端提前把 token 写入本地缓存或消息队列,重连时从缓存里补发。

如果模型服务不支持断点续传,至少要在业务层做幂等:用一个业务请求 ID 标识整轮对话,重连时服务端检查是否已有部分结果,要么丢弃重来,要么从最近一次成功写入的 checkPoint 继续,具体取舍看成本和产品体验要求。

5.4 把“流式”当作一等公民来建设

经历了手写、封装、虚拟线程三个阶段后,我发现流式接口的工程化不能只停留在“能通”层面,要当作独立的传输基础设施来建设。

监控层面要单独统计 SSE 连接数、SseEmitter 完成数、超时数、异常断开数。连接数突增可能意味着有人在刷接口,错误数突增往往和上游模型服务质量波动直接相关。日志层面不要把所有 token 都打出来,完整对话内容体量太大,只记录“第几条事件、多少字节、耗时多久”即可。存储层面流式数据是增量到达的,要设计好从临时缓冲到最终落库的路径,不能等全部生成完才写库。

这些细节决定了流式功能上线后是“能看”还是“真的好用”。

最后再分享一点个人经验

我最早写 SSE 客户端时也很自信,觉得协议看着简单,几行代码就能搞定,结果一碰真实网络环境就出各种问题。现在再看,SSE 真正的复杂度不在协议本身,而在长连接的可靠性和并发承载。

对准备入场的团队,我的建议是:第一步先把 text/event-stream 的协议格式吃透,亲手解析一次事件流;第二步再考虑引入框架封装,用 WebClient 或 okhttp-sse 把脏话细节挡在业务层之外;第三步才是上虚拟线程,在长连接并发量真正起来之后做优化。顺序不能反,反了会在排查问题时无从下手。

虚拟线程和 SSE 确实是目前 Java+AI 组合里很舒服的一对搭档。同步代码风格不变,却能享受到接近响应式架构的高并发红利。不过要记住,虚拟线程解决的是 IO 等待,不是 CPU 计算,选型时把这层账面算清楚,后面才不会踩坑。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询