首页 / 资讯中心 / 文章详情

Java对接大模型SSE流式接口:从HttpClient到虚拟线程

Java对接大模型SSE流式接口:从HttpClient到虚拟线程 ★ FEATURED ARTICLE
最近在折腾Java接入大模型流式接口绕不开的就是SSE这条路。不管你是调哪家大模型的chat接口只要开streamtrue服务端返回的就是text/event-stream格式的SSE流也就是Server-Sent Events。Java这边处理这种流式推送从最早的HttpClient手动逐行读流、自己解析data字段到后面把整套逻辑封装成通用组件再到JDK 21虚拟线程出来之后把并发量真正扛上去解决方案的形态发生了三次比较大的变化每一波我都踩了不少坑。这篇文章想把这套演进完整梳理一遍正在用Java对接大模型、准备做AI Agent的朋友应该会有共鸣。1. 从SSE说起AI流式输出为什么要选它1.1 一次真实的SSE响应长什么样SSE本质上就是HTTP服务端向客户端“挤牙膏”式地返回数据。普通HTTP响应是拿到全部数据才结束而SSE在响应头里通过Content-Type: text/event-stream告诉客户端“我要一点一点推你别急着断开”。每次推送的消息是纯文本以data:前缀开头以两个换行符\n\n作为消息分隔客户端每收到一组完整消息就解析一次。一个最简单的SSE响应长这样HTTP/1.1 200 OK Content-Type: text/event-stream Transfer-Encoding: chunked data: {content:你} data: {content:好} data: [DONE]这个格式背后没有复杂的协议栈不需要WebSocket那套Upgrade握手就是普通HTTP长连接。浏览器端原生支持EventSource对象服务端用Spring的SseEmitter或者直接操作Servlet的OutputStream都能输出。正是因为它门槛低、标准明确大模型服务端推送增量内容时几乎都选它。有一点要特别提醒SSE的data:字段也可以是多行的规范允许一条消息由多个data:行拼接中间用单个\n连接以\n\n结束。做解析的时候不能想当然认为“一行就是一个消息”。好在大模型厂商的流式接口普遍只用单行data:推送但如果你接的是自研推理服务就得多留个心眼。1.2 对比WebSocketAI场景为什么选SSE很多人一听说“服务端推送”就想到WebSocket但在AI流式输出这个场景里SSE几乎是最优解。我整理了一张对比表方便直接看对比项SSEWebSocket协议基础普通HTTP无需协议升级独立协议需要Upgrade握手数据方向服务端单向推送客户端另发请求交互双向实时通信自动重连原生支持浏览器自动带Last-Event-ID重连需要自己实现断线重连服务端实现成本低Spring SseEmitter即可高需要维护连接状态和帧解析中间链路兼容性网关、代理、CDN天然兼容常需要特殊放行和粘滞会话AI流式输出场景完全匹配单向token推送能做但明显过重AI场景选择SSE有个很务实的理由中间链路几乎不需要适配。在真实生产环境里客户端与后端之间往往隔着网关、负载均衡器、反向代理SSE基于普通HTTP这些中间件天然支持不需要额外开启特殊协议能力。我在实际项目中甚至能直接通过普通HTTP代理把SSE流一路转发出去这在WebSocket部署时经常被各种超时和握手问题卡住。还有一个产品层面的原因大模型生成是token级别的几十毫秒就能出一个token如果等全部生成完再一次性返回用户面对的就是十几秒的“转圈”体验极其糟糕。SSE能把首字延迟从十几秒压到几百毫秒配合前端逐字渲染打字机效果就出来了。这已经是对话式AI产品的标配体验。2. 第一层演进用Java HttpClient显式调用SSE接口2.1 显式调用的最小实现最早接大模型接口时Java这边能用的是java.net.http.HttpClient。JDK 11引入的HttpClient有一个很实用的特性可以用BodyHandlers.ofLines()把响应体映射成一个逐行读取的StreamString。这对SSE这种按行推送文本的协议非常契合。基础实现长这样HttpClient httpClient HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(30)) .build(); HttpRequest request HttpRequest.newBuilder() .uri(URI.create(https://api.example.com/v1/chat/completions)) .header(Authorization, Bearer sk-xxxx) .header(Content-Type, application/json) .header(Accept, text/event-stream) .POST(BodyPublishers.ofString( { model: gpt-4o, stream: true, messages: [ {role: user, content: 用一句话介绍SSE} ] } )) .build(); HttpResponseStreamString response httpClient.send(request, HttpResponse.BodyHandlers.ofLines()); response.body() .filter(line - line.startsWith(data:)) .map(line - line.substring(5).trim()) .forEach(line - { if ([DONE].equals(line)) { System.out.println(流结束); } else { JsonObject json JsonParser.parseString(line).getAsJsonObject(); String delta json.getAsJsonArray(choices) .get(0).getAsJsonObject() .getAsJsonObject(delta) .get(content).getAsString(); System.out.print(delta); } });代码本身不难但这个方案里全是细节。先说data:前缀的处理。substring(5)之后最好再trim()一下因为有的网关会在冒号后面补一个空格不处理的话你拿到的字符串带着前导空格JSON解析直接报错。这不是臆想我真实遇到过而且排查了很久。还有一个很多人会踩的坑httpClient.send()在服务端返回非2xx状态码时并不会抛异常只有连接建立失败比如DNS解析错误、连接被拒才会抛IOException。也就是说如果大模型网关返回了429限流或503超载你照样会进入下面的流处理逻辑。所以封装时第一件事就是检查response.statusCode()否则你会拿一坨错误页当成SSE去解析然后出现各种莫名其妙的字符。2.2 流解析的三个关键细节第一个是“心跳行”。很多网关和负载均衡器会对空闲连接自动断开所以服务端会定期推一个注释行:开头的行比如: keep-alive或者干脆推空行来保持连接。如果解析逻辑里只认data:开头这些空行和注释行就必须过滤掉。上面代码用了filter把这个处理掉了但如果你是拿BufferedReader.readLine()手动读空行就是事故多发地。第二个是流式响应结束标记。上面代码里的[DONE]是OpenAI协议约定俗成的终止符但不同厂商实现不一样有的是推送data: [DONE]有的是推一个特殊JSON字段然后直接关闭连接还有的根本不推结束标记由客户端读到流末尾判断。封装的时候这个判断必须做成可配置的开关不能硬编码死在业务代码里。第三个是“流必须消费完”。很多人用HttpClient发SSE请求后发现程序不退出其实是StreamString没有被完整读到底底层连接没有释放线程一直被占着。正确做法是确保在所有场景下包括异常路径都让流消费干净或者显式调用response.body().close()释放连接。在下面是实现的封装里我会保证try-finally逻辑完整避免这类僵尸连接。更多细节我在后面的封装部分展开讲。2.3 没有封装的代码痛点排着队来这种显式调用的方案能跑通但用起来非常难受。我总结了三个痛点第一个是样板代码爆炸。一个工程里只要有两个地方要调大模型流式接口你就得把HttpClient构建、请求头设置、流处理逻辑复制粘贴一遍。一旦要调整超时时间、代理配置、日志埋点所有副本都要同步改漏一个就出线上问题。第二个是重连和容错完全裸奔。SSE连接在网络抖动或服务端主动断开时HttpClient不会自动重连。做知识库问答时用户问到一半流断了整个对话流程直接崩掉体验非常差。要恢复只能自己在外层套重试逻辑还得考虑幂等性防止同一段内容被重复写入。第三个是线程模型伤不起。httpClient.send()在拿到响应体之前是阻塞的BodyHandlers.ofLines()返回的Stream底层也是阻塞读取。每一个并发SSE请求就意味着一个线程长期被占在那里等数据。用固定线程池去并发调多个模型线程池大小就是你的并发上限。同样一台8核机器Jetty可以轻松挂几千个连接但业务逻辑如果每个连接占一个平台线程几百个就到顶了再往上就是线程切换地狱。这个问题直接引出了第三阶段的虚拟线程方案后面再细说。3. 第二层演进把SSE流式调用封装成通用组件3.1 从业务反推接口和回调怎么设计痛过之后自然就会想到封装。但封装不能拍脑袋要先想清楚业务侧到底关心什么。我的结论是业务侧只关心三件事——给我输入、告诉我内容的增量是什么、告诉我整个流结束了没。至于怎么建立连接、怎么发请求、怎么解析事件流、怎么处理异常、要不要重连这些细节业务方都不想管。基于这个思考我设计了一个非常轻量的回调接口public interface StreamCallback { void onMessage(String content); void onDone(); void onError(Throwable t); }onMessage里收到的content是已经解析好的增量文本不是原始的SSE帧。这里有个设计取舍为什么不直接把整帧JSON丢给业务方因为业务方根本不关心choices[0].delta.content这种层级结构那是底层协议的内部细节。直接吐出纯文本增量业务方拿来拼接、转发、渲染都省事得多。如果你想暴露更多结构信息比如思考链、token用量可以在框架里额外提供一个带泛型参数的onChunk(T data)回调让高级用户自己parse但默认还是给纯文本。3.2 通用SseClient落地细节有了回调接口核心的调用逻辑就能收拢到一个SseClient类里public class SseClient { private final HttpClient httpClient; public SseClient() { this.httpClient HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(30)) .build(); } public void stream(String url, String apiKey, String requestBody, StreamCallback callback) { HttpRequest request HttpRequest.newBuilder() .uri(URI.create(url)) .header(Authorization, Bearer apiKey) .header(Content-Type, application/json) .header(Accept, text/event-stream) .POST(BodyPublishers.ofString(requestBody)) .build(); try { HttpResponseStreamString response httpClient.send(request, HttpResponse.BodyHandlers.ofLines()); if (response.statusCode() ! 200) { callback.onError(new RuntimeException(SSE请求失败状态码 response.statusCode())); return; } response.body() .filter(line - line.startsWith(data:)) .map(line - line.substring(5).trim()) .forEach(line - { if (isDoneMarker(line)) { callback.onDone(); } else { callback.onMessage(parseContent(line)); } }); } catch (IOException | InterruptedException e) { callback.onError(e); Thread.currentThread().interrupt(); } } }调用方现在只需要写sseClient.stream(url, apiKey, body, new StreamCallback() { Override public void onMessage(String content) { System.out.print(content); } Override public void onDone() { System.out.println(\n--- 生成完成 ---); } Override public void onError(Throwable t) { t.printStackTrace(); } });这种回调式封装最大的收益不在于少写几行代码而在于把“流的生命周期”和“业务逻辑”彻底解耦。回调可以永久地绑定到WebSocket转发、消息队列投递、MongoDB持久化、前端应答通道等不同出口业务代码完全不用改。还有一个容易漏掉的地方parseContent()里最好做空字符串保护。因为大模型流式响应里经常出现delta.content为null或空串的情况比如只返回了role信息或者思考过程但并不产出可见文本直接调getAsString()会拿个空串然后业务侧每收到空串就刷新一次前端渲染白白增加开销。3.3 进阶封装重连、限流与监控封装不能只停留在“把代码收拢”这一步。真正的隐式封装至少还要解决三个进阶问题。第一个是重连。我在封装层里加入了一个重试机制遇到IOException或非2xx状态码时根据配置决定是否重试重试次数和退避间隔做成构造参数。注意重连不能无脑重发相同请求如果业务侧已经消费了一部分流重连可能导致数据重复。稳妥的方案是把请求里的会话ID做成幂等键由服务端去重或者状态允许时重连前通过回调通知业务方“连接中断准备重试”让业务决定是续传还是重置。第二个是背压控制。流式接口的数据是源源不断的如果下游消费速度跟不上要么在内存里堆积要么丢数据。我在封装中会用一个信号量Semaphore控制最大在途请求数而不是简单地把线程池设得很大。对于AI Agent这种需要并发调多个模型的场景建议按模型维度拆出独立额度防止一个慢模型把整个请求池拖垮。第三个是可观测性。SSE流式调用的排障比普通接口难得多因为没有一份独立的完整响应体日志。我的做法是在封装层统一埋点连接建立耗时、首包延迟第一个token到达耗时、总耗时、推送消息条数、总字节数、断流时的最后一条消息内容。这些指标对AI应用的价值极大相当于把大模型的“慢”归因到了网络层、服务端排队层、还是Token生成层。你可以用LongAdder攒计数再用定时任务上报到监控系统整套下来成本很低但排障时能救命。做完这套封装业务侧的调用形态就从“面向HTTP细节”变成了“面向意图”。更进一步你还可以向外暴露一个chatSync方法内部调用chatStream把增量内容用StringBuilder攒起来流结束时返回完整文本。这样同一个底层引擎同时撑起流式对话和普通问答两个产品形态一套代码通吃。这个设计我在实际项目里用了很久效果非常稳。4. 第三层演进虚拟线程让SSE并发量飞跃4.1 平台线程池扛不住SSE长连接封装解决了“代码层的重复”但“线程层的浪费”依然存在。假设你的AI Agent要并发向10个大模型发流式问答每个模型的响应都要推给用户那么在平台线程模型下至少需要10个专用线程如果每个用户在对话过程中有3个并发流1000个在线用户就是3000个线程。平台线程的栈内存默认1MB起步3000个线程光栈内存就超过3GB这还不算上下文切换的开销。这里需要理解平台线程的底层困境。Java传统的平台线程是一比一映射到操作系统线程的创建和销毁都要经过内核成本高。线程池只是缓兵之计它解决的是“线程复用”不是“线程变轻”。对于SSE这种长连接、高阻塞、低计算量的场景理想的线程模型应该是连接数轻松上万每个连接有一个线程在等着读流但等待的时候不能白白占着内核线程。平台线程做不到这一点。你可能说可以用CachedThreadPool线程不够就新建。但大量平台线程在高并发下会有严重的上下文切换开销而且OS层面线程数是有上限的超过几千个之后调度器就开始吃力。我实测过一台8核16G的机器平台线程撑到800个并发时系统负载明显升高接口延迟翻倍。这不是Java的锅是平台线程模型的天花板。4.2 虚拟线程核心原理挂载与卸载Java 21正式发布的虚拟线程Virtual Threads把这件事做对了。虚拟线程由JVM调度不直接映射到操作系统线程而是架设在平台线程载体线程Carrier Thread之上。虚拟线程执行到阻塞操作时JVM会自动把它从载体线程上卸载unmount载体线程转而去执行别的虚拟线程等IO就绪了再把虚拟线程挂载回某个载体线程继续执行。打个比方平台线程是“一人一坑蹲守”虚拟线程是“一个客服同时排队多个客户谁有消息了接谁的电话”。对于SSE这种“等数据”占大头、计算量极轻的场景这个模型几乎完美因为虚拟线程等待网络数据时不需要占用任何线程资源。创建虚拟线程的姿势很简单// 方式一直接建 Thread.ofVirtual().name(sse-thread).start(() - { sseClient.stream(url, apiKey, body, callback); }); // 方式二推荐用虚拟线程执行器 ExecutorService executor Executors.newVirtualThreadPerTaskExecutor();Executors.newVirtualThreadPerTaskExecutor()是实践中用得最多的方式它的语义是“每个任务一个全新的虚拟线程”没有池化、没有排队创建成本在微秒级。你完全可以用它替换传统的固定线程池代码其他部分基本不动。4.3 虚拟线程加持下的SSE实战用虚拟线程重写并发SSE调用的代码变化小到让人怀疑是不是真的起了效果平台线程版本ExecutorService executor Executors.newFixedThreadPool(200); for (String prompt : prompts) { executor.submit(() - sseClient.stream(url, apiKey, buildBody(prompt), callback)); }虚拟线程版本ExecutorService executor Executors.newVirtualThreadPerTaskExecutor(); for (String prompt : prompts) { executor.submit(() - sseClient.stream(url, apiKey, buildBody(prompt), callback)); }表面看只是改了一行但底层逻辑完全变了。1000个SSE流同时建立时底层真正占用的OS线程可能只有几十个绝大多数虚拟线程在等待网络数据时都处于挂起状态不占载体线程。线程数不再是你需要担心的瓶颈你可以放心地为每个请求、每个Agent任务、每个子任务单独开一个虚拟线程代码保持同步顺序语义不需要引入响应式编程那套复杂的背压和订阅机制。一个典型场景是多Agent协作。Agent A要同时问3份文档Agent B要调2个工具API每个子调用都是SSE流式返回。用虚拟线程的话每个子调用直接一个虚拟线程各自阻塞等待谁先返回谁先处理完全不用手写异步编排。这在平台线程时代是不敢想的——你根本不敢为每个子任务分配一个平台线程。实现上要注意一点虚拟线程不是池化的所以不要用ThreadLocal做线程间数据传递更不要指望线程池那种“复用同一个线程所以ThreadLocal能跨任务共享”的行为。JVM只在任务执行期间挂载任务结束就销毁ThreadLocal的清理完全依赖你手写finally。这个细节在长期运行的Agent服务上尤其要重视否则会有内存泄漏风险。4.4 实测对比500路SSE流式请求的数据口说无凭放一组我实测的数据。测试环境是同一台8核16G的开发机本地部署一个模型服务并发向它发起500路SSE流式请求每路请求约生成300个token单路平均耗时约3秒。线程模型并发数总耗时CPU峰值备注固定线程池100线程50015~18秒中等大量任务排队线程全被阻塞读流占住CachedThreadPool5006~8秒偏高线程数到几百后切换开销明显虚拟线程执行器5004秒左右低底层载体线程数十个等待不占用资源平台线程用100个固定线程时500路请求被切成5批左右总耗时被拉到15秒以上。改用虚拟线程之后500路几乎同时发起总耗时降到单个请求的耗时附近而且CPU和内存没有显著上升。后来我把并发开到2000路虚拟线程方案依然稳定平台线程方案在这个量级上已经跑不动了。有一点要说明这组数据针对的是典型的SSE流式场景也就是“等网络数据”时间远大于计算时间。如果你的逻辑里混着大量CPU密集型计算或持锁操作虚拟线程的提升幅度会大打折扣。所以别指望换一个执行器所有场景都起飞它对IO密集型负载才有显著红利。5. SSE接入常见问题与虚拟线程避坑指南5.1 连接被服务端中断Stream disconnected before completion这类报错在流式接口的日志里非常常见网上经常搜到“idle timeout waiting for sse”的报错本质上都是连接空闲超时被中间层切断。遇到“Stream disconnected before completion”时排查顺序应该是先看服务端有没有配置retry字段或心跳机制再看客户端有没有把空行和注释行消费干净最后检查网关和负载均衡器的空闲连接超时参数。我的习惯是在封装层内置一个心跳兜底客户端在空闲超过N秒时主动做一次连接探测或者让服务端在每个SSE块的间隔里插入一个注释行。注释行格式就是:加任意文本客户端解析时直接忽略但连接因为一直在收发数据就不会被中间层判定为空闲。这个方法成本几乎为零但能解决绝大多数无规律断流问题。提示排查SSE连接问题先别急着上抓包工具。直接用curl -N --no-buffer命令行打一次请求看响应头Content-Type是不是text/event-stream、首包什么时候到、流结束符长什么样。这一步能快速区分是服务端根本没发流还是客户端解析层出了问题。5.2 帧粘连与数据乱序有些自研大模型网关的SSE实现不规范可能出现一个心跳行和真正的data内容挤在同一批里发送。解决的办法是不要逐行做业务消费而是先把以\n\n分帧的原始消息提取出来再对帧内做data:前缀解析。伪代码思路是这样的StringBuilder frame new StringBuilder(); for (String rawLine : rawLines) { if (rawLine.isEmpty()) { // 遇到空行说明一个SSE帧结束了 handleFrame(frame.toString()); frame.setLength(0); } else { frame.append(rawLine).append(\n); } }如果服务端同时又开了多个连接推送同一个响应比如按分片并行生成还得在客户端做基于消息序号的乱序缓冲。乱序问题在主流大模型API中不常见但一旦你开始接自研推理框架或私有化部署的服务就有可能撞上。提前在封装层里留好按event id排序的插槽能省掉后期不少事。5.3 虚拟线程的钉扎与ThreadLocal泄漏这里要把“钉扎pin”问题单独拎出来讲。在JDK 21中虚拟线程进入synchronized同步块后执行阻塞操作会被钉扎在当前载体线程上导致载体线程被这个虚拟线程长期占用。如果你在SSE调用链里用了synchronized并且刚好在锁内触发了阻塞网络读那么大量虚拟线程会把载体线程池耗光性能打回平台线程的原形。解决办法有两个一是把synchronized换成ReentrantLockJDK 21的虚拟线程对Lock接口的阻塞是支持卸载的二是确保synchronized块里只做快速操作把真正的阻塞IO移到锁外。JDK 24通过JEP 491改进了虚拟线程的细粒度载体线程释放但生产环境还在用21的话还是要主动规避。另一个容易被忽视的是ThreadLocal。虚拟线程数量可以非常多如果每个虚拟线程都往ThreadLocal写大对象累积起来就是严重的堆内存泄漏。处理SSE流式解析时状态更应该跟着回调对象走而不是依赖线程局部变量。实在要用必须try-finally里清掉。5.4 背压与流量控制虚拟线程带来的新问题虚拟线程让“创建大量并发”变得几乎无成本但下游服务不是无成本扩容。你一次拉起10万个SSE连接目标服务可能直接被打挂。所以并发控制要从“线程池大小”思维转换成“信号量额度”思维。我现在一律在封装层用Semaphore做全局限流配合按模型或租户拆分的独立额度使用public class SseThrottler { private final Semaphore globalLimit; private final MapString, Semaphore perModelLimits; public boolean tryAcquire(String model) { boolean globalOk globalLimit.tryAcquire(); boolean modelOk perModelLimits.get(model).tryAcquire(); if (globalOk) { if (modelOk) { return true; } globalLimit.release(); } return false; } public void release(String model) { globalLimit.release(); perModelLimits.get(model).release(); } }这个经验对AI Agent场景尤其重要因为Agent框架经常会并行调用多个工具和模型不加控制很容易形成流量脉冲瞬间把网关或模型服务拖垮。信号量方案的优点是它在任何线程模型下都成立不会因为换了虚拟线程而失效。6. 几条实战心得先放在这里梳理一下我在这个演进过程中沉淀下来的经验不保证对所有人通用但都是实际踩出来的。第一SSE封装一定要放在独立的底层模块里不要混进具体的Agent业务代码。我第一次做这个就是把SSE逻辑直接写在ChatService里结果后面要接第二个模型、要加重连、要加监控全部得动业务代码。后来重构出独立的SseClient这些能力才开始对业务透明。第二对AI流式接口首包延迟比总时延更影响用户体验。用户在打字机效果下第一个字能在100ms内出现后面慢一点也能接受首字等3秒后面再快也很难救回来。所以我在封装层里固定埋了首包延迟指标服务端选型时也优先选支持流式、首包快的模型服务。第三虚拟线程确实适合SSE但它不适合所有IO。如果你的场景是短平快的CPU密集型计算或者对虚拟线程调度开销敏感的超低延迟场景传统平台线程池反而更稳。对AI应用对接来说SSE流式调用是典型的IO密集型负载虚拟线程的收益是实打实的。第四建议把SSE、WebSocket、MQTT这类推送通道抽象成统一的“消息流”接口上层业务不要感知底层走的是SSE还是WebSocket。我在一个多Agent协作系统里试过这个思路效果很好同一个回复渲染器既消费SSE推送的token增量也消费WebSocket推来的推理引擎事件切换通道只是换一个适配器业务层完全不动。最后再分享一个小技巧排查SSE连接问题时多用命令行工具先做隔离测试再进代码排查。curl -N --no-buffer配合time命令能看到建立连接时间、首包到达时间、总传输时间这组数据能帮你快速聚焦问题到底在网络层还是在解析层。我自己靠这个习惯省下了大量排障时间也推荐你试试。
阅读完成 · 觉得有帮助?
咨询建站