Java接入大模型SSE流式输出:从手写解析到虚拟线程的实践演进
最近帮团队重构了一个 AI 网关的流式接入模块翻出之前的代码发现 SSE 相关的实现已经迭代了三个版本第一版是在业务代码里手写 HttpClient 逐行解析响应流第二版抽成了通用 Client调用方一句话就能拿到流式消息第三版切到 JDK 21 的虚拟线程之后单机长连接数直接上了一个量级。这个演进过程挺有代表性也踩了不少坑今天把这几个阶段的技术细节和设计思路完整梳理一遍。这篇文章主要写给正在接入大模型流式输出的 Java 后端同学、想自己封装 SSE SDK 的人以及在做高并发长连接方向选型的朋友。内容会覆盖 SSE 协议本身的格式细节、客户端显式解析时的常见陷阱、怎么做一次不过度设计的封装以及虚拟线程在 SSE 场景下的原理和实测效果最后会聊几个生产环境里高频出现的故障排查链路。1. 为什么大模型流式输出SSE 成了默认选项1.1 从一次真实的 AI 对接需求说起事情起源于团队要接入一个大模型的对话接口。接口的返回方式和传统 REST 完全不同用户问一句话服务端不会等全文生成完再返回而是把 token 一个一个、或者一小批一小批地推过来客户端边收边显示效果上就是打字机模式。当时前端同事第一反应是上 WebSocket理由是听说 WebSocket 支持双向、实时性更好。但后端这边其实倾向于更简单的方案因为需求只有一个方向服务端向客户端推送 token客户端不需要频繁上行打断。这种单工流式场景SSEServer-Sent Events服务器发送事件几乎是成本最低的解。那会儿我做了个简单的对比把 SSE、WebSocket、轮询三种方案在 AI 场景下的表现列出来维度SSEWebSocket轮询数据方向服务端单向推送配合 POST 请求实现先上行后下行全双工双向客户端主动拉取协议基础普通 HTTP响应类型为 text/event-stream独立的 ws:// 协议需要升级握手普通 HTTP代理穿透性好普通 Nginx、负载均衡都能处理中链路里任一代理不支持就得改造好自动重连浏览器 EventSource 原生内置没有内置需自行实现没有内置服务端复杂度低本质是保持一个响应流高需要管理连接生命周期最低但实时性差AI 场景契合度高token 单向流动即可满足偏高适合需要客户端随时中断/改写的 Agent 场景低token 延迟和请求压力都不可接受结论很直接如果只是让大模型的输出实时显示出来SSE 的复杂度是三者里最低的。前端用 EventSource 直接连或者用 fetch API 去读 ReadableStream都行。对后端来说它甚至不需要引入任何额外框架Servlet 容器里开个异步响应就能写。1.2 SSE 协议格式最少要懂的四个细节SSE 本质上不是新协议它是 HTTP 响应的一种特殊形态。服务端把响应头里的Content-Type设为text/event-stream然后不断往响应体里写事件文本。每个事件块由若干字段行和一个空行组成常用的字段就这几个retry: 5000 id: msg-1 event: message data: {content:第一段} data: {content:第二段}data:事件的数据内容一般是一行 JSON。规范允许一个事件有多行 data多行在客户端侧会以换行符连接成一个整体。event:自定义事件名客户端可以根据这个名字分发到不同的处理逻辑。不写就默认是 message 事件。id:事件 ID。客户端重连时可以通过Last-Event-ID请求头告诉服务端我收到哪了服务端据此续传。retry:重连间隔毫秒。客户端断线后按这个时间重试。注释行以冒号开头的行比如: keep-alive属于心跳客户端要忽略但不能当成异常。空行一个事件结束的标志。客户端必须等到空行才认为一个完整事件到达。这个格式看着简单但它决定了客户端解析的基本思路不是读一行处理一行而是读行、拼装、遇空行提交。很多人手写解析时忽略了这个细节数据一多就出事。1.3 Java 服务端看到的 SSE 长什么样在 Java 侧做 SSE 客户端很多人第一次看到的是这样的现象一次HttpClient.send(...)调用发出去返回的不是一个完整 JSON Body而是一个永远不会结束的 InputStream。你从流里读到的是一段一段的文本每个块对应一个 token 事件。服务端不主动断开这个流可以持续几分钟甚至几十分钟。这意味着任何要把响应整体读进内存再解析的写法在这里都不适用。你必须在一个持续存活的连接上做增量解析边读边把事件抛给上层业务。而持续存活这四个字就引出了显式调用阶段那些让人头皮发麻的坑。2. 手写 SSE 客户端显式调用阶段踩过的坑2.1 第一版代码长什么样刚开始我用的 JDK 自带 HttpClient没有用第三方库。核心代码其实不长HttpClient client HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(10)) .build(); HttpRequest request HttpRequest.newBuilder(uri) .header(Authorization, Bearer apiKey) .header(Accept, text/event-stream) .header(Content-Type, application/json) .POST(BodyPublishers.ofString(payloadJson)) .build(); HttpResponseInputStream response client.send(request, HttpResponse.BodyHandlers.ofInputStream()); 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(5).trim(); handleChunk(payload); } } }这段代码在一切正常的时候能跑通但离生产可用差得很远。问题不是它短而是它把 SSE 当成了逐行 JSON来处理忽略了这个协议真正的事件边界语义。2.2 第一个坑事件边界不是行是空行SSE 协议里一个事件可以在data:后跨多行。比如服务端把一个长 JSON 拆成了两行data: {content:先输出一部分然后继 data: 续输出剩余部分}上面两行属于同一个事件客户端要把它们拼接起来才能得到完整的 JSON。如果只按行处理第一行会被当成一个不完整的 JSON解析直接抛异常。更隐蔽的是某些大模型服务在流式输出时单个事件的 JSON 字段可能被分片发送而代理层或服务端框架在写响应流的时候又不能保证一整个事件就写在一个 TCP 包里。我当时的修复方式是引入一个简单的状态机遇到data:行就追加到缓冲遇到空行且缓冲非空才把累积内容作为一个事件提交StringBuilder eventBuffer new StringBuilder(); String line; while ((line reader.readLine()) ! null) { if (line.isEmpty()) { if (eventBuffer.length() 0) { dispatchEvent(eventBuffer.toString()); eventBuffer.setLength(0); } } else if (line.startsWith(data:)) { eventBuffer.append(line.substring(5).trim()); } else if (line.startsWith(:)) { // 注释行心跳忽略 } }这段逻辑是后来所有封装的地基。很多人直接把readLine结果逐个丢给 JSON 解析器等到线上出现偶发解析失败才排查到这一层。2.3 第二个坑心跳注释行不能当数据处理没有业务数据的时候服务端为了保住连接会定期发一个注释行比如: ping这一行按协议属于注释客户端必须忽略。但是在前面那段按行处理的第一版代码里line.startsWith(data:)为 false它会被静默跳过——这没问题。问题出在一些封装得比较粗的代码里有人用line.contains(data)或者按字符串分割的方式取数据把: ping也当成数据去解析了轻则日志多一堆警告重则把心跳内容塞进业务消息队列。所以我在解析器里明确把三行类型分开处理data:行进缓冲event:/id:/retry:行记录元信息冒号注释行直接忽略空行提交事件。这个分类思路后来成了封装版本里SseEventListener的基础。2.4 第三个坑UTF-8 字符被网络包切开这是显式解析阶段最难排查的一个问题。流式响应是Transfer-Encoding: chunked的网络层把字节流分成任意大小的块。一个完整的中文字符在 UTF-8 里占 3 个字节如果服务端把{content:你好}这个 JSON 的字节按 7 字节、7 字节这样切分恰好把一个你字的三个字节切到了两个 chunk 里。如果服务端那边是按行 flush那么readLine通常不会遇到这种问题因为 BufferedReader 会等读到换行符才返回。但如果用的不是 BufferedReader而是自己攒字节数组、直接new String(bytes)那就一定会出现乱码。在显式调用阶段我坚持用BufferedReaderInputStreamReader并显式指定 UTF-8把按字节切分的问题交给 JDK 的流式解码器去缓冲处理不去自己拼字节。2.5 第四个坑断开场景要分清不能一刀切重连SSE 连接断开至少有四种情况处理方式完全不同断开场景现象正确动作正常结束读取到一个标记事件如data: [DONE]随后服务端关闭流按正常流程收尾不重连HTTP 错误连接建立阶段返回 401/429/500响应体不是 event-stream记日志、按业务错误处理不盲目重试网络中断read()返回 -1或抛 IOException按退避策略重连带上 Last-Event-ID空闲超时长时间无数据代理层断开读方向表现为 EOF 或超时异常触发心跳重连并推动服务端加心跳我当时一度把所有IOException都当成网络抖动于是加上无脑重试。结果遇到了 HTTP 429限流客户端每 5 秒重试一次服务端的限流状态被持续加重最终把服务端打到雪崩。后来改成只有网络层异常IOException / EOF才重试HTTP 状态码异常直接上抛才算稳住。2.6 显式调用的真正痛点不是代码长而是每个业务方都有一份自己的实现显式解析跑通之后我发现团队里每个需要接流式接口的人都各自写了一份有人抄第一版有人随手用了Scanner有人把重试写成了 while(true)。代码的语义不统一出了问题要互相翻代码对照。这时候才意识到SSE 的解析、重连、超时、事件分发这些横切逻辑必须收口成一个独立的模块也就是标题里说的隐式封装阶段。3. 隐式封装把 SSE 流式调用收敛成一行代码3.1 封装前先想清楚对外暴露什么抽象封装最怕一开始就堆 API。我在设计通用 SSE Client 之前先明确了一点业务方需要的是什么他们不关心text/event-stream不关心注释行甚至不关心重连——他们只想说我要调这个模型的流式接口给我一个能拿到 token 序列的东西。我把内部实现拆成两层底层SseStreamReader只干一件事——从 InputStream 里读行、拼事件、触发监听器。它不关心数据是来自 OpenAI、Claude 还是自研模型。上层StreamingChatClient负责拼接请求、注入认证头、调用底层读取器并把事件按业务类型反序列化为ChatChunk这样的 Java 对象。对调用方暴露的接口设计成了回调式sseClient.chatCompletion(request) .onMessage(chunk - { // 业务侧只处理 ChatChunk底层解析细节全部隐去 buffer.append(chunk.getContent()); emit(buffer.toString()); }) .onError(e - log.error(stream error, e)) .onComplete(() - log.info(stream finished));为什么选回调而不是返回FluxChatChunk因为团队技术栈还是 Spring MVC 为主没有全面引入 WebFlux强行上 Reactor 会让没有响应式经验的人写出更难维护的代码。回调语义直白排查问题成本低。如果团队已经在用 WebFlux那封装成FluxString是更好的选择这一点取决于上下文没有绝对的优劣。3.2 解析器与业务解耦才方便单元测试封装版本里最重要的一个设计是把解析单独抽成了纯逻辑组件public class SseEventParser { private StringBuilder dataBuffer new StringBuilder(); private String eventName message; private String eventId; private long retryMs; public void feedLine(String line) { if (line.isEmpty()) { if (dataBuffer.length() 0 || eventName ! null) { SseEvent event new SseEvent(eventId, eventName, dataBuffer.toString()); dispatch(event); reset(); } return; } if (line.startsWith(:)) { return; // 注释行 } int colonIndex line.indexOf(:); String field colonIndex 0 ? line.substring(0, colonIndex) : ; String value colonIndex 0 ? line.substring(colonIndex 1).trim() : ; switch (field) { case data - { if (dataBuffer.length() 0) dataBuffer.append(\n); dataBuffer.append(value); } case event - eventName value; case id - eventId value; case retry - retryMs Long.parseLong(value); } } }这个类的价值在于它是纯 Java 对象不依赖网络、不依赖 HttpClient可以喂字符串做单测。我后来的测试用例全是拿模拟的 SSE 文本串跑测试比如String sample id: 1 data: {token:你} id: 2 data: {token:好} ;这种设计看起来简单但它把最容易出错的部分从 IO 线程里剥离出来让解析逻辑可以被反复验证。很多封装之所以越改越乱就是因为解析和 IO 耦合在一起无法独立测试。3.3 泛型反序列化的边界消息可能是错误大模型服务商的流式接口正常情况下每个data:都是一个 token 块。但服务端压力大或参数非法时可能在一个事件里返回错误信息——同一个事件流里混着正常数据和不正常的错误结构。封装时如果只做一种ChatChunk类型错误事件会被反序列化失败然后走onError但原始的错误 JSON 已经丢了。我的方案是添加一个原始事件通道public record SseEvent(String id, String event, String rawData) {}解析器先抛SseEvent只带 raw JSON 字符串上层再决定如何反序列化。如果反序列化为ChatChunk失败保留 rawData 记录到日志同时发起告警。这样排查问题时能看到服务端到底回了什么而不是只看到一句JsonParseException。3.4 重连策略指数退避加抖动且要区分能不能重试重连是封装里最需要工程化处理的部分。我最终采用的做法private void connectWithRetry() { int attempt 0; while (true) { long delay baseDelayMillis * (1L Math.min(attempt, 6)); delay delay ThreadLocalRandom.current().nextLong(0, delay / 5); // 加抖动 try { Thread.sleep(delay); openStream(); // 阻塞直到流正常结束 return; } catch (InterruptedIOException | SocketException e) { // 网络异常继续重试 attempt; } catch (HttpErrorStatusException e) { // 4xx/5xx 不重试直接上抛 throw e; } } }几个关键设计重试只针对网络层异常。HTTP 状态码异常直接抛出去即使要重试也由上层决定而不是底层无限循环。重连时必须带上Last-Event-ID否则服务端无法知道你之前收到哪条事件重连后可能丢数据。抖动jitter必须有。同一个服务端故障时所有客户端如果按同样的退避节奏重连会产生周期性尖峰流量加抖动可以打散这个尖峰。3.5 封装阶段的不过度设计原则我见过有人把 SSE 封装做成一个支持多路复用、动态路由、断线持久化的重型框架。说实话在大多数 AI 接入场景里这些东西 90% 用不上。我最后保留的能力就四样事件解析、事件分发、自动重连、通用请求注入。其他全部留给上层用组合的方式解决。这一阶段解决了解析和重连的重复劳动问题但有一个核心矛盾没有解决每个流式连接在等待数据时仍然阻塞着一个宝贵的平台线程。这个矛盾在并发量上去之后变成了新的瓶颈也就走到了第三个阶段——虚拟线程。4. 虚拟线程介入把 SSE 长连接的性能天花板捅破4.1 传统线程池为什么扛不住大量 SSE 长连接SSE 连接有一个特点大部分时间连接是空闲的。服务端生成 token 需要时间tokens 之间的间隔从几十毫秒到几秒不等客户端大部分线程阻塞在InputStream.read()上干等。一个线程池如果只有 200 个线程那么最多只能同时维持 200 个流式连接在等待。一旦达到上限后面新的连接请求会进入线程池的排队队列连接建立并发出 POST 后线程已经被前面的长连接占死用户表现为请求迟迟不返回。有人提出把等待线程换成连接池复用但 Java 里 InputStream 的 read 是同步阻塞的线程池不解决一条连接一个等待线程的问题除非引入 NIO 和事件循环重写整个客户端。这就引出了虚拟线程——它正是为了打破一个阻塞 IO 占一个线程的困局而生的。4.2 虚拟线程的调度原理为什么阻塞 IO 不再是灾难虚拟线程Virtual Threads是 JDK 21 正式引入的轻量级线程。平台线程Platform Thread直接由操作系统调度一个平台线程对应一个内核线程而虚拟线程是由 JVM 管理和调度的它不直接占用一个内核线程。一个平台线程可以被当作载体线程carrier在上面跑多个虚拟线程。当虚拟线程执行到阻塞 IO比如 socket read时JVM 的调度器会自动把这个虚拟线程挂起释放载体线程去执行其他虚拟线程。也就是说阻塞不再意味着一个 OS 线程被占住而是变成了一次便宜的上下文切换。我用一个出租车类比来帮助团队理解平台线程是一辆出租车一次只能载一位乘客乘客不上下车这辆车就不能去接别人虚拟线程是乘客遇到路口红灯阻塞 IO时调度器让乘客先下车出租车立刻去接另一位乘客。红灯结束乘客在原地继续上车前进。虽然每位乘客的到达时间受调度影响但整体运力大幅提升。使用方式极其简单try (var executor Executors.newVirtualThreadPerTaskExecutor()) { executor.submit(() - { // 这里是阻塞读平台线程版本会占死一个线程 sseStreamReader.read(); }); }newVirtualThreadPerTaskExecutor会为每个任务新建一个虚拟线程任务结束后自动回收。在 SSE 场景里每个长连接独立放到一个虚拟线程里等待读流连接数上限从线程池大小变成了JVM 堆内存能容纳多少虚拟线程——实测单机几万连接是常态。4.3 同一个网关上平台线程池与虚拟线程的实测对比我在自己维护的流式网关服务上做了一组压测。网关的作用是接收前端浏览器的 SSE 请求再以 SSE 客户端方式向上游大模型服务拉取 token 流中间做鉴权、限流和转发。这个服务是最典型的大量长连接等待场景。对比两套实现平台线程Executors.newFixedThreadPool(200)每个请求分配一个线程执行阻塞读。虚拟线程Executors.newVirtualThreadPerTaskExecutor()每个请求分配一个虚拟线程执行阻塞读。压测条件模拟 3000 个客户端同时连接每个连接持续 60 秒期间服务端每 500ms 推送一个 token。指标平台线程200 线程虚拟线程最大并发连接约 200超出的排队3000 无压力客户端 P95 首 token 延迟8.6s排队导致480msCPU 占用3.2 核4.1 核线程栈内存占用200 x 1MB ≈ 200MB3000 x 16KB ≈ 48MB故障表现RejectedExecutionException 频发无拒绝异常平台线程的两百个连接很快就打满了后续连接全部在队列里排队任务实际开始时间被无限拉长虚拟线程版本在 3000 连接下没有出现连接拒绝CPU 也没有明显飙升。虚拟线程的栈默认是懒加载的不像平台线程那样初始就分配 1MB所以线程数量上来之后内存反而可控。这个对比充分说明一个结论在 IO 密集、连接长期处于等待态的 SSE 场景里虚拟线程是比池化平台线程更合适的调度模型。它不是让你少写代码而是让你可以继续保持一个连接一个任务线程这种最简单的编程模型却不再受制于线程数量。4.4 虚拟线程 SSE 时容易翻车的三个细节虚拟线程不是万能药在 SSE 场景里我踩了几个比较隐蔽的问题第一个是 synchronized 锁的 Pin 问题。JDK 21 早期版本里如果虚拟线程执行到了synchronized块并且该锁被其他线程持有虚拟线程会钉住pin在载体线程上无法让出载体线程导致载体线程被白白占住。解决方法是优先使用ReentrantLock等显式锁或者在涉及锁的代码里尽量缩短持有时间。第二个是 ThreadLocal 的内存放大。每个虚拟线程都是一段独立的调用栈如果代码里大量往 ThreadLocal 里塞大对象那几万连接就会产生几万份副本。我在网关里用 ThreadLocal 存储请求上下文压测时发现 GC 压力明显上升。换成从参数里显式传递上下文之后问题才缓解。第三个是第三方库的兼容性。有些老牌 HTTP 客户端、数据库驱动还在用自己的线程池虚拟线程只在你自己写的阻塞代码上有收益。如果连接建立和读取都发生在 Spring 的 WebClient基于 Reactor Netty 的线程模型里那虚拟线程的好处就体现不出来。所以迁移之前一定要先确认连接链路里最底层的阻塞等待发生在哪里。5. 生产环境 SSE 故障排查从报错到根因的完整链路5.1 经典报错stream disconnected before completion: idle timeout waiting for sse这是我在网上平台和团队群里都多次见过的报错src 场景是 AI 流式响应到一个固定时间点比如正好 1 分钟或 2 分钟就断了下游拿到的是不完整的输出。这个报错本质上就是四个字空闲超时。排查它别急着改代码先回答一个问题这个超时是谁发起的我按下面的链路逐步排查看客户端日志确认异常抛出的位置。如果来自 HTTP 底层解析器通常是代理网关先断了如果来自业务代码超时那是我们自己设了 socket read timeout。看中间的代理层如果请求经过 Nginx重点看proxy_read_timeout它的默认值是 60 秒而大模型生成第一个 token 可能就要十几秒后续 token 间隔也可能超过 60 秒。这个参数不调SSE 长连接一定会被 Nginx 掐断。看服务端有没有心跳SSE 协议里即使没有业务数据服务端也应当周期性发送注释行: keepalive加空行来维持连接。很多 AI 服务商的服务端没有做这一步就只能靠代理层加大超时来兜底。客户端侧的 read timeout 检查有些 JDK HttpClient 配置了request.timeout或者底层 socket 的 soTimeout这也会导致等待下一行超时。在 SSE 场景下读超时应该设置得很大或者干脆不设置由重连机制兜底。最终我给出的标准配置是Nginx 的proxy_read_timeout调到300s服务端每 15 秒发一个: ping心跳客户端不设读超时。三条同时做这个报错就从线上绝迹了。5.2 服务端重启导致的重连风暴SSE 连接断开的瞬间所有客户端几乎同时发现连接断了如果它们的重连策略里只有固定间隔比如 3 秒那么 3 秒后所有客户端会同时发起重连服务端刚刚起来就被流量打懵。我处理这个问题的两个措施客户端重连必须加指数退避和随机抖动上一节封装里已经做了服务端发布时先主动关闭新连接让健康检查探针先通过再恢复服务。这样能错开客户端的重连时间。重连风暴典型的日志特征是服务端所有 worker 线程在同一个时间点全部 busy客户端日志里全是 connection refused 和 retrying 交替出现。5.3 连接池被 SSE 长连接占满Java 的HttpClient默认带连接池同一个目标域名复用一个连接。SSE 长连接的特点是一次连接长期占用如果服务本身同时在调用普通 REST 接口普通请求就可能拿不到复用连接不断建立新连接直到目标主机的连接数打满。对这种问题最简单有效的手段是把 SSE 调用和普通调用拆成两个独立的HttpClient实例SSE 实例的Executor使用虚拟线程的同时给它一个独立的连接池上限普通 REST 实例保持默认配置。不要让长连接和短连接抢同一个池子。5.4 文件描述符泄漏虚拟线程版本下更隐蔽平台线程时代连接泄漏大概率表现为线程数飙升jstack 一眼看到。虚拟线程时代线程数不再是一个敏感指标连接泄漏更容易表现为文件描述符持续增长直到触发器触发too many open files。排查方法# 查看进程打开了多少文件描述符 lsof -p pid | wc -l # 持续观察 10 秒 for i in {1..10}; do lsof -p pid | wc -l; sleep 1; done如果数量单调递增基本就是SseEventStream没有在onComplete或onError时关闭。虚拟线程案例里线程会被快速回收但如果底层 socket 没有 closefd 依然在增长。所以封装里try-with-resources不只是语法习惯它在长连接场景是保护伞。真实体会三个阶段的取舍不是层层替代如果你问我这三个阶段是不是必须按顺序经历我的答案是不用。我现在做新项目会直接采用解析器 回调式封装 虚拟线程的组合起步没有人需要先手写一遍痛苦的第一版才能理解封装的价值。但如果你接手的是老代码从显式调用往封装演进时先守住解析器这个纯逻辑核心再谈重连和虚拟线程这个顺序千万别反。另外给正在做这个方向的同学一个建议别把 SSE 客户端写成只适配某一家的协议数据格式会有差异、错误码会有差异、结束标记也有差异有的用data: [DONE]有的用event: done。在封装层留好适配接口上层按服务商注入各自的解析规则后面接第二家模型时能省一大半功夫。我自己的体会是Java 处理 SSE 的关键从来不是某个神秘的高深技巧而是把协议细节理解透之后用最朴素的方式把边界条件守住。虚拟线程把性能问题解决掉之后剩下的就是工程问题解析、重试、心跳、连接生命周期一个一个处理好这套系统就能稳定跑很久。
上一篇/下一篇内容由系统自动关联
返回资讯列表 →