ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java+AI的SSE流式输出:从显式调用到虚拟线程实战

Java+AI的SSE流式输出:从显式调用到虚拟线程实战 做AI应用的Java后端绕不开SSEServer-Sent Events。眼下大模型产品基本都默认用SSE做流式输出——用户提问、模型逐字返回、页面像打字机一样滚动。而Java服务端要接这些大模型接口、要把自己内部的AI能力透出给前端都得跟SSE打交道。这篇文章我就从实战角度把JavaAI场景下SSE的完整演进路线拆开讲透先是原始的显式调用怎么一步步手搓再到怎么封装成一套统一的流式接口调用逻辑最后讲讲为什么虚拟线程一上手整个吞吐量直接上一个数量级。内容都是我在真实项目里趟过坑之后的总结适合那些用过SSE但没系统做过封装的人也适合并发一上来就被线程池卡死的同学参考。1. SSE到底在AI场景里解决了什么问题1.1 SSE的协议本质与AI场景的天然契合SSE全称Server-Sent Events翻译过来就是“服务器发送事件”。它本质上还是HTTP协议只是把响应头的Content-Type设置成text/event-stream然后连接不关闭服务端可以持续往这个响应流里吐数据。我打个比方普通HTTP请求就像你去餐厅点菜服务员把菜端上来这一单就结束了。而SSE相当于服务员端上来一个火锅锅底一直在加热菜品可以源源不断往锅里加。这个特性对AI场景简直是量身定做。大模型生成回答是一个token一个token往外蹦的一个完整回答可能要15到30秒才能生成完。如果用传统HTTP前端得一直等着用户看着空白页面心里直发慌。用SSE的话生成一个token推送一个token用户立刻就能看到反馈体验完全是两个级别。AI场景为什么都在SSE和WebSocket之间选了SSE我总结出几个关键点。第一SSE是单向的服务端推给客户端就够了。AI对话场景里客户端只需要把用户问题发一次剩下的全是接收方向压根不需要双向通信。WebSocket解决的是双向实时通信问题拿它来干单向推送的活不是不行但属于高射炮打蚊子握手复杂度、协议复杂度、二进制帧解析全都白白增加了成本。第二SSE天然跑在HTTP上穿透性好。WebSocket需要专门的握手升级中间有代理网关还需要额外配置支持。SSE就是普通的HTTP响应Nginx、负载均衡、云厂商的网关全都能直接转发不需要做任何特殊处理。第三SSE自带断线重连机制。规范里规定了retry字段客户端断线后会按这个时间间隔自动重连而且服务端可以通过Last-Event-ID字段告诉客户端上次收到哪了。WebSocket的重连逻辑得自己手写。我见过不少团队在技术选型时纠结SSE还是WebSocket这个表应该能帮大家快速下定决心对比维度SSEWebSocket前端轮询连接类型普通HTTP长连接独立协议升级连接普通HTTP短连接通信方向服务端单向推送给客户端双向客户端主动拉取数据格式纯文本UTF-8按行协议二进制帧或文本帧任意断线重连协议自带retry机制需要自己实现自己实现网关穿透直接透传无需额外配置需要支持Upgrade直接透传复杂度低JS原生支持EventSource高需要WebSocket客户端库最低适合场景内容流式下发、实时通知双人互动、游戏、即时聊天低频状态查询1.2 一个绕不开的坑idle timeout waiting for SSE刚才说到SSE穿透性好但它也有一个非常经典的坑就是标题里那个报错stream disconnected before completion: idle timeout waiting for sse。这个报错的本质是连接虽然是长连接但中间的网络设备Nginx、云负载均衡、API网关都有空闲超时设置。默认情况下Nginx的proxy_read_timeout是60秒如果60秒内这个连接上没有数据传输网管设备就认为连接已经死了主动把它断开。SSE在实际AI调用里特别容易触发这个问题。因为AI生成内容不是匀速的有些大模型在思考阶段可能要憋十几秒不说话尤其是复杂推理任务中间会有很长时间的空窗期。这个空窗期一旦超过网管设备的idle threshold连接就被掐断了。前端表现就是等到一半冒烟了报个红字就中断了。后面我在封装章节会专门讲怎么解决这个问题这里先给两个最直接的思路服务端定时往连接里塞心跳注释行SSE协议里注释行为:开头的行浏览器会自动忽略或者前端拿到idle timeout报错后立即自动重连。两种方案配合使用效果更好。2. Java侧显式实现SSE调用的实战记录2.1 服务端怎么往SSE连接里写数据在讲客户端封装之前得先说清楚服务端是怎么吐数据的不然客户端解析时容易踩坑。我用Spring框架举例子最常见的做法是用SseEmitter。RestController public class ChatController { GetMapping(/chat/stream) public SseEmitter chat(RequestParam(question) String question) { // 0表示不超时实际生产里建议设置一个合理的超时时间 SseEmitter emitter new SseEmitter(0L); // 用线程池异步执行避免占满HTTP请求线程 ExecutorService executor Executors.newFixedThreadPool(8); executor.execute(() - { try { // 模拟AI逐字生成 for (int i 0; i 10; i) { emitter.send(SseEmitter.event() .name(message) .data(第 i 块内容)); Thread.sleep(1000); } emitter.send(SseEmitter.event() .name(done) .data([DONE])); emitter.complete(); } catch (Exception e) { emitter.completeWithError(e); } }); return emitter; } }这段代码往连接里发送的内容在HTTP层面长这样event: message data: 第0块内容 event: message data: 第1块内容 event: done data: [DONE]注意每两个事件之间用空行分隔event:是事件类型data:是实际数据载荷。客户端解析就是按这个格式来拆的。这里有个细节值得说为什么好多AI服务商都习惯用[DONE]作为结束标记因为SSE协议本身没有定义“消息结束”的事件服务端关闭连接就代表流结束了但客户端希望显式地收到一个“我说完了”的信号这样它知道该停止渲染、把流式数据拼装成完整结果。所以[DONE]是各家默认的约定俗成标记我在封装层里也沿用了这个设计。2.2 显式调用手写HTTP客户端逐步解析数据流在没有封装层之前客户端接入SSE是很痛苦的。我得手动发HTTP请求、手动读输入流、手动按行拆事件、手动拼JSON。最早我用JDK 11自带的java.net.http.HttpClient来实现代码大致长这样HttpClient client HttpClient.newHttpClient(); HttpRequest request HttpRequest.newBuilder() .uri(URI.create(http://localhost:8080/chat/stream?question你好)) .header(Accept, text/event-stream) .GET() .build(); HttpResponseInputStream response client.send(request, HttpResponse.BodyHandlers.ofInputStream()); try (BufferedReader reader new BufferedReader(new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; StringBuilder jsonBuffer new StringBuilder(); while ((line reader.readLine()) ! null) { if (line.startsWith(event:)) { // 事件类型当前场景基本用不到但可以记录 String eventType line.substring(6).trim(); } else if (line.startsWith(data:)) { // data内容可能被分块需要累积拼接 jsonBuffer.append(line.substring(5).trim()); } else if (line.isEmpty()) { // 空行代表一个事件结束此时才是一个完整体 String json jsonBuffer.toString(); if ([DONE].equals(json)) { break; } // 真正的业务解析逻辑写在这里 handleMessage(json); jsonBuffer.setLength(0); } } }这里面的坑我列一下全是真实的血泪教训第一data:后面的JSON不一定是完整的一行。大模型输出的长文本JSON可能被TCP分片拆成多行每行都是data:开头需要累积到buffer里等空行出现再合并解析。新手最容易犯的错就是每拿到一行data:就当作完整JSON去解析结果就是JsonParseException满天飞。第二readLine()是阻塞的。这个阻塞有个问题网络正常时没事但服务端长时间不发送数据时这个线程会一直卡在这里。如果服务端不主动关闭连接整个线程就永远挂在那了。在平台线程时代这个问题极其致命。第三异常处理非常棘手。连接被网管断开、服务端报错、JSON解析失败各种异常都要单独处理完全没有统一出口。项目里如果每个业务方都这么写一遍那代码风格能乱到天际。2.3 显式调用的四大痛点总结这段经历让我把显式调用的痛点摸得一清二楚总结下来就是四个字各自为政。第一个痛点是代码高度重复。每次对接一个新的AI服务商我都要把上面这段读取流、拼接buffer、处理空行的逻辑重新写一遍换一个服务商就换一个URLJSON解析还得跟着对方的响应结构走。一个项目里同时存在三四个AI服务商的接口时代码里能翻出三四份长得几乎一模一样但又不完全相同的SSE读取逻辑维护起来让人崩溃。第二个痛点是断线重连逻辑缺失。手写的客户端基本上没有重连能力一旦连接断开就抛出异常整个流程彻底结束。而AI流式接口的可用性通常达不到100%内网抖动、服务端重启、机房网络波动都会导致连接中断没有重连用户就得反复重发问题。第三个痛点是超时控制基本没有。readLine()方法默认无限期阻塞一旦服务端不返回数据也不断开连接这个线程会占着资源直到地老天荒。平台线程被这种僵尸连接占久了Tomcat线程池迟早枯竭。第四个痛点是和虚拟线程、线程池等基础设施的耦合。每个业务方都自己开线程、自己管理生命周期代码里全是new Thread(...).start()或者私自创建的ExecutorService资源完全失控。我在一个老项目里见过一个类一个线程池的情况十几个类就是十几个线程池jstack一拉全是睡眠线程看得人血压飙升。3. 从显式到隐式统一SSE流式调用封装的设计3.1 封装层要解决的终极问题被显式调用反复折磨之后我开始设计封装层。这个封装层要达成的目标很纯粹调用方不需要知道SSE协议细节不需要处理断线重连不需要关心线程管理只需要提供URL、参数和几个回调方法剩下的事全由封装层搞定。设计上我参考了OkHttp的Callback模式把整个SSE消费过程抽象成生命周期连接建立成功时触发onOpen每收到一个完整事件时触发onMessage流式输出结束时触发onComplete任何异常导致流中断时触发onError基于这个设计最终调用方的代码应该是这样清爽SseClient client SseClient.builder() .connectTimeout(Duration.ofSeconds(3)) .readTimeout(Duration.ofSeconds(0)) // 0表示不限制读取超时 .autoReconnect(true) .retryInterval(Duration.ofSeconds(2)) .build(); client.subscribe(https://api.llm.example.com/v1/chat, paramMap, new SseCallback() { Override public void onOpen(Response response) { // 连接建立可以在这里做日志记录 } Override public void onMessage(SseEvent event) { // 每收到一个完整事件这里就是业务处理入口 // event.eventName() 获取类型event.data() 获取JSON字符串 String json event.data(); if ([DONE].equals(json)) { return; } // 业务方只关心这个把增量内容渲染到页面 render(json); } Override public void onComplete() { // 整个流结束在这里做收尾工作 } Override public void onError(Throwable t) { // 发生异常在这里做降级或告警 log.error(SSE stream error, t); } });3.2 核心实现拆解SseClient的内部工作机制这个封装层的核心实现我拆成四个模块来设计和实现。模块一连接管理。这部分用JDK自带的HttpClient发起异步请求返回CompletableFutureHttpResponseInputStream然后设置一个专门的线程去读取流。连接超时用connectTimeout控制默认给3秒比较合理——AI接口如果3秒连不上说明网络基本有问题了没必要干等。模块二流式解析器。这是整个封装层的灵魂。我把它设计成一个状态机从BufferedReader里逐行读取遇到event:开头记录事件名遇到data:开头把内容追加进StringBuilder遇到空行说明一个事件结束触发回调。private SseEvent parseLineByLine(BufferedReader reader) throws IOException { StringBuilder dataBuffer new StringBuilder(); String eventName message; // 默认事件名 String line; while ((line reader.readLine()) ! null) { if (line.isEmpty()) { // 空行事件结束 if (dataBuffer.length() 0) { return new SseEvent(eventName, dataBuffer.toString()); } continue; } if (line.startsWith(:)) { // 注释行心跳包直接忽略 continue; } if (line.startsWith(event:)) { eventName line.substring(6).trim(); } else if (line.startsWith(data:)) { // 追加并加换行符防止JSON字符串中的换行被吞掉 dataBuffer.append(line.substring(5).trim()); dataBuffer.append(\n); } // 其他字段如id、retry这里暂不处理 } return null; // 流正常结束 }这个模块有个细节我吃了好几次亏才注意到data:后面的JSON本身可能包含多行比如JSON字符串里有个\n它在SSE协议里实际上会被拆成两行data:所以追加进buffer时要把换行符也带进去否则拼出来的JSON是不完整的。模块三心跳保活线程。这是专门用来对抗idle timeout waiting for sse的。我开启一个定时任务每隔15秒往连接里写一个空行或者注释行。SSE规范里注释行以:开头发送后服务端和客户端都会忽略它但它能让中间网管设备感知到连接是活的从而重置空闲超时计时器。ScheduledExecutorService heartBeatScheduler Executors.newSingleThreadScheduledExecutor(r - { Thread t new Thread(r, sse-heartbeat); t.setDaemon(true); return t; }); private void startHeartbeat() { heartbeatTask heartBeatScheduler.scheduleAtFixedRate(() - { try { outputStream.write(: heartbeat\n\n.getBytes(StandardCharsets.UTF_8)); outputStream.flush(); } catch (IOException ignored) { // 连接可能已断开心跳发送失败没关系读取线程会感知到 } }, 5, 15, TimeUnit.SECONDS); }模块四自动重连机制。当流被中断且服务端没有返回明确的[DONE]标记时自动重连就有用武之地了。我设置了重试次数上限默认3次超过上限就放弃并触发onError回调防止无限重连把服务端打挂。同时实现了指数退避策略第一次重连等待1秒第二次2秒第三次4秒避免断线风暴。3.3 事件模型与业务解耦的细节封装层做完后我意识到一个关键问题onMessage回调拿着的是原始JSON字符串业务方拿回去还是要自己解析。不同的AI服务商返回的结构不一样有的返回{content: 你好}有的返回{choices: [{delta: {content: 你好}}]}这个差异不能指望业务方去适配。所以我在回调参数上加了一层轻量级抽象——SseEvent对象往里面塞两个方法eventName()返回事件类型data()返回原始数据。至于这个data()要不要解析成结构化的DTO我的建议是不要。因为AI服务商的接口结构变化太快今天给你content字段明天可能就改成delta字段封装层强绑定对方的协议反而失去灵活性。业务方在自己这层做协议适配改起来成本最低。3.4 隐式封装后生产环境的接入效果封装完的这套东西上线后最直接的感受就是新接一个AI服务商从原来的半天工作量减到半小时。业务方只需要写SseCallback的实现类在里面处理onMessage回调把增量内容推给前端。断线重连、心跳保活、超时控制这些基础设施层面的东西完全透明了。我拿生产环境的数据来说话之前显式调用时期线上平均每天能收到七八次“流式输出中断”的用户反馈前端页面直接白屏报错。封装层上线后配合前端的重连逻辑这个问题基本消失。就算偶尔网络抖动断线后端自动重连也会悄悄接上用户感知差异小很多。4. 虚拟线程彻底干掉SSE阻塞IO的吞吐天花板4.1 平台线程时代的阻塞地狱封装层解决的是代码复用问题但它没有解决一个更根本的问题——线程资源消耗。平台线程也就是传统意义上的Java线程本质上是对操作系统线程的包装。OS线程是稀缺资源创建要分配栈空间切换要保存和恢复CPU上下文一个OS线程的默认栈大小在1MB左右。而JVM创建线程时虽然不需要OS去分配物理内存但背后依然是直接映射到系统调用。在平台线程时代我们处理SSE流式调用时是这样的一个SSE连接分配一个线程这个线程阻塞在readLine()上等数据。如果同时有200个用户发起AI对话就需要200个线程池里的线程全部卡在读数据上。等到请求再多一点线程池满了新的请求只能在队列里排队。更可怕的是队列里的请求多了之后用户的等待时间成倍增加。我自己遇到过一次真实的线上事故AI对话功能上线后下午两点高峰期瞬时并发200多Tomcat默认200个线程直接被打满。所有线程全卡在SSE的阻塞读上健康检查请求反而抢不到线程。监控面板上一片飘红最后只能紧急重启并限制并发接入数。4.2 虚拟线程为什么是“性能飞跃”的答案JDK 21正式发布虚拟线程Virtual ThreadsJEP 444之后这个问题迎刃而解。虚拟线程是JVM层面的轻量级线程它的创建成本极低生命周期极短和OS线程的映射关系是动态的。当一个虚拟线程执行到阻塞操作比如readLine()时JVM会自动把它从底层平台线程上摘下来unmount让这个平台线程去执行其他就绪的虚拟线程当阻塞操作完成时再把这个虚拟线程挂到某个空闲的平台线程上继续执行mount。这个过程用大白话说就是虚拟线程不是一个人占一个办公位而是一堆人在一个工位上轮流干活。谁在等数据就让出工位谁有数据了就回来坐着干活。工位数量平台线程数永远不必等于干活的人数虚拟线程数几十个平台线程就能支撑成千上万个虚拟线程。这对SSE场景是绝配。因为SSE的readLine()是典型的阻塞IO阻塞时正好触发虚拟线程的unmount机制完全不占资源。改造成本还低代码都不用大改核心就是把线程创建方式换掉// 旧方案每个SSE连接一个平台线程 ExecutorService executor Executors.newFixedThreadPool(200); // 新方案每个SSE连接一个虚拟线程 ExecutorService executor Executors.newVirtualThreadPerTaskExecutor();或者更彻底一点凡是启动独立的SSE消费任务直接Thread.ofVirtual().name(sse-consumer- connectId).start(() - { sseClient.subscribe(url, params, callback); });只改了这两处其他代码完全不用动。4.3 实测数据与压测结果对比我在测试环境做了一轮对比压测机器配置是4核8G跑一个Spring Boot应用模拟SSE服务端每500毫秒推一个事件持续40秒。客户端分别用平台线程池和虚拟线程发起SSE订阅观察最大并发和资源占用情况。指标平台线程池200线程虚拟线程按需创建提升幅度最大支持并发SSE连接数200线程耗尽后排队500025倍以上内存占用200并发时约2.1GB约800MB60%下降线程栈内存占用200MB按1MB/线程算基本可忽略数量级下降阻塞时的CPU占用高频繁上下文切换低unmount无上下文切换开销数量级下降代码改动成本无线程创建方式替换即可极低真实业务场景里的变化更明显之前线上只能同时支持180路AI对话流再多就要排队虚拟线程上线后同时跑了3000路也没看到线程池告警压测到5000路时应用还是稳的内存也就多了几百MB。这个提升幅度用“性能飞跃”来形容毫不夸张。4.4 虚拟线程在SSE场景的三条使用铁律虚拟线程确实强但也不是无脑用的。我在实践中总结出三条铁律写出来给大家避坑。第一条不要在synchronized块里做阻塞IO。虚拟线程遇到ReentrantLock时是正常unmount的但遇到synchronized同步块时会直接钉死在平台线程上pinned这会阻塞平台线程导致整体性能下降。如果需要加锁保护共享资源优先用ReentrantLock或者干脆把IO移到锁外面。第二条不要池化虚拟线程。虚拟线程的价值就在于“按需创建、用完即弃”池化反而抹掉了它的优势。我见过有人把虚拟线程塞进线程池复用结果白白多了一层队列和调度逻辑还容易造成信号量冲突。直接每次调用Thread.ofVirtual().start()就行。第三条谨慎使用ThreadLocal。虚拟线程的ThreadLocal在任务结束时会自动清除但问题是它会在虚拟线程切换载体时带来额外的存储开销。JDK 22JEP 464推出ScopedValue以后推荐用ScopedValue替代ThreadLocal来传递上下文尤其在高并发SSE场景下这个优化效果明显。5. 常见问题速查与避坑指南最后整理一份我在生产环境里踩过的坑和对应的解决方式算是给各位的抄作业清单。问题现象根本原因解决方案报错idle timeout waiting for sse中间网管设备空闲超时断开连接服务端定时写心跳注释行或调大Nginx的proxy_read_timeout收到多个data:行但JSON解析失败SSE数据被TCP分片拆分JSON跨多行累积到buffer里空行出现时再合并解析线程池被打满请求全部排队平台线程阻塞在SSE读取IO上全链路替换为虚拟线程断线后没有收到任何错误直接卡死readLine()无限期阻塞且没有超时控制设置合理的readTimeout并启用自动重连服务端已经返回完数据但客户端还在等服务端忘记调用emitter.complete()关闭连接服务端显式发送[DONE]标记后关闭连接虚拟线程压测时吞吐反而下降synchronized钉死虚拟线程导致平台线程阻塞检查代码把同步块里的阻塞IO移出去心跳注释行导致前端业务解析异常心跳数据混入业务数据流确认心跳行格式为:开头浏览器/客户端自动忽略还有个容易被忽视的点SSE与前端EventSource的兼容性。如果前端直接用浏览器原生EventSource它只会监听message事件自定义事件名必须用addEventListener(自定义名, handler)来监听。所以封装层的SseEvent.eventName()一定要原样透传别自作主张把事件名改写掉否则前后端对不上。另外我建议在SSE接入上设置全链路超时控制。我用的参数是这样的连接超时3秒、总读取超时5分钟、心跳间隔15秒、自动重连次数3次、重连间隔指数退避从1秒开始。这套参数在生产环境跑了大半年效果稳定。每个项目的数据量级和网络环境不一样参数需要微调但思路和原则是通用的宁可连接主动断开重来也不要让一个僵尸连接占着服务器资源不干活。从最初手写HttpClient逐行读流到封装出带心跳和重连的SseClient再到虚拟线程彻底释放阻塞IO的隐患这一路走下来我自己最大的体会是SSE本身不复杂真正决定一个AI服务稳定性的是那些协议边界之外的工程细节。把这些细节沉淀成可复用的组件Java AI的后端开发才能真正做到“少踩坑、多睡安稳觉”。
返回列表