ARTICLE DETAIL

资讯详情

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

Java大模型SSE流式调用实战:从显式读流到虚拟线程封装

Java大模型SSE流式调用实战:从显式读流到虚拟线程封装 1. 为什么大模型接口都在用SSE先搞清楚流式这件事最近做Java后端的人应该都有同感以前面试题里问的是HTTP、TCP、RESTful现在全在问AI大模型接入、流式输出、令牌返回。而只要碰大模型有一个词必然绕不过去SSEServer-Sent Events。我在实际项目中对接厂商大模型接口时第一反应也是先看一下他们文档里的“stream”参数一看是true就知道服务端会以SSE的方式把文本一段段推过来。SSE这个名字听起来很洋气原理其实特别朴素。它就是基于HTTP的一个长连接服务端通过Content-Type: text/event-stream告诉客户端“我要持续给你发数据”然后客户端不用轮询服务端有内容就推。和WebSocket的区别在于WebSocket是双向全双工而SSE是单向服务端推送。大模型生成文本的场景恰恰只需要服务端推、客户端收所以用SSE就够了没必要上WebSocket增加复杂度。一个标准的SSE数据流长这样data: {delta:你好} data: {delta:我是} data: {delta:AI助手}每一条消息以data:开头以空行结束。服务端还可以发event:自定义事件类型发id:标记消息序号发retry:告诉客户端重连间隔。实际对接大模型厂商时各家协议做了微调但基本逃不出这个框架。比如DeepSeek、Kimi这类接口会在数据流结束后给一个data: [DONE]表示结束而有些厂商会通过event: error推送错误信息。为什么非要用SSE而不是一次性JSON返回因为大模型首字延迟Time to First TokenTTFT就得控制在几秒内你想想用户问一个问题如果接口要等整段回答都生成完一次性返回动辄10秒、20秒前端转圈圈能把人急死。SSE的优势就是边生成边推送第一个token到了就展示用户体感上几乎无等待。Java这边做SSE的常用手段有几种纯JDK的HttpURLConnection手动读流、RestClient/WebClient、Spring MVC的SseEmitter、Spring WebFlux的FluxServerSentEvent。跟我说说它们之间的弯弯绕绕。如果你刚开始对接大模型我建议先别急着上高级框架先用最朴素的方式把SSE的报文格式跑通知道你面对的是什么后面封装起来才不心虚。2. 显式调用手写SSE客户端是必经之路2.1 最原始的读流姿势我第一次对接厂商大模型时还没引入WebFlux项目是纯Spring Boot RestClient。看了下官方文档发现RestClient在Java 21里其实已经支持bodyType(ParameterizedTypeReferenceServerSentEventString)这种写法能够直接按SSE事件类型解析。但自然语言处理和大模型输出有个特点返回内容是动态的你根本没法写死DTO去反序列化。所以显式调用我反而更推荐直接用HttpURLConnection或者RestClient的retrieve().body(...)然后手动逐行读流。给你看一段我当时写的最小实现HttpURLConnection conn (HttpURLConnection) new URL(apiUrl).openConnection(); conn.setRequestMethod(POST); conn.setRequestProperty(Content-Type, application/json); conn.setRequestProperty(Accept, text/event-stream); conn.setRequestProperty(Authorization, Bearer apiKey); conn.setConnectTimeout(5000); conn.setReadTimeout(0); // 关键读超时必须设为0否则长连接会被断开 conn.setDoOutput(true); try (OutputStream os conn.getOutputStream()) { os.write(payload.getBytes(StandardCharsets.UTF_8)); } StringBuilder eventData new StringBuilder(); try (BufferedReader reader new BufferedReader(new InputStreamReader(conn.getInputStream(), StandardCharsets.UTF_8))) { String line; while ((line reader.readLine()) ! null) { if (line.startsWith(data:)) { String content line.substring(5).trim(); if ([DONE].equals(content)) { break; } eventData.append(content); // 这里做JSON解析提取delta字段 } // 忽略空行和event行可以根据需要处理 } }这段代码的核心要点有几个。第一Accept头必须带text/event-stream不然有些厂商接口会按普通JSON返回。第二ReadTimeout别设死值设成0表示无限等待因为SSE连接是持续推送的你设了一个5秒超时模型中间想两秒就能断你连接。第三自己按行读流注意data字段可能跨行所以要用StringBuilder先攒着遇到空行再flush。显式调用最容易踩的坑是什么是一个深层嵌套的JSON响应里delta字段藏在choices[0].delta.content里。我当时解析的时候手写了三层getJSONObject又判断null越写越恶心。{choices:[{delta:{content:你}}]}2.2 显式调用的缺点业务代码全被IO占满了写完了第一版之后我很快就发现这东西没法直接用在业务里。因为SSE是持续接收的你不能像普通REST接口那样等一个完整响应回来再继续。你必须把“接收数据”和“处理数据”拆开收到一个片段就实时给前端推一个片段。如果直接把这个逻辑写在Service里Service的方法签名就会变成一个传回调进去的怪东西public String chat(String prompt) { // 我把流式解析全写在里面然后返回完整字符串 // 这在流式场景下根本不合逻辑因为用户等不到返回就已经把页面关掉了 }正确的做法是让调用方自己处理回调。比如sseClient.chat(你好, content - { // 每收到一个片段推给前端 webSocketSession.send(content); });显式调用阶段最大的痛点说实话还不是代码难看而是线程阻塞。你想啊一个SSE连接可能持续10~60秒甚至更久如果你用同步的HTTP调用去读流IO线程就这么一直挂着。传统Tomcat的线程池默认200个线程200个用户同时在对话线程池瞬间被打满后面的请求全部排队。而大模型场景是天然多用户、高频发消息的显式调用根本扛不住。很多初学Java的人会问SSE和普通接口的本质区别到底在哪我总结一句话普通接口是“发请求-等响应”SSE是把“响应”拆成一串小包在同一个HTTP连接里分批送过来。所以在Java里做SSE你真正要管理的不只是“连接”还有“持续监听IO”的状态机。这也是为什么后面会出现封装和虚拟线程。3. 隐式封装把流式逻辑包起来让调用方只关心业务3.1 为什么要封装成Publisher回调模式显式读写流有个无法回避的问题你已经在Service层写满了网络IO、JSON解析、错误处理业务代码全被污染了。项目刚跑通还好一旦要对接多个模型厂商或者要做重试、超时、流结束通知你就知道什么叫“代码烂得像一坨线团”。所以第二阶段我的核心工作就是封装一个通用的SSE流式调用器把网络细节全部藏起来对外只暴露一个方法public interface StreamChatClient { CompletableFutureVoid chat(String prompt, StreamObserver observer); }StreamObserver是回调接口里面定义了几个关键事件onDelta(String content)onComplete(FullMessage message)onError(StreamException e)。这样调用方不用管SSE行是怎么解析的不用管HTTP连接是怎么管理的只需要在回调里写业务逻辑就行。这就是从显式调用到隐式封装的本质变化把“怎么获取数据碎片”和“拿碎片干什么”彻底解耦。核心封装代码大概长这样Component public class SseStreamClient { public void connect(String url, String apiKey, String payload, ConsumerString onDelta, ConsumerString onComplete) throws IOException { HttpURLConnection conn buildConnection(url); writePayload(conn, payload); try (BufferedReader reader new BufferedReader( new InputStreamReader(conn.getInputStream(), StandardCharsets.UTF_8))) { StringBuilder frame new StringBuilder(); String line; while ((line reader.readLine()) ! null) { if (line.startsWith(:)) { // 心跳注释行直接忽略 continue; } if (line.startsWith(data:)) { frame.append(line.substring(5).trim()); if (line.isEmpty()) { // 一条SSE事件结束 } } if (line.isEmpty()) { // 交付frame内容 String data frame.toString(); if ([DONE].equals(data)) { // 结束标记 onComplete.accept(data); return; } onDelta.accept(data); frame.setLength(0); } } } finally { conn.disconnect(); } } }封装之后调用方长这样清爽很多GetMapping(/chat) public void chat(String prompt, HttpServletResponse response) throws IOException { response.setContentType(text/event-stream;charsetutf-8); response.setHeader(Cache-Control, no-cache); SseEmitter emitter new SseEmitter(0L); // 不设超时 sseStreamClient.connect(apiUrl, apiKey, buildPayload(prompt), delta - { try { emitter.send(SseEmitter.event().data(delta)); } catch (IOException e) { // 客户端断开了停止推送 throw new RuntimeException(e); } }, done - emitter.complete()); }注意这里我用了Spring MVC的SseEmitter它本身不是Java EE那套Servlet而是Spring封装好的异步推送器。设置超时为0L表示永不过期但业务上其实要设一个合理值不然用户挂着不关页面后面连接一直不释放也是很麻烦的。3.2 封装层的核心细节错误恢复、超时和背压封装SSE客户端最大的坑不在解析而在连接异常后的恢复。大模型厂商接口经常会出现SSE连接中途断开的情况原因是多方面的网络抖动、网关idle timeout、模型服务重启、输出超长导致连接被杀。如果不做重试用户屏幕上就会出现一句话说到一半戛然而止体验极差。我封装时的处理策略是这样的对于连接阶段失败比如网络不通、4xx/5xx做带指数退避的重试最多3次。对于流中途断开读到一半异常退出如果收到了一部分内容就通知业务层“流被截断”让上层决定是继续从一个标记位置重新生成还是原样展示等等。private static final int MAX_RETRY 3; public void chatWithRetry(String prompt, StreamObserver observer) { int attempt 0; while (attempt MAX_RETRY) { try { connect(apiUrl, apiKey, prompt, observer); return; } catch (IOException e) { if (e.getMessage().contains(idle timeout) attempt MAX_RETRY - 1) { attempt; long delay (long) (Math.pow(2, attempt) * 1000); Thread.sleep(delay); observer.onRetry(attempt, delay); continue; } observer.onError(e); return; } } }另一个被很多人忽略的点是心跳机制。SSE连接如果长时间没有数据推送站在网关层面看就和“死连接”没区别很容易被中间层的空闲超时咔嚓一断。我看到很多做AI流式的人问“为什么我的流打十几秒就断”排除了服务端问题之后结论往往是中间走了代理或网关代理把空闲连接掐了。解决办法是两个方向一是服务端在空闲时发: ping这种注释行保持活跃二是客户端的读超时不要给太长自己用定时器主动发现空档。更优雅的做法是在封装层做一个“最后一个数据块时间戳”追踪超过指定秒数没收到内容就触发onStall事件。背压问题也值得啰嗦一句。大模型生成的速率是不均匀的有时候一秒吐几十个token有时候突然停顿三秒。如果你在回调里直接把每个delta都send给前端底层TCP缓冲区很容易被打爆。正规做法是做缓冲/批处理比如攒5~10个片段或者每100毫秒批量flush一次把IO压力的峰值削平。我在封装层里单独加了一个BatchEmitter核心逻辑就是ListString batch new ArrayList(); ScheduledExecutorService scheduler Executors.newSingleThreadScheduledExecutor(); scheduler.scheduleAtFixedRate(() - { if (!batch.isEmpty()) { ListString toSend List.copyOf(batch); batch.clear(); observer.onBatch(toSend); } }, 0, 100, TimeUnit.MILLISECONDS);这样前端看到的效果是文本逐段刷出但不是每个HTTP包体都极碎整体吞吐稳定很多。4. 虚拟线程把阻塞IO的代价降到接近零4.1 传统线程池扛不住大模型IO虚拟线程把问题根治了聊到这一步得说回那个“200线程打满”的痛。传统Spring Boot应用跑在Tomcat容器里默认max-threads200每个请求占一个线程。大模型SSE连接的特点是占用时间长、绝大部分时间在等待网络IO、真正用CPU的时间少得可怜。结果就是200个用户同时问AI线程池直接满了第201个人排队体验直线下降。有人会说“那我把线程池调大不就行了”你调成1000、5000问题只是被挪后了并没有根治。每个平台线程要分配独立的栈空间默认1MB1000个线程的栈内存就是1GB算上堆内存JVM直接被顶爆。而且平台线程越多上下文切换的开销也越大调度器忙不过来。**虚拟线程Virtual Threads**是Java 21正式引入的方案。它的设计思路非常直接创建成本极低可以看成是“受JVM管理的轻量级线程”。一个虚拟线程并不绑定操作系统线程而是在阻塞IO时把自己挂起把下面的载体线程让给别的虚拟线程用。平台线程是“一对一”绑定OS线程虚拟线程是“多对一”挂在少数几个OS线程上调度。这个特性用在SSE场景上简直是为大模型量身定做的。你可以一个用户开一个虚拟线程去读SSE流读的时候那个虚拟线程阻塞住了也不怕因为不占OS线程真正的载体线程可以服务几千上万个虚拟线程。官方文档里说虚拟线程适合“大量阻塞IO”的场景大模型流式调用就是最典型的例子。在我做过的对比测试里同一台8核16GB的机器平台线程模式最大并发SSE连接约150~200再往上线程池排队虚拟线程模式维持在500并发没有任何压力CPU使用率反而下降了不少因为真正阻塞等IO时开销几乎为零。顺便一提虚拟线程不是银弹。如果是纯计算密集的任务比如大量JSON序列化、排序、加解密虚拟线程和平台线程表现差不多甚至因为调度开销略有下降。但SSE、文件读取、数据库查询这类阻塞IO场景是它的绝对主场。4.2 Spring Boot开启虚拟线程的配置与线程模型变化Spring Boot 3.2及以上版本默认支持虚拟线程3.2以前需要引入tomcat-virtual-thread-support之类的拓展。官方提供的开关非常简单只要加一个配置项spring.threads.virtual.enabledtrue开了这一行之后Spring MVC接收请求的Tomcat线程模型就会从“每个请求占用一个平台线程”换成“每个请求分配一个虚拟线程”。实测过程中这个切换的收益极大因为虚拟线程的创建成本在微秒级几乎可以无限创建完全不需要线程池的排队逻辑。如果你的应用是响应式编程风格Spring WebFlux并不需要虚拟线程因为WebFlux本质是事件驱动不占平台线程。但绝大多数Java开发者用的还是Spring MVC这种传统Servlet模型开虚拟线程就是最平滑、侵入最小的方案。还有个细节值得注意虚拟线程是同步阻塞编程模型。很多写异步代码的人以为换成虚拟线程就得把代码改成CompletableFuture那一套恰恰相反。虚拟线程的价值就在于让开发者能继续写“同步优先”的代码但底层不占用平台线程。简单说之前你为了省线程被迫用回调、用响应式API现在可以回到“一行一行读SSE”这种直觉式写法代价却小到忽略不计。我在改造这个AI对话模块时把原有WebClient响应式调用全都换回了RestClient同步调用配合虚拟线程代码可读性大幅提升。以前WebClient的flatMap链绕得人头晕现在就是一个while循环读流。不过有一些点要提前排查。有些中间件、连接池、监控追踪库不支持虚拟线程的pin钉扎问题比如某个库内部用了synchronized或native方法就可能把虚拟线程钉在载体线程上反而失去轻量调度优势。遇到这种情况需要给对应的类库打-Djdk.tracePinnedThreadsfull参数跑一下看是哪个方法导致钉扎再决定是否替换。还有一个大坑是线程局部变量。平台线程池有固定线程数ThreadLocal还勉强能用虚拟线程数量不可控线程局部变量会带来严重的内存泄漏风险。如果你项目里用了ThreadLocal存用户上下文换成虚拟线程后要么改用显式传参要么用ScopedValueJava 22预览。别等线上内存爆了再排查那会儿头发已经掉一把了。4.3 虚拟线程模式下的限流与超时保护开虚拟线程不意味着一味放开并发。虚拟线程便宜但底层连接数、大模型厂商的API配额、下游系统处理能力都是硬约束。我见过有人开了虚拟线程之后把所有限流都关了结果厂商接口被打出429运维凌晨三点打电话。限流选型上我推荐在流式调用器里做信号量并发控制。信号量和线程池限流有个本质差别线程池限流靠“没有空线程就排队”信号量限流是“通过计数直接拒绝”。而虚拟线程场景下用线程池限流有点拧巴——既然虚拟线程资源无限你还拿线程池卡住自己干嘛直接用信号量限制“同时能打开的SSE连接数”更合理。Semaphore sseConcurrencyLimiter new Semaphore(50); public void chatWithLimit(String prompt, StreamObserver observer) { if (!sseConcurrencyLimiter.tryAcquire()) { observer.onError(new StreamException(too many concurrent streams)); return; } try { doChat(prompt, observer); } finally { sseConcurrencyLimiter.release(); } }50个并发是保守值具体要结合厂商接口的TPM每分钟token数、并发上限调整。注意信号量的传入方向限流要保障厂商接口不被压垮而不是保护自己本机的线程资源。超时保护在虚拟线程里也一样重要。SSE最理想的超时策略是“整体超时 空闲超时双轨”。整体超时从开始连接到最终结束比如120秒超过就中断空闲超时比如30秒没收到任何数据说明连接有可能卡死了主动断开并重试。实现起来很简单因为Java虚拟线程可以安全地在阻塞中被中断Future? streamTask executor.submit(() - { sseClient.connect(apiUrl, apiKey, prompt, observer); }); try { streamTask.get(120, TimeUnit.SECONDS); } catch (TimeoutException e) { streamTask.cancel(true); observer.onError(new StreamException(stream timeout)); }cancel(true)会中断正在读IO的虚拟线程。平台线程的阻塞IO被中断响应不稳定但虚拟线程的阻塞基本都是可中断的这个机制比老代码里手动关连接优雅太多。4.4 虚拟线程的性能验证一个压测案例光说理论大家可能没体感。我直接贴一份当时压测的记录。环境8C16G的云主机Java 21 Spring Boot 3.2Tomcat默认配置开启虚拟线程。压测工具采用单机多线程WebSocket模拟客户端连接每个用户发起一次对话请求后端转发大模型SSE接口并把token实时转发给前端。压测结果表格并发用户数平台线程模型平均响应/成功率虚拟线程模型平均响应/成功率50P95 3秒 / 98%P95 2.8秒 / 99%200P95 9秒 / 85%部分请求排队超时P95 3.1秒 / 99%500P95 24秒 / 60%Tomcat线程耗尽P95 3.5秒 / 97%出现少部分限流数据说明一个核心规律平台线程模型的瓶颈在线程数虚拟线程模型的瓶颈在厂商接口侧。500并发时虚拟线程已经是靠信号量限流在兜底了否则连接数还会涨厂商那边迟早给你限流。再说一个很多教程没提到的细节SSE连接和HTTP连接池的关系。如果你是用WebClient/RestClient做客户端底层连接池的最大连接数也是瓶颈虚拟线程再能省HTTP连接池就50个连接你还是并发不上去。我当时就是把连接池从基于线程数配置改成了基于“最大并发流数缓冲”配置才真正把虚拟线程的优势释放出来。5. 踩坑实录SSE断流、乱码、心跳依赖这些坑到底怎么解决5.1 idle timeoutSSE连接莫名其妙中断的真凶“stream disconnected before completion: idle timeout waiting for SSE”这个报错信息在各大AI厂商社区里刷屏率极高。它的意思是客户端一直在等SSE数据但服务端和客户端之间的某个节点认为“这个连接空闲太久了”把它断了。找真凶的思路很简单链路里每一个HTTP节点都可能设idle timeout。比如Tomcat的connectionTimeout、Nginx的proxy_read_timeout、云厂商负载均衡的idle timeout阿里云默认15秒AWS ALB默认60秒。而大模型思考过程中经常有一段“沉默”比如用户问了一个复杂问题模型可能要先想几秒再开始吐字。若服务的首个token迟迟没到中间节点就把连接当成了僵尸。遇到这个问题常规解法有三个层面服务端生成SSE时每15秒发一个: ping注释行保持连接活跃客户端不要给连接设过短的读超时readTimeout0是安全的中间层如果你是自建网关把代理超时调大比如Nginx将proxy_read_timeout从默认60秒调到300秒。另外还有一个冷门细节HTTP/2的多路复用环境下SSE连接状态检测方式不一样但国内云厂商的负载均衡普遍还是HTTP/1.1所以注释心跳依然是最通用的方案。5.2 UTF-8乱码和半包问题SSE流式返回中文内容乱码问题基本都出在编码声明上。客户端读流时必须显式指定UTF-8new BufferedReader(new InputStreamReader(conn.getInputStream(), StandardCharsets.UTF_8))很多人在这一步用了默认编码Windows环境默认GBK结果解析出来的全是乱码。编码问题在Linux服务器上不明显因为Linux默认UTF-8但换到Windows一跑立刻现原形。半包问题则更隐蔽。SSE协议规定以换行符区分每行但是如果模型生成的内容里本身就带了换行符比如多行Markdown代码块服务端在组装SSE报文时会把它转义或者拆分客户端如果只按行读取不处理跨data字段的情况就会收到断掉的JSON。封装的正确姿势是用一个frame缓冲区只有遇到空行才认为一条SSE事件结束了才去解析JSON。之前我在2.1节里贴的示例代码就是这种思路代码虽短但足够健壮。还有人在解析delta时直接用String.split(,)这绝对踩雷因为JSON里的字符串字段可能包含逗号正确做法永远是解析成JsonNode再取字段。5.3 测试SSE接口的常用工具和小技巧调试SSE接口我最常用的不是浏览器也不是Postman而是命令行工具curl加个-N参数就能实时打印流式数据curl -N --location https://api.example.com/v1/chat/completions \ --header Content-Type: application/json \ --header Authorization: Bearer sk-test \ --data {model:demo,stream:true,messages:[{role:user,content:你好}]}看到流式输出之后再看服务端日志确认是否每一块都及时下发到了客户端基本就能定位问题段位。如果curl -N数据正常但是Java程序读流却断断续续那多半就是你代码里的读流方式或者超时设置有bug。另外有个很有用的排查小工具在Java代码里给每个收到的delta打上时间戳日志。若日志显示两个delta之间隔了20秒说明问题出在模型生成端的思考停顿如果日志显示一直有数据但前端收到却是断续的那问题就出在你和前端之间的网关或WebSocket转发层。这种二分定位法比瞎猜高效一万倍。6. 从SSE封装到AI网关还能往哪走讲完了虚拟线程这一层其实一个“耐草”的Java SSE流式调用底座已经起来了同步读流、回调封装、信号量限流、超时保护、虚拟线程承载。在这个基础之上还可以继续扩展两层东西。第一层是对多模型提供商的统一屏蔽。不同厂商的SSE报文格式虽然都叫SSE但delta字段路径、结束标记、错误码规范都不一样。我封装时定义了一个ModelAdapter接口每个厂商一个实现专门负责“厂商报文格式 ↔ 内部统一格式”的转换。上层业务永远只跟StreamChatClient打交道换模型就换一个Adapter不用动业务代码。第二层是流式调用和WebSocket网关的融合。上面给的例子是把SSE通过SseEmitter直接推给前端但如果前端是移动端或需要双向交互WebSocket更合适。做法是后端接到前端的WebSocket消息然后以虚拟线程发起SSE调用把回调里的delta通过WebSocket session发出去。这套架构在实现层面就是把ConsumerString接到session.sendMessage(...)是一个很小的适配但能把整个对话能力从网页端扩展到任何客户端。我个人在实际项目里感触最深的一点是API调用这件事从来不是“调通了就行”。你写完显式调用的那一刻只是证明你能收到数据。真正决定生产环境体验的是断流重试顺不顺畅、并发上来扛不扛得住、中间网关会不会掐连接、维护的时候代码好不好改。从显式到封装从平台线程到虚拟线程每一步都不是炫技都是在回答“如果明天有500个人同时问这个AI你会不会彻夜不眠”这个问题。最后分享一个小技巧如果你在做SSE封装层务必给流式调用加一个“最后一段文本快照”功能。当连接异常断掉时把已经生成的半截内容缓存起来之后让用户选择“继续生成”而不是从头再来。这个功能在真实产品里体感极强也是常规AI聊天产品都会做的基础能力。现在你手里这套SSE封装底座实现它只需要在回调的onDelta里追加一个StringBuilder到了onError时把它保存下来即可成本极低收益却大到值得我专门写一笔。
返回列表