尧图精选

Java后端SSE实战:从手写解析到虚拟线程,搞定AI流式输出

🕒 发布时间:2026/10/1 9:33:22 📁 来源:尧图网络
去年年中我接到一个任务把公司大模型对话平台从等完整回答改成边生成边显示。技术选型摆到面前——轮询、WebSocket、SSE我最后选了SSE。这个决定本身不难难的是后面一连串的事手写SSE解析逻辑时的枯燥和踩坑把它封装成通用组件时的设计取舍以及最后引入虚拟线程后吞吐量直接翻了几倍的惊喜。这篇文章就把这条路径完整记录下来Java后端对接AI流式输出时SSE从显式调用到隐式封装再到虚拟线程带来的性能飞跃。对于正在做AI应用、Java后端和SSE集成的同学这应该是条能直接照着走的路。1. 为什么大模型应用偏偏看上了SSE1.1 AI对话场景对传输通道的三个硬要求AI对话和传统后端接口最大的区别在哪传统接口是请求-响应模型一次把结果给完等多久都无所谓用户反正盯着loading转圈。AI对话完全不是这个体验大模型是逐token生成的用户看到第一个字越早体感越好。这带来三个硬要求第一是增量传输。模型的生成过程是流式的后端必须把生成中的token一段一段推给前端不能等全文生成完再一次性返回。第二是单向通道。整场对话里用户的有效输入基本只有最开始的那条Prompt之后全是服务端往下推内容不存在频繁的客户端上行消息。第三是中断可控。用户随时可能点停止生成服务端收到中断信号后要立刻停掉推送并且给出一个明确的收尾状态。这三个要求放到一起选型的天平就已经倾斜了。1.2 轮询、WebSocket、SSE一次现场对比很多团队的第一反应是WebSocket毕竟它听起来实时、全双工、高级。但冷静看一下AI场景根本用不上双向能力——流式生成过程中客户端几乎不往服务端发消息WebSocket一半的功力是浪费的而你要付出的代价却是实打实的协议升级、心跳保活、帧边界处理、断线重连全都要自己写代码。一个细小的帧没处理好连接就废了。轮询就更不用说了两三秒查一次结果要么延迟大要么服务器被空转的查询打垮。AI生成的回答动辄十几秒轮询期间的大多数请求都是无效请求。SSEServer-Sent Events服务器发送事件恰恰踩在需求点上。它基于普通HTTP连接建立后服务端持续往客户端写数据协议非常简单——MIME类型是text/event-stream消息格式就是几个固定字段按空行分隔。它甚至自带断线重连机制靠Last-Event-ID就能续传。我把三者的对比整理成一张表给团队做调研汇报时就用的这个特性轮询WebSocketSSE数据方向客户端主动请求双向全双工服务端单向推送协议复杂度最低高低实时性差好好断线重连自己实现自己实现协议内置兼容性最好一般好典型场景低频状态查询在线聊天、互动游戏消息推送、AI流式输出结论很直接——AI对话就是SSE的主场。这也是为什么主流大模型开放平台的流式接口几乎全在走SSE这套格式。所以你在搜索引擎里搜Java 实现 SSE翻出来的结果十有八九都在做AI流式接入这不是偶然。1.3 为什么主流大模型平台都在用SSE可以这么理解大模型服务商面对的客户端五花八门有浏览器、有移动端、有后端服务。SSE只需要一个HTTP客户端就能接任何语言都能轻松解析而WebSocket在部分企业网络环境里还会被网关拦。对服务商来说SSE还有一个好处——底层就是HTTP所有现有的鉴权、限流、负载均衡体系都可以直接复用不用为流式协议单独造一套基础设施。这个成本低、兼容广的组合让SSE在AI流式接口里成了事实标准。理解了这一点再去看各家平台的接入文档会发现骨架出奇地一致。2. 显式调用第一次手写Java SSE客户端2.1 SSE线上数据的真实长相在动笔写代码之前先把SSE线上的数据长相看清楚。一个典型的事件流是这样的event: message data: {delta:{content:你好},index:0,finish_reason:null} event: message data: {delta:{content:},index:0,finish_reason:null} event: message data: {delta:{content:今天},index:0,finish_reason:null} event: message data: {delta:{},index:0,finish_reason:stop} data: [DONE]几个要点每条事件之间用空行分隔data:后面跟的数据可以跨多行多个data:行合起来算一条完整消息事件类型由event:指定不写就是默认的message最后的[DONE]是流式结束标记。解析逻辑说穿了就是把data:后面的内容拼起来去掉注释行按空行切事件。但真正手写的时候就会发现坑都在细节里。2.2 用HttpClient逐行读流Java 11开始JDK原生HttpClient已经够用了不需要额外引依赖。最直观的写法是用BodyHandlers.ofInputStream()拿到原始字节流然后包装成BufferedReader逐行读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 sk-xxxx) .header(Accept, text/event-stream) .header(Content-Type, application/json) .POST(HttpRequest.BodyPublishers.ofString({\model\:\gpt-4o-mini\,\messages\:[{\role\:\user\,\content\:\你好\}],\stream\:true})) .build(); HttpResponseInputStream response client.send(request, HttpResponse.BodyHandlers.ofInputStream()); try (BufferedReader reader new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; StringBuilder dataBuffer new StringBuilder(); while ((line reader.readLine()) ! null) { if (line.isEmpty()) { // 空行表示一条事件结束统一处理 if (dataBuffer.length() 0) { handleSseEvent(dataBuffer.toString()); dataBuffer.setLength(0); } } else if (line.startsWith(data:)) { String payload line.substring(5).trim(); if ([DONE].equals(payload)) { handleComplete(); } else { dataBuffer.append(payload); } } // 注释行以冒号开头和event:行按需忽略 } }这段逻辑跑通之后模型生成的文字就能一段一段打到终端里了。第一次看到增量输出出来的时候那种成就感是真实的但紧接着问题就来了——这段代码没法直接上生产。2.3 显式调用暴露出的三个工程问题用了一个礼拜我在这段能跑的代码上陆续发现了三个问题。第一个问题是解析逻辑和业务逻辑完全耦合。上面这段解析代码一旦散落到各个业务方法里接一个模型服务就要复制粘贴一次。后面同时接两家、三家大模型平台每个平台的字段结构还略有差异复制粘贴出来的代码就彻底失控了。第二个问题更致命异常、重连、超时全靠自觉。AI生成过程里网络抖动、模型服务端重启时有发生。一旦流中途断掉客户端默认就是异常退出日志里留一句没人看得懂的堆栈。用户那边看到的就是回答到一半卡死了。没有统一的重连机制没有幂等的续传体验就没法保证。第三个问题是没有监控和降级的抓手。线上看到流断了到底是网关超时、服务端异常还是客户端处理太慢每一条长连接的生命周期状态、吞吐量、错误率如果没有统一记录排查一次事故需要翻好几个系统的日志。这三个问题指向同一个方向SSE接入不能停留在每个接口自己写一套解析的层面得把它封装成一个独立的、可复用的组件。这就进入了下一阶段。3. 隐式封装把SSE变成业务无感知的能力3.1 先定调用契约事件回调还是返回式流封装之前先把调用契约想清楚。市面上常见的封装方式有两种一种是Reactive风格的把SSE流抽象成FluxSseEvent或者StreamSseEvent业务方能拿到一个流对象自己订阅另一种是回调风格的用一个SseListener接口注册事件、错误、完成三个回调。我最终选了回调风格原因很实际项目里对接AI流式接口的业务方大多是写惯了传统Spring MVC的团队他们对Reactive编程模型不熟悉上手成本高。回调接口虽然老派但心智模型简单——连接由组件管事件来了我处理就行。契约定义成这样public interface SseListener { void onEvent(SseEvent event); void onError(Throwable cause); void onComplete(); } public record SseEvent(String id, String event, String data) {}SseEvent保存原始字段业务方拿到data后自己决定怎么解析JSON。组件不做过度封装因为不同模型平台的JSON结构差异太大统一收敛反而会制造一堆条件分支。3.2 内部结构连接管理、消息解析与状态流转组件内部拆成三层连接层负责建立HTTP连接、设置超时、发送请求解析层负责把字节流切成SseEvent处理注释行、多行data:拼接、[DONE]标记调度层负责把事件分发给SseListener同时维护连接状态。这里有一个容易被忽略的细节——状态流转。SSE连接不是简简单单连接中/已断开两种状态我实际维护的状态有这么几个CONNECTING - CONNECTED - STREAMING - IDLE - COMPLETED / FAILED / CANCELLEDIDLE状态特别关键当模型在长时间思考推理模型已经是常态连接上没有数据流动但连接本身还活着。这时候组件既不能误判为失败也不能干等——需要配合心跳机制来判断连接是真的断了还是暂时空闲。我把心跳检查设计成如果超过30秒没有任何数据行到达主动探测一次连接状态连续3次探测无响应才判定为FAILED并触发重连。这个设计直接避免了很多假断连误报。3.3 业务代码的最终形态封装完成后业务方接入的代码收敛成了这个样子SseRequestSpec spec SseRequestSpec.builder() .url(https://api.example.com/v1/chat/completions) .header(Authorization, Bearer sk-xxxx) .body(bodyJson) .build(); sseClient.execute(spec).subscribe(new SseListener() { Override public void onEvent(SseEvent event) { if ([DONE].equals(event.data())) { chatSession.finish(); return; } JsonNode node objectMapper.readTree(event.data()); String token node.path(choices).path(0).path(delta).path(content).asText(); chatSession.append(token); // 推给前端 } Override public void onError(Throwable cause) { chatSession.error(生成中断请重试); } Override public void onComplete() { chatSession.complete(); } });注意重连、心跳、超时这些逻辑全部在execute()内部处理完了业务方感知不到连接层发生的任何事。这轮重构之后新接入一家模型服务只需要写一个Listener解析逻辑从复制粘贴地狱变成了写一次适配。3.4 多模型协议的适配层设计说到适配这是封装SSE时真正的进阶题。不同大模型平台虽然都用SSE框架但data:里的JSON结构千差万别。有的平台直接套OpenAI格式choices[0].delta.content有的用output.text字段有的会在事件中间塞usage信息有的把工具调用结果放在delta.tool_calls里。我的处理方式是把解析SSE传输层和解析业务payload彻底分开。传输层统一出SseEvent这是稳定契约业务payload的解析放到MessageParser接口里每个模型一个实现public interface MessageParser { String parseContent(String data); ListToolCall parseToolCalls(String data); boolean isEnd(String data); }新的Model接入注册一个Parser即可。这个分层拯救了我后面很多次对接——平台方改了字段名我只需要改对应的Parser传输层和业务层完全不受影响。4. 虚拟线程SSE长连接性能瓶颈的破局点4.1 平台线程池模型下SSE长连接是怎么把服务拖垮的SSE封装好了业务方接入顺畅了但性能问题接踵而至。过去Java Web服务的并发模型是一个请求占用一个平台线程。以Tomcat默认配置为例线程池上限通常200左右。普通接口毫秒级返回线程用得快还得快但SSE长连接完全不同——一个连接从建立到结束往往要持续30秒甚至几分钟。假设200个用户同时在等AI生成回复200个线程就全部被占满第201个用户只能在队列里干等。更揪心的是这些被占用的线程大部分时间都在等网络IO——等模型服务端吐出下一个token。线程没有在计算没有在读写数据库就是在那里阻塞着等待。这是对平台线程的极大浪费。有一次压测300路并发SSE直接让接口的响应时间从200ms恶化到3秒以上问题就出在线程池被打满连健康检查的请求都排不上队。4.2 虚拟线程把阻塞变回一件便宜的事JDK 21正式发布了虚拟线程这个局就被破掉了。虚拟线程是什么简单说它是由JVM调度而非操作系统调度的轻量级线程。平台线程和操作系统线程一一对应贵虚拟线程则是挂在平台线程上的用户态线程创建和切换成本低了几个数量级。关键点在于阻塞语义的变化。平台线程一旦阻塞比如socket.read()等数据整个线程就被OS挂起而虚拟线程阻塞时JVM会把它从载体平台线程上卸下来那个平台线程立刻可以去跑别的虚拟线程。等IO数据到了虚拟线程再被调度回去继续执行。所以对虚拟线程来说阻塞不再是浪费资源的罪过——阻塞一个虚拟线程的成本约等于让出一个CPU时间片。放到SSE场景里我们可以给每一条SSE连接分配一个虚拟线程让它大大方方地阻塞着等数据服务端同时挂几千上万个长连接平台线程池依然是空闲的完全不影响其他普通接口的处理。这正好解决了SSE长连接、低计算场景下的线程耗尽问题。4.3 接入虚拟线程的具体改动如果你用的是Spring Boot 3.2及以上版本开启虚拟线程的执行器非常简单spring.threads.virtual.enabledtrue这一行配置会同时影响Tomcat的请求处理执行器和Spring的Async执行器。也就是说进来的HTTP请求会直接在虚拟线程上跑SSE长连接自然也就跑在虚拟线程上了。如果不是Spring Boot需要手动配置Tomcat的协议处理器执行器Bean public TomcatProtocolHandlerCustomizer? protocolHandlerCustomizer() { return protocolHandler - protocolHandler.setExecutor(Executors.newVirtualThreadPerTaskExecutor()); }这里有个容易踩的坑虚拟线程不要用线程池。Executors.newVirtualThreadPerTaskExecutor()虽然名字里带Executor但它不是一个池化执行器每次execute()都会创建一个新的虚拟线程用完即弃。池化虚拟线程是错误用法也没有意义——创建虚拟线程本身开销极小池化反而增加了不必要的复杂度。4.4 性能实测与必须回避的坑我在项目里做了对比压测同样一台32核64G的机器平台线程模式在并发数超过300时接口的P99延迟立刻恶化线程池打满后大面积超时切到虚拟线程后并发升到1200P99依然平稳瓶颈已经不在线程维度而是转移到了网络带宽和CPU的JSON解析上。这个提升对于以SSE长连接为主的AI网关服务来说几乎是质变。但虚拟线程不是银弹有四个坑我在实践中踩过或者观察过第一synchronized会把虚拟线程钉死在底层载体线程上。如果在持锁期间做了阻塞IO虚拟线程无法让出载体线程性能优势就没了。建议把持锁范围内的阻塞操作剥出来放到锁外。第二ThreadLocal在虚拟线程里要慎用几万条虚拟线程的ThreadLocal累积起来内存很可观。JDK 24开始提供了无平台线程限制的ThreadLocal处理模式但生产环境如果还在JDK 21建议直接用ScopedValue或者把事情搬到参数里传递。第三数据库连接池的大小不会因为虚拟线程而变大底层物理连接还是受数据库限制别把连接池配置里最大连接数按虚拟线程数量调整。第四不要试图复用虚拟线程它的设计哲学就是一次性消耗品。5. 生产环境里最要命的坑空闲超时断连与保活方案5.1 流断了一半一次真实的SSE超时事故封装上线后的第三个星期线上报警群里弹出一条让人头皮发麻的消息用户在AI对话里问到第3轮回答到一半流断了。翻开日志核心错误是这么一行stream disconnected before completion: idle timeout waiting for sse那阵子恰好接入了几个带深度思考能力的模型它们有个特点收到问题后会先思考很久十几秒甚至几十秒不出一个字。我们的SSE连接上没有任何数据流动但连接本身活得好好的。结果在某个网关节点上连接被判定为空闲超时直接断掉了。5.2 断连根因链网关、心跳与客户端的三角关系这个报错的根因链条其实很好理解。请求链路是客户端 - 我们的Nginx网关 - 应用服务器 - 模型服务商。Nginx对上游和下游之间的数据转发有一个空闲超时设置默认值通常是60秒——意思是60秒内如果Nginx没有从上游获取到任何数据发给客户端它就认为连接没有意义了主动掐断。这带来一个矛盾对于传统接口60秒没数据确实不正常但对于带思考能力的AI模型来说60秒内不出token是常态甚至可能是特征而非异常。于是在模型沉默期Nginx比客户端还着急先动手把连接断了。客户端这边如果没有做断线重连用户看到的就是回答到一半没了。那SSE不是有内置的断线重连吗协议层面确实有但前提是客户端得实现Last-Event-ID的重连逻辑并且在重连时重新走一遍模型接口把历史消息带上——我们的显式调用版本和第一版封装都没有把重连逻辑做成闭环。5.3 一套总结SSE生产配置清单解决思路是三个角色一起配合服务端发心跳保活、网关放宽超时、客户端设对读超时并做断线重连。第一服务端主动心跳。SSE规范里注释行是合法的数据帧客户端会自动忽略它。我的做法是每30秒往连接里写一个注释行// 定时任务每30秒执行一次 PrintWriter writer response.getWriter(); writer.write(: ping\n\n); writer.flush();注释行让Nginx看到链路上确实有数据流动就不会触发空闲超时了。这个方案不用改任何客户端协议代码成本极低。第二网关侧配置。Nginx针对SSE接口要单独调一组参数location /v1/chat/completions { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; proxy_send_timeout 3600s; }proxy_buffering off很关键——如果开着缓冲Nginx会把上游数据攒够一整块再发给客户端SSE的增量显示效果就直接没了。proxy_read_timeout从默认60秒调到3600秒给模型思考留足空间。第三客户端侧超时需要区分。连接超时connectTimeout可以设短一点比如10秒但读超时readTimeout一定要设大或者干脆不设以心跳探测为准。很多客户端库默认读超时60秒这在传统接口下没问题在AI流式场景下就是事故源。我最后把readTimeout设为0即不超时靠心跳机制来保证连接可靠性。这套配置组合下来断流问题基本绝迹。要验证SSE接口是否正常我推荐一个很土但很好用的方法——用curl -N直接看原始字节流curl -N -H Authorization: Bearer sk-xxxx \ -H Content-Type: application/json \ -d {model:xxx,messages:[{role:user,content:你好}],stream:true} \ https://api.example.com/v1/chat/completions-N参数取消缓冲能看到每个data:帧实时刷出来。如果服务端心跳正常你会看到即使模型在思考也会每隔一段时间刷一行冒号开头的注释。最后再分享一点个人经验SSE这套东西纸上谈兵看着简单真正生产环境的复杂度全在协议之外——网关、心跳、线程模型、监控告警每一个都能让你半夜爬起来。我现在的做法是把SSE连接的生命周期事件也打进日志和指标系统连接建立数、事件吞吐量、平均首包延迟、断连原因分布。有了这些数据AI流式接口的健康状态才真正变得可观测。如果后续要做AI Agent的实时工具调用展示或者想给前端推更细粒度的生成过程状态这套封装的扩展点也都已经留好了。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →