尧图精选

Java大模型SSE流式输出实战:从手写解析到虚拟线程优化

🕒 发布时间:2026/10/1 9:44:33 📁 来源:尧图网络
1. 从一次大模型联调说起SSE为什么成了Java AI开发的“隐形地基”如果你最近半年写过Java后端对接大模型的服务大概率绕不开一个词SSEServer-Sent Events。别管你接的是国内的开源模型、商用API还是自己部署的推理服务只要产品里有“打字机效果”——也就是AI回复一个字一个字往外蹦的那种体验底层几乎都是SSE在干活。SSE不是新东西2011年前后就进了HTML5标准但过去十年它在Java后端圈子里的存在感一直不强。原因很简单传统企业级接口追求“请求-响应”一次到位没人愿意把响应拆成一串碎片再拼回去。直到大模型时代来了情况彻底反转——一个完整的大模型回答可能需要十几秒甚至几十秒如果等全量返回再展示用户早就流失了。流式输出成了刚需SSE这个老技术一夜之间被推到聚光灯下。但真正让Java开发者头疼的不是SSE本身而是它在Java生态里的“水土不服”。Java原生的HTTP客户端本来就对响应式协议支持一般Servlet 3.1之前的容器甚至没法很好地处理异步输出。再加上大模型接口的各种“小脾气”——超时、断连、半截消息、数据格式不标准——一套能用、能扛、能上生产的SSE封装远没有想象中那么简单。这篇文章就是围绕这条主线展开的从最早一批Java工程师手写SSE客户端开始到后来封装通用流式解析层再到JDK 21虚拟线程带来的性能质变。我会把每一阶段的代码逻辑、设计取舍、踩坑记录都摊开讲希望能帮你少走几段弯路。2. SSE原理再梳理Java大模型场景下你必须吃透的几个细节2.1 SSE的协议格式没那么神秘SSE本质上是HTTP协议上的一层文本流协议。服务端设置Content-Type: text/event-stream然后持续向客户端推送形如这样的数据块data: {id:1,content:你} data: {id:1,content:好}每两条消息之间用空行分隔每条消息以data:开头以\n结尾。很多大模型接口还会带event:字段来区分消息类型——比如有的模型用event: message表示正常内容event: done表示结束。如果只是按data字段一刀切解析遇到带event的接口就会出问题。有些刚接触SSE的同事会把它和WebSocket搞混。简单区分WebSocket是双向全双工SSE是单向服务端推送WebSocket走的是独立的ws/wss协议SSE还是标准的HTTP连接。在AI场景里客户端只需要“接收”服务端的生成结果几乎不会反向发数据所以SSE的“单向”反而成了优点——没必要为了一个单向流引入WebSocket的复杂度。2.2 流的边界消息拆分的核心难点SSE的难点不在读数据而在“切数据”。HTTP底层是字节流TCP会做分包和粘包你拿到的InputStream不会恰好一次给一条完整消息。有可能一次读到了三段半消息也有可能读了四次才凑出一条完整SSE消息。这个“半条”的处理逻辑如果写不好轻则丢内容重则JSON解析直接崩溃。我见过不少从Python转过来的同事写Java SSE解析喜欢用BufferedReader.readLine()按行读。这个思路本身没错但readLine()是阻塞的而且如果服务端长时间没有新数据比如模型在“思考”这个线程就会一直挂在IO上。单连接、单线程测试没问题一旦并发上来线程池分分钟被打满。2.3 心跳与超时一辈子踩不完的坑SSE还有一个隐性问题连接空闲。大模型接口有时候会在“思考”阶段停顿几十秒——不是断连就是单纯没产出内容。这时候如果中间有负载均衡器、网关或者TCP层做空闲超时连接就会被静默掐断。客户端视角就是消息流到一半突然抛了IOException: unexpected end of stream或者像Spring WebClient里常见的报错——stream disconnected before completion: idle timeout waiting for sse。这个问题在标题里那串热搜词里也出现了确实太典型了。后面我会专门用一节讲这个问题的排查和修复方案这里先记住一个结论处理SSE超时策略必须和普通HTTP调用区别对待不能用一套老配置走天下。3. 显式调用阶段手写SSE客户端的那段“野蛮生长”岁月3.1 第一版代码用JDK原生API硬啃SSE在最早期的大模型对接项目里多数Java团队都没有现成的SSE库或者不敢直接用毕竟接口太新型担心不稳定。最朴素的方案就是用HttpURLConnection或HttpClient拿到输入流然后手动解析SSE格式。我摘一段当年项目里真实用过的代码做了简化这是典型的显式调用实现public void testSSE() { HttpClient client HttpClient.newHttpClient(); HttpRequest request HttpRequest.newBuilder() .uri(URI.create(http://localhost:8080/llm/chat)) .header(Content-Type, application/json) .timeout(Duration.ofSeconds(30)) .POST(BodyPublishers.ofString({\prompt\:\讲个笑话\})) .build(); client.sendAsync(request, HttpResponse.BodyHandlers.ofInputStream()) .thenAccept(response - { try (BufferedReader reader new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; StringBuilder eventData new StringBuilder(); while ((line reader.readLine()) ! null) { if (line.isEmpty()) { // 空行代表一条消息结束 String message eventData.toString(); processMessage(message); eventData.setLength(0); } else if (line.startsWith(data:)) { eventData.append(line.substring(5).trim()); } else if (line.startsWith(event:)) { // 记录事件类型 currentEvent line.substring(6).trim(); } } } catch (IOException e) { log.error(SSE stream error, e); } }).join(); }这段代码可以跑通一个最简单的SSE接口但问题也很明显单条消息如果很大跨多行的data拼接容易丢失换行符readLine()阻塞导致线程占用内存大HttpClient默认的io线程是ForkJoinPool阻塞会拖累其他非阻塞任务而且没有任何心跳处理连接空闲长了直接被掐。3.2 显式调用带来的三个血泪教训显式调用阶段我最大的体会是三个字不抗造。第一个教训缓存区没有边界控制。大模型偶尔会输出一个很长的JSON比如带一大堆历史上下文回显。如果StringBuilder无限拼接一个异常响应就能撑爆内存。后来我们加了MAX_MESSAGE_LENGTH限制超过阈值强制断连避免被异常数据拖死。第二个教训错误处理太粗暴。上面那段代码catch到异常就只是打日志但实际生产里SSE流断掉之后是需要决定“这条消息怎么收尾”的——是重跑一次全量请求还是把已收到的部分兜底存掉这个决策逻辑没有的话用户看到的就是半截回答停在页面中间。第三个教训线程模型落后。HttpClient的sendAsync底层用ForkJoinPool做回调你用阻塞式readLine()就相当于把一个异步线程池卡成了同步线程池。并发20个SSE请求ForkJoinPool的线程数默认等于CPU核数8核机器8个线程一个模型思考卡住20秒剩下12个请求全部排队。更气人的是ForkJoinPool的线程从代码里还不好单独调。显式调用不是不行但它要求每个接入方都自己处理解析、失败、调度、超时这一整套逻辑。项目一多同样的代码复制粘贴四五个版本每个版本还都留着各自的bug——这时候“封装”就成了顺理成章的需求。4. 隐式封装提炼SSE流式解析层的设计思路与核心实践4.1 封装的目标让业务代码不感知“流”做封装之前我先把需求拆了一遍。团队里有四五条业务线都在调大模型接口有的是RAG客服问答有的是代码生成助手还有的是文档摘要。每条业务线的请求参数、提示词模板、后处理逻辑都不一样但它们有一个共同点都需要调用SSE接口都需要逐字符获取AI输出都需要处理“中途断了怎么办”的问题。封装的第一个目标因此很明确——把SSE的读流、解析、心跳、超时、断连重试全部收进一个通用层业务侧只需要传入一个回调函数像订阅消息一样消费“增量内容”和“结束事件”。业务开发不必知道SSE协议长什么样更不用在代码里看到任何一个data:前缀。第二个目标是——保留灵活逃逸口。封装不是为了锁死总有一些场景需要拿到原始事件对象做特殊处理所以我会在回调里暴露SseEvent完整对象而不仅仅是一个干巴巴的字符串。4.2 通用SSE解析器的核心实现下面这段代码是封装层的骨架也是我后来在新项目里反复复用的一套东西。我先给核心部件public class SseClient { private final HttpClient httpClient; private final ExecutorService executor; public SseClient(ExecutorService executor) { this.executor executor; this.httpClient HttpClient.newBuilder() .executor(executor) .connectTimeout(Duration.ofSeconds(10)) .build(); } public void stream(SseRequest request, SseListener listener) { HttpRequest httpRequest HttpRequest.newBuilder() .uri(URI.create(request.getUrl())) .timeout(Duration.ofSeconds(request.getWriteTimeout())) .header(Accept, text/event-stream) .header(Authorization, request.getAuthToken()) .POST(BodyPublishers.ofString(request.getPayload())) .build(); httpClient.sendAsync(httpRequest, HttpResponse.BodyHandlers.ofInputStream()) .thenAccept(response - { if (response.statusCode() ! 200) { listener.onError(new SseException(HTTP response.statusCode())); return; } parseStream(response.body(), listener); }) .exceptionally(ex - { listener.onError(ex); return null; }); } private void parseStream(InputStream body, SseListener listener) { try (BufferedReader reader new BufferedReader( new InputStreamReader(body, StandardCharsets.UTF_8), 8 * 1024)) { String line; SseEventBuilder builder new SseEventBuilder(); while ((line reader.readLine()) ! null) { if (line.isEmpty()) { SseEvent event builder.build(); listener.onEvent(event); if (SseEventType.DONE.equals(event.getEvent())) { listener.onComplete(); return; } builder.reset(); } else if (line.startsWith(data:)) { builder.appendData(stripLeadingSpace(line.substring(5))); } else if (line.startsWith(event:)) { builder.setEvent(line.substring(6).trim()); } } listener.onComplete(); } catch (IOException e) { // 这里并不仅仅是捕获异常而是要把半截消息兜住 listener.onError(e); } } }看到new BufferedReader(... , 8 * 1024)这个参数了吗这是缓冲区大小默认是8K实际上大模型一条SSE消息很少超过这个尺寸够用。但更关键的是我在SseEventBuilder里做的处理——它不只是简单拼接字符串还做了UTF-8断点处理。当一条消息被TCP拆成两个半片时第二个半片可能以开头这时候用String直接拼接会乱码。builder内部用ByteArrayOutputStream攒字节到最后统一转String确保没有半个码点的问题。4.3 回调接口设计增量、结束、错误各走各的通道回调接口我总结了三段式的设计public interface SseListener { // 每收到一条完整消息就触发一次 void onEvent(SseEvent event); // 流正常结束包括服务端发了done事件或者EOF void onComplete(); // 异常结束带原因 void onError(Throwable cause); }为什么必须把onComplete和onError分开因为业务上两者的处理逻辑可能完全不一样。正常结束前端可以把光标从“生成中”切成“已完成”异常结束前端可能需要提示重试还可能要把已生成的部分内容存进草稿。我见过有人把这两种情况合并处理最终用户经常看到“AI回答到一半悄无声息消失”的情况——那就是因为没有区分正常终止和异常断开。还有一点经验onEvent里的SseEvent对象除了data、event字段之外我还会带一个sequence序号。这个序号不是协议里的是我在解析器里自己递增的——它的作用有两个一是方便排查消息顺序有没有错乱比如某些网关重发了数据二是业务侧做“跳过前面N个历史消息”的增量处理时可以直接替换不必再解析data里的游标字段。4.4 隐式封装带来的开发体验变化封装层落地后的效果可以从两个维度看。站在业务同学的角度接入一个SSE接口从原来至少写200行“管道”代码压缩到这样几行sseClient.stream( SseRequest.builder() .url(llmService.getChatUrl()) .authToken(llmService.getToken()) .payload(buildPayload(userMessage)) .build(), new SseListenerAdapter() { Override public void onEvent(SseEvent event) { String delta event.getData(); chatSession.appendDelta(delta); } } );SseListenerAdapter是适配器模式的小技巧——让调用方只需覆写自己想处理的方法不想管的onComplete、onError用默认空实现兜着代码就更简洁了。站在系统层的角度解析、心跳、超时的改动不再需要动业务代码。有一次网关调整了空闲超时策略导致SSE连接频繁被断我们只是修改了封装层里的心跳发送逻辑给所有SSE连接统一加了20秒的心跳保活业务侧零改动就修复了线上问题。这就是“隐式封装”真正的价值——你可以把底层治理能力集中起来而不是散落在每个业务模块各自为政。5. 虚拟线程落地SSE场景的性能实测与架构改造5.1 SSE为什么是虚拟线程的理想试验场JDK 21正式发布了虚拟线程Virtual Threads这个特性在Java社区引起了不小的讨论。有人觉得是“响应式编程的救星”有人觉得只是换了个线程池名字。但如果你做的是SSE这种线程密集阻塞型IO的服务虚拟线程几乎是量身定做的。先理解一下传统方案在SSE高并发场景下的困境。假设你的服务作为网关同时对接上游大模型向下游WebSocket推送内容。每接一条SSE连接你需要一个线程专门负责读流——这个线程大部分时间都在readLine()上阻塞等待。用传统的平台线程Platform Thread每个线程占1MB左右的栈空间线程切换成本也不低。一台8核16GB的机器保守估计只能支撑几百到一千个并发SSE连接多数线程还都“无所事事”地等着数据。虚拟线程把这个问题彻底解掉了。虚拟线程是JVM内部实现的轻量级调度单位一个平台线程可以挂载成千上万个虚拟线程。当虚拟线程执行到阻塞IO时JVM会自动把它从平台线程上“卸载”掉把平台线程让给其他可运行的虚拟线程。这种机制下“一个连接一个线程”的编程模型重新变得可行——哪怕一万个SSE连接也只需要几十个平台线程就能扛住。5.2 把SseClient重构到虚拟线程我实际做的改造其实就是把SseClient里的executor从传统的Executors.newFixedThreadPool(50)换成了虚拟线程。改起来非常简单ExecutorService executor Executors.newVirtualThreadPerTaskExecutor(); SseClient sseClient new SseClient(executor);不止这些还有两处需要配合调整。第一HttpClient连接器必须复用。HttpClient内部有自己的连接池虚拟线程环境下连接池大小需要调大否则虚拟线程再多连接池不够也是白搭。我把HttpClient.Builder的connectionPoolSize调到了5000具体要看业务并发量别拍脑袋。第二sendAsync的调用链里尽量不要用传统CompletableFuture的thenApplyAsync。虚拟线程的核心理念是“用普通阻塞代码代替异步链”所以我把回调里那套CompletableFuture链路简化成同步代码让每个虚拟线程从头到尾串行处理一个SSE流。代码可读性反而提升了——没有乱七八糟的thenCompose嵌套了。5.3 压测数据从线程枯竭到指数级扩容改造完成后我做了一轮对比压测。场景是模拟2000个客户端同时请求AI对话每个SSE流持续30秒每秒钟产出一条消息。压测机器的配置是8核16GBJDK 21。使用传统线程池固定50线程时大约在400个并发连接左右线程池就接近饱和RT开始飙升大量请求排队部分请求直接超时。改用虚拟线程后2000个并发SSE流稳定运行RT的P95从原来的4.2秒降到1.8秒——这1.8秒还不全是线程调度开销而是上游模型推理本身的时间。线程池那组数据没法继续测更高的并发因为固定50个线程意味着最多50个连接同时被处理其余全部排队。虚拟线程那组我继续往上压到5000并发依旧稳定痛点反而转移到了机器负载和上游接口的限流策略上。这个结果其实在意料之中。SSE场景里线程的“拥有者”数量决定了并发上限——平台线程时代物理资源决定了你有多少个“拥有者”虚拟线程时代只要你愿意可以给每个连接配一个专属虚拟线程连接数就不再是瓶颈了。5.4 虚拟线程改造的三个提醒改造的甜头虽然大但也有几个坎要过。提醒一synchronized块别滥用。JDK实现虚拟线程时已经优化了synchronized在JDK 21里如果-XX:VMContinuations开启遇到ReentrantLock会避免pin住载体线程但如果你用了一些老库里面有大量synchronized虚拟线程阻塞时可能还是会把底层平台线程“钉住”导致并发上不去。压测时如果发现虚拟线程表现异常先排查有没有ClassLoader锁、Netty内部锁之类的老代码。提醒二ThreadLocal要慎用或换成ScopedValue。传统代码用ThreadLocal存用户上下文挺常见的但虚拟线程的数量是海量的如果每个虚拟线程都往ThreadLocal里塞大对象内存会被吃掉很多。JDK 21提供了ScopedValue作为替代方案但API还在孵化期多数项目可能没有迁移动力。折中做法是把ThreadLocal里的对象改小别放全量上下文只放一个ID需要时再查。提醒三JDK版本和Spring Boot版本的兼容性。如果你项目还在JDK 17虚拟线程的API是Preview默认是关闭的。Spring Boot 3.2开始才完整支持虚拟线程。线上压测和生产切换要分两步走——先在测试环境用-Dspring.threading.virtual.enabledtrue开关试跑再平滑切流。6. 常见问题与排查技巧这套体系跑了一年的实录6.1 那个著名的“idle timeout waiting for sse”排查实录里这绝对是我遇到频次最高的问题。Spring WebClient下报错信息长这样stream disconnected before completion: idle timeout waiting for sse翻译成人话就是WebClient在等SSE数据但等到超时了。这个超时不是HTTP连接超时也不是响应超时而是读空闲超时——Spring的JettyClientHttpConnector或Reactor Netty的通道层在“规定时间内没有读到任何字节”时自动断开连接。大模型场景下为什么会触发最常见的情况是用户发了一个复杂问题模型开始“思考”前10秒、15秒、20秒没有产生任何文本输出。如果中间网关的读空闲超时设置在15秒这条连接在模型“开口说话”之前就被掐断了。解决办法有两层第一层把超时调大。如果是Reactor Netty在WebClient构建时这样配置ConnectionProvider provider ConnectionProvider.builder(custom) .maxConnections(500) .pendingAcquireTimeout(Duration.ofSeconds(60)) .maxIdleTime(Duration.ofSeconds(30)) .maxLifeTime(Duration.ofSeconds(30)) .build(); HttpClient httpClient HttpClient.create(provider) .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 5000) .doOnConnected(conn - conn.addHandlerLast(new ReadTimeoutHandler(60))) .protocol(HttpProtocol.HTTP11); WebClient client WebClient.builder() .clientConnector(new ReactorClientHttpConnector(httpClient)) .build();第二层让Server端主动发心跳。很多大模型服务端在空闲时会定期发送一个空注释: ping或空data:来维持连接活跃。但你不能假设所有服务端都会发。在客户端封装层里主动设一个读超时探测不行读超时只会报错不会保活。客户端能做的其实是针对一批已知不发送心跳的模型在业务层加一个“无结果超时承诺”——比如允许模型思考最多2分钟超过就主动截断并提示用户。这个问题的核心认知是“空闲超时等待SSE”不是网络故障是协议层面的保活机制缺失。修复它是服务端和客户端双方的事情单靠哪一边都治标不治本。6.2 消息粘连导致JSON解析报错这种问题多发生在没有严格按SSE规范实现的服务端。比如服务端把两条data消息用一个空行隔开但是忘了带\n后的\r或者在某些代理层被重新包装过。我遇到过最离奇的场景消息内容是data: {text:\u4f60\\n好}反序列化后丢了一个反斜杠导致JSON字符串里的转义符失效。排查半天才定位到是网关在处理\n时自作主张做了一次“规范化”把\\n变成了实际换行符。排查套路先用原始TCP抓包tcpdump或Wireshark看真正在线上传输的字节是什么确认格式无误后再到本地产一条相同载荷用socat或nc模拟服务端逐字节发送对比Java端处理结果定位是解析层还是传输层多做了手脚。6.3 连接池被SSE长连接占满这问题在封装SSE之前特别常见。使用HttpClient默认的HttpURLConnection方案时连接池默认5个连接——并发5条SSE流就直接堵死了。即使换到HttpClient默认的最大连接数也是受限的。解决思路很简单给SSE客户端单独配置一个连接池不跟普通请求混用同时给连接池设置一个“最大空闲时间”防止SSE连接断开后池里残留一堆半开连接。我在封装层里直接给SseClient单独初始化一个HttpClient实例业务侧普通HTTP调用用的是另外的HttpClient实例两者互不干扰。6.4 后台线程泄漏导致内存涨还有一类偶发问题SSE流结束后回调线程没有正确终止。在最早的显式调用版本里如果服务端既不发送done事件也不关闭连接这种情况在某些自建模型服务里真的存在readLine()会永远阻塞。如果每次请求都泄漏一个线程系统线程数就会缓慢但稳定地增长。封装后我加了一个“兜底定时器”每条SSE连接启用后启动一个守护线程计时90秒内没有新数据且没有结束就主动close()连接强制退出阻塞。这个“最后防线”上线后再没出现过线程泄漏的告警。问题类型直接表现根因处理手段空闲超时断开流中断、报idle timeout网关/服务端空闲保活缺失调大读空闲超时服务端心跳消息粘连JSON解析失败、内容错位TCP粘包或代理改写按data:空行规范解析抓包对比连接池占满请求排队、RT飙升默认连接池太小独立连接池配置maxConnections线程泄漏线程数持续增长、内存告警服务端不关流/客户端不超时兜底空闲定时器强制关闭6.5 一个小工具SSE调试器最后分享一个日常调试的实用技巧。排查SSE问题时用第三方工具比自己写代码快得多。我常用的是在本地跑一个小脚本模拟服务端输出#!/bin/bash # 模拟SSE服务端每秒输出一条消息 for i in $(seq 1 20); do echo data: {\seq\:$i,\content\:\hello\} echo sleep 0.5 done然后用客户端的封装层去连这个本地服务复现解析问题。这个脚本能帮你把“网络传输问题”和“客户端解析问题”快速隔离开省去跟网络组反复扯皮的工夫。7. 实践总结从显式到隐式我踩过的坑希望你绕过去回看这段从手写SSE到虚拟线程性能改造的历程最深的体会不是某一个技术细节多厉害而是整个过程中的决策逻辑。显式调用暴露了问题推动了封装封装稳定后又遇到性能瓶颈推动了线程模型升级。每走一步都是前面问题的自然延伸。如果让我给正在做同类项目的团队三个建议我会选这三个第一SSE客户端不要只做“能用”要考虑“能抗”。一个能跑通的解析器打包上线和能在网关抖动、模型卡顿、代理改写、连接池枯竭下还能优雅降级的解析器中间隔着一整层工程能力的积累。封装层是把这些能力集中起来管理的关键节点值得认真投入设计。第二虚拟线程值得用但别指望一行配置拯救所有性能问题。SSE场景性能提升明显的核心原因是IO密集与阻塞天然契合虚拟线程模型。如果你的服务CPU密集型虚拟线程的收益会打折如果存在大量synchronized老代码虚拟线程还可能要调优很久才能发挥优势。第三日志和可观测性要从第一天就做。SSE流式接口的排查难度比普通接口高一个数量级——错误发生在“流的中间”而不是“请求的结束”。我在封装层里给每条SSE流打了一个traceId并把“开始时间、首包耗时、总耗时、结束原因正常/超时/异常、累计字符数”全打出来。这套日志后来帮我们在一次线上事故里十分钟定位到是上游某台机器间歇性抽风而不是我们自己代码的问题。如果你刚起步不需要一步到位做虚拟线程改造——先把SSE解析封装层做好跑通业务再逐步优化。把基础打牢了后续的任何性能飞跃都不是问题。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →