Java消费SSE流全实践:从HttpClient解析到虚拟线程并发

发布时间:2026/10/2 4:44:20
Java消费SSE流全实践:从HttpClient解析到虚拟线程并发 接大模型接口的时候我第一次在终端里看到一行行data:往外蹦说实话心里是有点激动的——一个普通 HTTP 请求突然变成了“活的”服务端每生成一个 token 就推一段过来。但对于我这种平时主要写 Java 后端、对 HTTP 细节不太较真的开发者来说激动之后紧接着就是问题Java 里到底怎么正确消费这种 SSEServer-Sent Events流这个问题的答案我从最初的裸写 HttpClient到后来封装成通用组件再到用虚拟线程解决并发瓶颈走了一个还挺完整的演进过程。这篇就完整记录这段实践给准备接 AI 流式接口、或者面试被问到 Java SSE 和虚拟线程的同学做个参考。1. SSE是个“老协议”但AI把它带回了聚光灯下1.1 从Comet、长轮询到EventSourceSSE其实早就存在SSE 全称 Server-Sent Events2010 年前后随着 HTML5 一起出现在规范里。它的目标很直接让服务器能把数据主动推给浏览器。那几年大家对“推”的执念很深因为 HTTP 请求-响应模型本质上是“拉”所以才有 Comet、长轮询这些土办法。SSE 的思路很朴素——你发一个普通 HTTP 请求服务器挂起响应想推数据了就往响应体里写一行写完不关连接。用生活化一点的类比普通 HTTP 是你去窗口问一句窗口给你一个完整答案然后关窗。长轮询是你守在窗口等窗口有答案了再给你但每次给完都关窗你只能再排队。SSE 则是窗口直接给你开了一条专属通道答案分了几个批次递出来通道一直留着后面有补充继续递。这个协议本身不复杂但在过去十几年里一直不温不火因为大部分 Web 场景用轮询就能凑合。直到 AI 大模型把“逐字生成”变成刚需SSE 才重新被推到前台。如今几乎主流的大模型推理接口都默认支持text/event-stream返回你会发现这个“老古董”反而成了 AI 应用接入的事实标准。1.2 AI流式输出为什么偏偏选SSE而不是WebSocket我在设计技术方案时最先想到的是 WebSocket毕竟这个词在很长一段时间里就是“服务器推送”的代名词。但深入对比之后会发现AI 流式场景选 SSE 并不是拍脑袋而是权衡后的结果。方案通信方向协议成本断线续传Java后端实现复杂度短轮询客户端单向拉低天然支持低WebSocket双向全双工高需要升级协议需手写实现高SSE服务端单向推低仍是普通HTTP自带 Last-Event-ID低关键点在于大模型生成内容是一个严格的单向过程客户端提交一次请求参数服务端不断回推内容片段。整个过程中客户端几乎不需要通过同一个连接往上发消息如果需要取消直接把连接关掉就行。SSE 的“半双工单向下行”模型和这个场景刚好贴在一起。而且 SSE 走的是标准 HTTP不涉及协议升级意味着 Nginx、网关、负载均衡这些基础设施几乎不用做额外适配。WebSocket 则要处理握手、心跳、代理透传、连接状态恢复等问题链路越长越容易出幺蛾子。1.3 SSE数据格式快速扫盲到底在解析什么后面所有代码都建立在一个基础上SSE 的传输格式。它本质上是 UTF-8 文本流每一行是一个字段字段名和值之间用冒号分隔常见字段如下data:真正的数据载荷event:事件类型默认是messageid:事件 ID用于断线续传retry:重连间隔毫秒:开头的行是注释行一般用于心跳两个空行之间的所有内容算一个完整事件。一个典型的 AI 流式返回长这样event: message data: {role: assistant, content: 你} event: message data: {role: assistant, content: 好} : heartbeat data: [DONE]注意最后那个data: [DONE]很多大模型服务会在全部内容推完后额外发一个结束标记这和流式事件格式有关但不属于业务 JSON解析时要把这种控制消息单独处理否则在 JSON 反序列化阶段会直接报错。我当时没注意这个细节第一版代码在[DONE]上报了不下十次异常。2. 第一版实现裸写HttpClient从InputStream里逐行抠数据2.1 为什么不能等整个JSON返回必须用InputStream最开始我犯过一个典型错误用HttpClient的BodyHandlers.ofString()去请求大模型接口结果确实能拿到完整 JSON但整个响应在服务端全部生成完之前不会返回。用户看到的效果就是点了发送之后转圈圈等三五秒后“唰”地一下整段冒出来完全失去了流式输出的体验。问题的本质在于ofString()会把整个响应体缓存到内存里等连接关闭后才一次性给你。流式的基础是先拿到响应体边下载边消费。Java 11 引入的HttpClient提供了一个很关键的BodyHandlers.ofInputStream()它不等响应体结束而是在响应头到达后就返回一个输入流后续数据靠我们从流里逐行读取。这就相当于把“整包取货”改成了“传送带取货”。2.2 最小可运行的显式调用代码下面这个例子是纯 Java 标准库实现不依赖任何第三方框架也是我第一版代码的简化版HttpClient client HttpClient.newHttpClient(); HttpRequest request HttpRequest.newBuilder() .uri(URI.create(https://gateway.example.com/v1/chat/completions)) .header(Content-Type, application/json) .header(Authorization, Bearer sk-xxxx) .header(Accept, text/event-stream) .POST(HttpRequest.BodyPublishers.ofString(payload)) .timeout(Duration.ofSeconds(30)) .build(); HttpResponseInputStream response client.send( request, HttpResponse.BodyHandlers.ofInputStream()); if (response.statusCode() ! 200) { // 错误响应的 body 不是 SSE 格式用普通方式读取 String err new String( response.body().readAllBytes(), StandardCharsets.UTF_8); throw new RuntimeException(请求失败: response.statusCode() - err); } try (BufferedReader reader new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; String eventType message; StringBuilder data new StringBuilder(); while ((line reader.readLine()) ! null) { if (line.isEmpty()) { // 空行代表一个事件结束这里把积累的 data 交给业务处理 handleEvent(eventType, data.toString()); eventType message; data.setLength(0); continue; } if (line.startsWith(event:)) { eventType line.substring(6).trim(); } else if (line.startsWith(data:)) { // 注意如果 data 后面直接跟内容要去掉一个前导空格 data.append(line.substring(5).trim()); } else if (line.startsWith(:)) { // 注释行一般是心跳忽略即可 } } }这段代码跑通之后的效果是大模型每生成一个 tokenhandleEvent就会被调用一次前端就能看到逐字输出的效果。这里有一个特别容易忽略的细节data:行官方规范允许同一个事件里出现多次而且多行 data 应该用换行符拼接。但主流大模型接口为了便于 JSON 解析基本都发单行 data所以我第一版直接用StringBuilder追加。如果你的上游是某个私有化部署模型建议还是先抓包看一眼确认它是单行 data 还是多行 data再决定拼接策略这个坑我后面还会提到。2.3 显式调用跑通之后代码很快变得没法维护满足感没有持续太久。第二周我把这个调用方式复制到第二个业务场景时问题集中爆发了。首个问题是解析逻辑和业务逻辑彻底纠缠在一起。每个需要调用大模型的服务都要复制一份这二三十行的循环解析代码里面既有 HTTP 状态码判断又有 SSE 事件拼装还有具体业务的 JSON 反序列化。改动协议细节时要全局搜索所有调用点很容易漏掉一处。第二个问题是没有任何重连机制。SSE 连接的特点是“长时间挂起”网络抖一下、服务端主动断开、代理层空闲超时连接随时会断。第一版代码一旦readLine()返回 null 或者抛 IOException整个流程就终止了用户只能重新发起一次请求体验非常糟糕。第三个问题是线程模型很粗糙。client.send()本身是阻塞调用一个请求从发出到流结束要占住一个线程。后面要讲并发问题源头其实在这里。这三个痛点加在一起让我觉得必须做一个统一的封装层把“怎么读 SSE”这件事从业务代码里彻底剥离出去。3. 现实马上给我上了三课断流、超时、乱码3.1 思考期的静默与“idle timeout waiting for SSE”报错第一课上得最疼。对接某个私有化部署的模型时客户端稳定收到一段内容之后突然卡住几十秒后抛出一个异常stream disconnected before completion: idle timeout waiting for sse。一开始我以为是 Java HttpClient 的读超时但排查后发现根本不是。这条报错的本质是中间链路对“空闲连接”的清理机制大模型生成内容不是匀速的它经常会在某个 token 之后进入“思考”阶段几秒甚至几十秒不吐数据。连接确实还活着但所有中间代理Nginx、API 网关、云负载均衡都默认一个连接太长时间没动静就是死连接直接掐掉。要验证这个判断其实有个土办法用curl -N直连上游服务观察思考停顿期间连接是否还保持。我试下来发现直连完全没问题但只要走网关到时间就被切断基本可以锁定是中间层空闲超时。解决思路分两端。服务端如果有控制权最标准的做法是定时发心跳注释行也就是:开头的行这样连接永远有流量代理层不会觉得它空闲。如果控制不了服务端客户端能做的主要是两件事一是把重连机制做好断线后自动携带上次的事件 ID 重新建立连接二是在前端交互上给用户一个“模型正在思考”的提示而不是看着界面死等。3.2 断线重连Last-Event-ID 到底怎么用SSE 协议里设计了一个非常实用的续传机制服务端可以在事件里带id字段客户端断线重连时把上次最后收到的 ID 放到请求头的Last-Event-ID里服务端看到之后可以从对应位置继续推而不是从头开始。浏览器原生的 EventSource 会自动完成这个过程但 Java HttpClient 没有任何这样的“自动重连”能力所有逻辑都要自己写。我当时的实现思路是在解析循环里维护一个lastEventId变量遇到id:行就更新重连时给 HttpRequest 加一个Last-Event-ID头每次readLine()返回 null 或者捕获到 IOException先判断是否请求被主动取消如果没有取消就进入重连退避逻辑。再补充一点退避策略如果服务端明确抛出了 429限流或者 5xx服务端错误直接退避重连往往是火上浇油应该先等待几秒再从短到长地尝试。我一开始犯的错是无脑立即重连结果正好撞上服务端还在重启白折腾了十几分钟。后来改成指数退避首次等 500ms翻倍到上限 10 秒情况立刻稳定很多。3.3 那些“: heartbeat”注释行和编码问题第三课的来源是心跳行和字符编码两个都是看起来不起眼但很会坑人的细节。先说注释行。某次对接另一个平台时我发现终端里偶尔会出现一行看似乱码的内容仔细看才发现是一行以冒号开头的: heartbeat。这行数据没有任何业务意义就是服务端用来维持连接活性的一口气。按照规范客户端遇到:开头的行必须完全忽略。但我的第一版解析器没有处理这条分支于是心跳行被打进 data 缓存里等空行一来整个 JSON 字符串就变成了{role: assistant}: heartbeat...这种四不像反序列化直接失败。再说编码。SSE 规范规定传输内容必须是 UTF-8但如果直接把响应流交给new InputStreamReader(response.body())默认会使用 JVM 平台默认编码。在 Linux 服务器上一般没问题可在 Windows 上稍有差池就会变成乱码。我后来统一用StandardCharsets.UTF_8显式指定编码从根上杜绝了这个隐患。还有一个小坑是和\r\n相关的。某些网关会在行尾追加回车符虽然BufferedReader.readLine()会把行终止符去掉但如果数据里混了\r空行判断可能会踩雷。稳妥的做法是对每一行做一次trim()再走下一步逻辑。4. 重构把SSE封装成“注册回调就能用”的通用组件4.1 封装目标业务层不碰任何协议细节在正式动手之前我先给自己定了几条验收标准调用方传入 URL、请求头、请求体和几个回调函数其他什么都不用管事件解析、心跳识别、断线重连、连接关闭全部在组件内部完成错误信息通过统一的回调通道暴露出去方便日志监控线程模型可替换为下一步引入虚拟线程留好扩展点。换句话说我想把“怎么读一个 SSE 流”这件事彻底封装成隐式能力业务方看到的只有“我发了一个流式请求内容是一段段回调给我的”。这种接口设计还有另一个好处测试的时候可以很方便地替换成一个本地 mock SSE 服务不依赖真实大模型接口。4.2 核心接口设计与解析器实现封装之后的调用方式长这样try (SSEClient client SSEClient.create(https://gateway.example.com/v1/chat/completions) .header(Authorization, Bearer sk-xxxx) .body(payload) .onEvent(message, (event, data) - { ChatChunk chunk JsonUtil.parse(data, ChatChunk.class); String content chunk.getChoices().get(0).getDelta().getContent(); if (content ! null) { frontendSink.onNext(content); } }) .onDone(() - log.info(流式输出完成)) .onError(err - log.warn(流式连接异常: {}, err.toString())) .build()) { client.connect(); }这里头最核心的是两个东西事件分发器和状态管理。事件分发器维护一个MapString, ConsumerSSEEvent解析器每解析完一个事件就根据事件类型查表并调用对应回调。这样新增一个事件类型只需要在调用方注册一个 handler组件本身不用改。状态管理则负责closed标志、lastEventId、重连次数、当前连接状态这些内部状态。特别强调一下closed所有循环和重连逻辑都先检查它这是后续支持主动取消的关键。解析器本身是一个纯文本扫描器核心逻辑和第二章的手写代码类似但改成了回调驱动并且把“注释行忽略”“多行 data 拼接”“空行触发事件”“[DONE]标记识别”这些边界情况全处理掉了。4.3 心跳检测、自动重连与生命周期管理自动重连放进组件之后细节比想象中多。我遇到过一个比较隐蔽的问题服务端因为业务异常主动断连时TCP 层可能不一定会立刻让readLine()抛异常而是进入半开状态。如果只靠解析循环退出触发重连可能要在几分钟后才发现连接已经死了。所以我在组件里加了个“读等待超时机制”每次进入readLine()之前用一个单独的调度线程倒计时如果在设定时间内没有读到任何新数据就主动关闭当前 InputStream让解析循环退出并触发重连。这里要特别注意一个语义问题重连不代表重新生成内容而是尽量接着上次的位置继续。配合之前维护的lastEventId组件在重建连接时会自动带上Last-Event-ID请求头。如果服务端不支持续传最坏情况是重新推一遍已生成的内容这在业务上也比“卡死到超时”好接受得多。生命周期管理上我让SSEClient实现了AutoCloseable这样调用方可以用 try-with-resources 确保关闭。真正关闭时要做三件事置closed标志、关闭 InputStream、把调度线程里未执行的倒计时取消掉。顺序不能乱先置标志再关流否则解析循环可能把关闭动作误判成“需要重连”。4.4 和WebClient、OkHttp EventSource对比为什么我选了自建很多同事看到这问我Spring WebFlux 的 WebClient 不是有现成的bodyToFlux(ServerSentEvent.class)吗OkHttp 也有 EventSource 封装为什么还要自己造轮子我的理由其实很实际。项目本身是传统 Spring Boot MVC 技术栈如果仅仅为了消费一个 SSE 流就把 WebFlux 和 Reactor 引入进来依赖复杂度、响应式编程的学习成本、线程模型的变化都是实打实的隐性开销。WebFlux 很好但它更适合整个链路都走响应式的场景而不是作为一个孤立的客户端工具去嵌入老项目。OkHttp 的 EventSource 封装得很省心但对底层线程模型拿捏得不够直接尤其是在想用虚拟线程做自定义调度的时候它内部还是按自己的线程池和回调机制来跑反而多了一层隔阂。自己做这个组件核心解析代码也就两百行左右每个线程模型、每条重连策略都能按自己的需求调整可测试性也高。现在回头看这个决策做对了。后面引入虚拟线程时我只需要在调用方加一个 executor整个组件几乎没动。5. 虚拟线程是SSE并发量的“胜负手”5.1 平台线程为什么扛不住SSE长连接封装完之后代码层面已经干净了但并发问题很快就来了。我们的业务需要在一条请求里同时调用多个 AI 会话每个会话都要维持一条 SSE 连接。最开始用的是传统线程池固定 200 个线程跑业务逻辑已经占了一批剩下一百多个线程最多只能同时挂一百多条 SSE 连接新请求立刻在队列里堆积。表面看这是线程池大小问题把线程数调到 500、1000 就行。但平台线程没有那么便宜每个线程栈默认要预留兆级内存线程多了上下文切换成本也非线性上升。在对延迟敏感的服务里靠堆平台线程解决问题迟早会把内存和 CPU 都拖垮。SSE 场景和普通请求还有一个本质差异普通请求在线程上“做事”做完就释放SSE 连接在线程上“等人”可能几十秒甚至几分钟都没有任何数据可读。用平台线程等数据等于雇了一百个人什么都不干就盯着电话等通知太奢侈了。5.2 虚拟线程的原理阻塞时“脱身而去”虚拟线程的调度逻辑用一个词概括就是“阻塞卸载”。它跑在一小组载体线程上一般就是 JVM 内部维护的一个 ForkJoinPool。虚拟线程执行到阻塞操作比如InputStream.read()时JVM 会把当前虚拟线程从载体线程上摘下来也就是 unmount载体线程立刻去执行另一个虚拟线程等之前的 I/O 事件就绪了再把虚拟线程挂回去继续跑。站在应用代码的角度这就是一个普普通通的线程在傻等readLine()。但底层视角完全不同等数据的那段时间里载体线程早就跑去干别的活了。用生活化的类比平台线程是“一个服务员只服务一桌客人干等客人思考”虚拟线程则是“服务员把菜单递给客人之后马上招呼下一桌客人要加水时再回来”。这里有个常见的误区是觉得虚拟线程能提高单条 SSE 的速度。并不会它提高的是并发承载能力。单条流该多快还是多快虚拟线程的价值在于让几千条流同时挂着服务端依然能喘得过气。5.3 接入虚拟线程的最小改动JDK 21 里虚拟线程正式可用接入方式非常直接ExecutorService bridgeExecutor Executors.newVirtualThreadPerTaskExecutor(); // 每个 SSE 连接分配一个虚拟线程阻塞读 InputStream for (AiSessionRequest request : sessionRequestList) { bridgeExecutor.execute(() - { try (SSEClient client buildClientFor(request)) { client.connect(); } catch (Exception e) { log.warn(会话处理失败: {}, e.toString()); } }); }关键点就一个Executors.newVirtualThreadPerTaskExecutor()。它创建的线程池没有固定上限每个任务创建一个虚拟线程任务结束后自动回收。对于 SSE 这种“一条连接一个阻塞线程”的模型虚拟线程几乎是量身定做的。如果你的代码是在SSEClient.connect()内部同步阻塞那么在调用处包一层虚拟线程池就是全部接入工作。不要用synchronized包住长阻塞读因为一些 JDK 版本里虚拟线程在synchronized块内会发生 pinning就是“钉在载体线程上不让走”这会让虚拟线程的优势打折扣。习惯上用ReentrantLock或者干脆避免持锁读流。5.4 参考性压测从500路拥挤到2000路从容我们当时在测试环境做了一组对比场景是模拟真实 AI 会话同时建立多条 SSE 连接每条连接按固定频率收到模型数据再转发给前端。硬件是普通的八核服务器JDK 21。固定线程池 200 个线程时跑到 300 到 400 路并发线程全部阻塞在读流上新增请求开始排队延迟明显上升。500 路左右时新的连接请求基本进不来了。切换到虚拟线程之后同样一套代码并发从 500 路一路加到 2000 路进程里虚拟线程对象数量确实涨到了两千多但载体线程始终维持在几十个CPU 主要花在 JSON 解析和数据转发上而不是线程切换。2000 路并不是上限只是我们当时的测试边界。强调一下这是测试环境参考量级不代表所有场景的绝对数值。但趋势非常明显平台线程的瓶颈在“等数据的线程占住了资源”虚拟线程通过阻塞卸载把这种浪费拿掉了。6. 实践中还要注意的几个小细节6.1 主动取消与正确关流上游内容生成到一半前端用户点了停止这个请求应该被真正取消掉而不是让底层继续傻读。实现方式是在SSEClient里暴露一个cancel()方法内部先置closed标志然后主动关闭InputStream阻塞中的readLine()会因为 IOException 退出解析循环检查closed后不会进入重连直接走结束流程。这里有个常见反模式为了中断读操作直接去interrupt()那个执行解析循环的线程。对平台线程粗暴 interrupt 会引发各种中断异常对虚拟线程来说更是不优雅。关闭流才是让阻塞读退出的正道。还要注意关闭操作要在 finally 块里进行。SSE 解析循环中任何一个 JSON 解析异常都不该让连接泄漏组件内部必须保证InputStream一定被关闭。6.2 不要在SSE回调里做CPU密集处理回调函数默认是在虚拟线程里执行的。虚拟线程擅长的是“等”不是“算”。如果你在onEvent回调里做大量计算、复杂正则、或者同步调用多个下游接口虚拟线程不会在这个阶段卸载多个这样的任务会互相抢载体线程的计算资源。合理的分层是回调里只做轻量解析、内容拼接、投递给消息队列或异步线程池重活全部丢到专门的线程池里执行。这样虚拟线程就能始终专注于它的核心价值——替大家把“等待”这件事扛下来。6.3 和前端配合的两个约定最后说两个前端联调时踩过的约定问题。第一浏览器原生的 EventSource 只支持 GET 请求但大模型接口几乎都是 POST所以前端要么用 fetch 加 ReadableStream 手动解析 SSE要么引入现成的解析库。后端对接时要有这个意识别默认对方也用 EventSource。第二服务端返回的Content-Type必须是text/event-stream不要擅自加charsetutf-8后缀或者改成text/plain否则有的前端解析库会直接拒绝识别为 SSE 流。别问我怎么知道的有一次联调时前端抓包发了三天最后发现就是 Content-Type 写斜了。如果让我重新走一遍我的总结可能只有三句话第一遍想的是“怎么把流读出来”第二遍想的是“如何让别人别再写一遍解析代码”第三遍才意识到“连接数才是真实瓶颈”。虚拟线程最大的价值不是让单条 SSE 变得更快它不会而是让几千条 SSE 同时挂在服务上时我们的进程依然能正常调度、正常处理其他业务。踩过几次坑之后我现在的原则是SSE 客户端这种纯粹的阻塞 I/O 场景方案越简单越好——一条连接一个虚拟线程一个线程只做一件事剩下的交给 JVM 调度。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询