ARTICLE DETAIL

资讯详情

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

Java+AI 流式输出:SSE 从手写解析到虚拟线程优化全攻略

Java+AI 流式输出:SSE 从手写解析到虚拟线程优化全攻略 做 AI 应用接入大模型流式输出只要你碰的是 Java 后端就一定绕不开 SSEServer-Sent Events。我第一次接这类需求时以为这只是个“普通接口多等几秒而已”结果项目从手写 HTTP 流式解析到后面接 Spring 的封装再到现在用虚拟线程撑起大批量长连接踩的坑比预想中多得多。这篇文章就把这条路径完整复盘一遍聊聊 JavaAI 场景下 SSE 从显式调用到隐式封装、再到虚拟线程性能优化的演进过程。这篇文章适合正在做 AI 应用开发的 Java 工程师也适合后端技术负责人评估流式接口改造方案。你会看到协议层细节、代码示例、封装思路、压测心得以及一些常规文档里不会写的问题排查经验。1. 为什么 AI 流式输出场景偏偏选了 SSE1.1 SSE 是一段不会结束的 HTTP 响应很多第一次接触 SSE 的人会下意识把它和“轮询”混在一起但 SSE 的本质是一次完整的 HTTP 请求只是服务端不立刻返回结果而是把连接一直开着持续向客户端推送文本数据。服务端响应的 Content-Type 必须是text/event-stream数据格式非常朴素每一行是一个字段字段和字段之间用空行分隔成一个事件块。一个典型的事件块长这样data: {token:你} data: {token:好}客户端收到的是普通 HTTP 响应体服务端通过不断追加这个体来持续传输数据。这种设计让人很容易上手不需要像 WebSocket 那样先走一通握手协议、升级连接SSE 只是“一个响应迟迟没读完”的普通请求。SSE 协议里定义了五种字段data表示数据内容event表示自定义事件名id表示事件编号retry表示断线重连的间隔时间以及用冒号开头的注释行。注释行看起来没用但它是服务端用来“刷心跳”的常用手段——发一行: ping就能让连接保持活跃又不会污染业务数据。1.2 和 WebSocket 对比后的选型逻辑很多人会问流式推送为什么不用 WebSocketWebSocket 确实更全能但 AI 对话这个场景实际需要的只是一条“服务器单向推到客户端”的下行通道。对比维度SSEWebSocket连接模型普通 HTTP 长响应TCP 专用长连接协议复杂度低基于文本行较高需要帧解析和掩码握手方式无需额外握手需要 101 升级握手数据流向单向下行双向自动重连内置基于 Last-Event-ID需要自己实现调试成本浏览器控制台直接看 Network 流需要专门的调试面板或工具AI 流式场景完全匹配功能过剩从工程视角看AI 对话基本是“客户端发一个问题服务端返回一段流式文本”几乎没有实时上行需求。用 WebSocket 服务端要维护会话状态、处理心跳帧、考虑消息分片而 SSE 只需要把响应写好就行。还有一个现实原因大模型厂商的服务端接口大多直接输出 SSE 格式Java 后端要做的不是发明协议而是接住上游的流再原样或加工后推给下游。这个单向链路用 SSE 最顺手。1.3 AI 应用里的完整 SSE 调用链路在一个典型的 JavaAI 项目中SSE 会贯穿两层链路。用户在前端点“发送”前端用EventSource发起请求Java 后端收到这个请求后再去调用大模型接口。大模型接口通常也是 SSE 流式返回Java 后端一边读上游 token一边通过自己的 SSE 连接把 token 推给前端。前端每收到一块新数据就把它渲染到对话框里形成“打字机”效果。这中间要处理三件事上游模型的流式数据解析、后端到前端的协议转换、以及两段连接的异常处理。很多人只关注“调用 API”忽略了传输链路上的断连和超时后面会详细讲。2. 显式调用手写 SSE 客户端的那段日子2.1 早期实现HttpClient BufferedReader 手动读流最早我写 SSE 客户端完全是“硬读”。Java 自带java.net.http.HttpClient支持把响应体作为InputStream接收然后就能用BufferedReader一行行往下读。代码看起来很简单HttpClient client HttpClient.newBuilder() .connectTimeout(Duration.ofSeconds(10)) .build(); HttpRequest request HttpRequest.newBuilder() .uri(URI.create(http://localhost:8080/v1/chat/completions)) .header(Content-Type, application/json) .header(Accept, text/event-stream) .POST(BodyPublishers.ofString( {model:demo,prompt:你好,stream:true} )) .build(); HttpResponseInputStream response client.send(request, HttpResponse.BodyHandlers.ofInputStream()); if (response.statusCode() ! 200) { // 这里要先把错误体完整读取出来记录再抛异常 } try (BufferedReader reader new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; while ((line reader.readLine()) ! null) { if (line.startsWith(data:)) { String payload line.substring(data:.length()).trim(); if ([DONE].equals(payload)) { break; } // 用 Jackson 或 Gson 把 payload 解析成 JSON // 从 choices[0].delta.content 里取出增量 token } } }这段代码放到真实环境里能跑通 demo但它真的太“显式”了——所有协议细节都暴露在业务代码里所有异常都要业务侧自己兜底。2.2 显式调用最容易忽略的协议细节SSE 的事件块解析并不是“读一行就完事”这么简单。按协议规范一个事件块里可能出现多行data多行内容会被拼接成一段完整数据中间用换行符连接。也就是说下面这种格式是合法的data: {delta:你} data: {delta:好}这段内容会拼成{delta:你}\n{delta:好}解析时不能只取第一行。event、id、retry字段也各有用途。id字段特别重要客户端断线重连时会带Last-Event-ID头服务端可以根据这个 ID 决定从哪条事件之后开始重推。retry字段则是服务端建议客户端下次重连的等待毫秒数。我后来写了一个通用解析方法负责把一个事件块的多行内容转换成结构化对象static SseEvent parseEventBlock(ListString lines) { String data ; String event message; String id null; Integer retry null; for (String raw : lines) { String line raw.endsWith(\r) ? raw.substring(0, raw.length() - 1) : raw; if (line.isEmpty() || line.startsWith(:)) { continue; } int sep line.indexOf(:); String field sep 0 ? line : line.substring(0, sep); String value sep 0 ? : line.substring(sep 1); if (value.startsWith( )) { value value.substring(1); } switch (field) { case data - data data.isEmpty() ? value : data \n value; case event - event value; case id - id value; case retry - retry Integer.valueOf(value); } } return new SseEvent(id, event, data, retry); }这套解析逻辑看着繁琐但一旦接入多家大模型供应商就会发现只处理data字段的代码完全不够用有的是事件名不同有的是断点续传需要id。2.3 显式调用阶段踩过的坑这一阶段最大的体会是SSE 客户端最难的从来不是“怎么读”而是“读的过程中连接挂了怎么办”。请求超时不能瞎设。很多人在HttpRequest上加.timeout(Duration.ofSeconds(30))然后发现流式响应一到 30 秒就被客户端主动断掉。因为 HttpClient 的 timeout 是“整个请求从开始到结束”的总超时而 SSE 恰恰是一个“持续时间很长”的请求。正确做法是把 connectTimeout 单独设置好不在请求级别设置总超时超时控制交给专门的空闲读超时逻辑。必须处理非 200 状态。模型服务如果鉴权失败、限流或参数错误会直接返回 4xx响应体里是一个普通 JSON 错误信息根本不是 event-stream 格式。流式解析代码遇到这种混入的响应直接懵。所以进入读流之前先检查状态码把错误体完整读出并记录。断线重连和心跳是刚需。早期我没有做重连逻辑结果一次网络抖动就让用户看到“生成中断”而且没有恢复机制。后来才理解 SSE 协议自带 Last-Event-ID 重连机制就是给这个场景用的客户端维护一个最后处理的事件 ID重连时带上去服务端决定要不要补发。背压很容易被忽略。如果上游生成速度快下游消费者处理不过来BufferedReader.readLine()会持续读到新数据内存里的消息队列就会膨胀。尤其在做 Agent 场景时中间的日志、状态更新、工具调用结果都要处理这一段消费链路要设计好缓冲和丢弃策略不能无脑 while 循环。2.4 显式调用留下的实际价值虽然显式调用代码最繁琐但我不建议新手一上来就套封装框架。因为后续排查线上问题比如“为什么读到一半断了”“为什么这个事件没触发”最终都要回到协议的原始格式去看。亲自动手实现一遍 SSE 解析之后再看各种封装库的源码就会清楚它内部到底在做什么。3. 隐式封装SSE 从手写到开箱即用3.1 服务端发送SseEmitter 和 WebFlux 怎么选Java 后端给前端推 SSE如果是 Spring MVC 项目最直接的方式是返回SseEmitter。PostMapping(/chat) public SseEmitter chat(RequestBody ChatRequest request) { SseEmitter emitter new SseEmitter(0L); // 不设超时由业务控制何时完成 modelCallExecutor.execute(() - { try { // 模拟模型逐步返回 emitter.send(SseEmitter.event() .id(1) .name(token) .data(你好)); emitter.send(SseEmitter.event() .id(2) .name(token) .data(我是 AI)); emitter.complete(); } catch (Exception ex) { emitter.completeWithError(ex); } }); return emitter; }SseEmitter.event()返回一个事件构建器id、name、data正好对应 SSE 协议里的字段底层会帮你拼成text/event-stream响应。这里emit线程的选择值得注意模型调用如果是阻塞式 HttpClient 请求就不要占用 Tomcat 的请求工作线程否则长连接一多容器线程很快被占满。这也是后面虚拟线程切入的入口。如果你已经用了 Spring WebFlux 的响应式栈可以直接返回FluxServerSentEventGetMapping(value /chat, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxServerSentEventString chat() { return Flux.interval(Duration.ofSeconds(1)) .map(i - ServerSentEvent.builder(token i) .id(String.valueOf(i)) .build()); }两种方式的选择建议很简单项目原本就是 Spring MVC用 SseEmitter不要为了 SSE 把整个技术栈换成 WebFlux项目本来就是全链路响应式用 Flux 更自然。强行混用会让团队同时维护两套编程模型代价远超收益。3.2 客户端接收WebClient 和 OkHttp 的封装差异Java 后端调用大模型接口时推荐优先用 Spring WebClient 的bodyToFlux。它对 SSE 做了内置支持几行代码就能把上游数据流变成一个响应式的 Flux。WebClient client WebClient.builder() .baseUrl(http://localhost:8080) .build(); FluxString stream client.post() .uri(/chat) .bodyValue(request) .accept(MediaType.TEXT_EVENT_STREAM) .retrieve() .bodyToFlux(String.class);默认的bodyToFlux(String.class)拿到的是data:字段的内容如果你还想拿到事件的id、event、retry元信息可以用bodyToFlux(ServerSentEvent.class)FluxServerSentEventString events client.post() .uri(/chat) .bodyValue(request) .accept(MediaType.TEXT_EVENT_STREAM) .retrieve() .bodyToFlux(new ParameterizedTypeReferenceServerSentEventString() {});如果你的项目已经大量使用 OkHttp不考虑引入 WebClient可以用okhttp-sse扩展包OkHttpClient okHttpClient new OkHttpClient.Builder().build(); EventSources.createFactory(okHttpClient).newEventSource(request, new EventSourceListener() { Override public void onEvent(EventSource eventSource, String id, String type, String data) { // 每收到一个事件回调一次 } Override public void onClosed(EventSource eventSource) { // 连接正常关闭 } Override public void onFailure(EventSource eventSource, Throwable t, Response response) { // 连接意外失败 } });OkHttp 的封装把断线重连逻辑做了一部分但它只负责客户端侧服务端实现要自己写。WebClient 的优势在于和 Spring 生态的响应式链路融合得更紧密适合“上游到下游全程流式”的架构。3.3 把流式接口统一抽象成 StreamResponse 模式工程里接入多个模型供应商后会发现每家模型的流式协议大体相同但细节各有差异。有的返回[DONE]有的返回data: [DONE]有的 token 在choices[0].delta.content里有的在data.choices[0].message.content里。这个时候最该做的不是到处写 if-else而是抽象一个统一回调接口public interface StreamCallback { default void onStart(StreamContext context) {} default void onToken(StreamContext context, String token) {} default void onFinish(StreamContext context) {} default void onError(StreamContext context, Throwable throwable) {} }每个供应商实现一个适配模块内部负责把各自的 SSE 事件解析成统一的onToken回调。业务层只需要关心“我收到了哪个字符串”不需要关心它是 OpenAI 的格式还是国产模型的自定义格式。这种封装的本质是把我前面手写的那套协议解析、断线重连、超时控制、背压处理逻辑沉淀成公共组件。它不只是省代码量更关键的是让所有接入方都拿到一致性的可靠性保障。4. 虚拟线程SSE 并发量的真实拐点4.1 平台线程池和 SSE 长连接是一个天然矛盾传统 Tomcat 默认请求线程池是 200 个平台线程SSE 却是“一个连接长时间占着一个线程直到断开”的模型。如果同时有 200 个用户挂着流式对话请求线程池被占满其他普通接口全部排队。这个问题的本质是等待大模型返回期间线程完全阻塞在 IO 上不干任何事但仍然占用内存和调度资源。一个平台线程的默认栈内存大约 1MB启动 5000 个线程就近乎 5GB 的内存开销这还没算操作系统上下文切换的成本。前面提到的方式无论是 NIO 还是响应式 WebClient都是非阻塞方案的变体。但这条路有两个代价一是代码风格从同步改造成链式回调整个团队要重新学习二是与现有 Spring MVC、MyBatis、事务管理等阻塞式技术栈耦合时非常别扭。4.2 虚拟线程怎么解决“阻塞等待”问题虚拟线程是 JDK 21 正式提供的能力它和平台线程最大的区别是平台线程是操作系统调度的虚拟线程是 JVM 内部调度的。JVM 会创建少数平台线程作为载体线程池虚拟线程运行在载体线程上。当虚拟线程执行到阻塞点比如读取网络流、Thread.sleep、LockSupport.parkJVM 会自动把这个虚拟线程从载体线程上摘下来让载体线程去执行另一个虚拟线程。整个过程由 JVM 调度器完成无需业务代码介入。生活化理解平台线程像固定工位一个人坐在工位上等快递工位就浪费了虚拟线程像临时工快递没到就先干别的活快递到了再回来接着干。SSE 场景里大量连接都在“等 token”正好是虚拟线程最擅长消化的一类负载。4.3 虚拟线程 SSE 的落地改造在 Spring Boot 3.2 及以后版本最简单的开启方式是在配置文件里设一个开关spring: threads: virtual: enabled: true开启后容器接收请求时会用newVirtualThreadPerTaskExecutor()默认创建的虚拟线程执行器处理请求。对 SSE 接口来说意味着每个 SSE 长连接都可以占一个虚拟线程而不是平台线程。如果你的项目不打算全局开启也可以只在流式接口的“模型调用”环节用虚拟线程执行器ExecutorService modelCallExecutor Executors.newVirtualThreadPerTaskExecutor(); PostMapping(/chat) public SseEmitter chat(RequestBody ChatRequest request) { SseEmitter emitter new SseEmitter(0L); modelCallExecutor.execute(() - { try { // 这里内部是阻塞式 HttpClient 调用大模型 modelService.streamChat(request, token - { try { emitter.send(SseEmitter.event().name(token).data(token)); } catch (IOException e) { throw new UncheckedIOException(e); } }); emitter.complete(); } catch (Exception ex) { emitter.completeWithError(ex); } }); return emitter; }虚拟线程的特点是“创建成本低、用完即丢”所以不要用线程池去池化它。每次调用就execute执行完自动销毁即可。我做过一个简单对比测试一台 8 核 16G 的开发机器模拟 300 个 SSE 客户端连接每个连接保持 60 秒服务端每 2 秒向下推一条消息。平台线程池模式下Tomcat 默认线程很快被打满普通接口响应开始出现几秒延迟切到虚拟线程后CPU 占用平稳普通接口响应时间基本没受流式连接影响。当然这个数据只是中午跑的小实验权当参考但虚拟线程在长连接场景的收益方向是明确的。4.4 虚拟线程使用边界和几个隐藏问题虚拟线程不是银弹有几类场景反而要格外小心。synchronized会让虚拟线程 pins 到载体线程上。如果在虚拟线程里锁竞争激烈虚拟线程无法被摘除会直接占用一个平台线程并发效果立刻退化。能用ReentrantLock的场景尽量替换。ThreadLocal 在虚拟线程里虽然能用但开启虚拟线程的 ThreadLocal 复制成本更高。不要在虚拟线程里传重量级上下文对象尤其跨线程池传数据时要重新设计。虚拟线程适合 IO 密集型不适合 CPU 密集型。如果你在流式任务里做大量 JSON 大字段解析、递归计算、正则回溯这些计算不能靠虚拟线程换并发反而会因为线程库增加而带来额外调度开销。另外不要在代码里为了“看起来用到了虚拟线程”而手动创建大量虚拟线程去做无限循环任务。SSE 长连接的本质是 IO 等待但如果连接建立后业务代码本身不阻塞、只拼 CPU 死等虚拟线程帮不上忙。5. 常见问题与排查技巧实录5.1 SSE 问题速查表现象可能原因排查方向快速处理前端 EventSource 不触发 open响应 Content-Type 不是 text/event-stream看 Network 响应头后端显式指定produces TEXT_EVENT_STREAM_VALUE数据不实时刷新一次性返回反向代理层开启了响应缓冲检查网关/Web 代理配置关闭代理缓冲添加X-Accel-Buffering: no响应头连接在固定时间后断开网关或客户端空闲超时看服务端日志、网关日志服务端定期推送: ping注释行心跳收到乱码字符编码不一致检查响应头 charset 字段统一使用 UTF-8显式指定编码断线重连后内容重复没有维护 Last-Event-ID检查客户端是否传重连 ID服务端按 ID 实现断点续推开启虚拟线程后性能下降synchronized pin 或 CPU 密集计算看线程 dump、锁竞争替换锁、拆分CPU密集任务5.2 遇到 stream disconnected before completion 怎么查这个报错本质是 SSE 流没有读完就断开了常见于客户端、网关、服务端三者的超时策略不一致。第一步看客户端。如果用的是 WebClient检查是否设置了全局超时或者连接空闲超时响应式框架有多个超时维度很容易误伤长连接。第二步看网关。Nginx 默认proxy_read_timeout是 60 秒SSE 连接超过这个时间没有新数据网关直接断开。配合服务端心跳注释行并调大proxy_read_timeout或者干脆在特定路径上关闭缓冲。第三步看服务端。服务端如果正常complete()了但客户端以为还该继续收数据也会出现这种报错。排查时要区分是“符合预期的正常关闭”还是“异常中断”日志里记录的堆栈是关键。5.3 断线续传和业务幂等怎么设计最常见的场景是用户提问后模型已经生成了 30 个 token第 31 个 token 发送时断网。重连后如果把整个请求重发一遍模型要重新生成全部 token体验很差还可能造成费用重复。SSE 自带的id字段就是干这个的。服务端持续为每个事件生成递增 ID客户端重连时带上Last-Event-ID服务端从该 ID 之后的事件继续推送。但这里有个复杂点如果模型服务侧没有缓存中间生成结果服务端拿到Last-Event-ID也无从恢复。所以更通用的做法是服务端提前把 token 写入本地缓存或消息队列重连时从缓存里补发。如果模型服务不支持断点续传至少要在业务层做幂等用一个业务请求 ID 标识整轮对话重连时服务端检查是否已有部分结果要么丢弃重来要么从最近一次成功写入的 checkPoint 继续具体取舍看成本和产品体验要求。5.4 把“流式”当作一等公民来建设经历了手写、封装、虚拟线程三个阶段后我发现流式接口的工程化不能只停留在“能通”层面要当作独立的传输基础设施来建设。监控层面要单独统计 SSE 连接数、SseEmitter 完成数、超时数、异常断开数。连接数突增可能意味着有人在刷接口错误数突增往往和上游模型服务质量波动直接相关。日志层面不要把所有 token 都打出来完整对话内容体量太大只记录“第几条事件、多少字节、耗时多久”即可。存储层面流式数据是增量到达的要设计好从临时缓冲到最终落库的路径不能等全部生成完才写库。这些细节决定了流式功能上线后是“能看”还是“真的好用”。最后再分享一点个人经验我最早写 SSE 客户端时也很自信觉得协议看着简单几行代码就能搞定结果一碰真实网络环境就出各种问题。现在再看SSE 真正的复杂度不在协议本身而在长连接的可靠性和并发承载。对准备入场的团队我的建议是第一步先把 text/event-stream 的协议格式吃透亲手解析一次事件流第二步再考虑引入框架封装用 WebClient 或 okhttp-sse 把脏话细节挡在业务层之外第三步才是上虚拟线程在长连接并发量真正起来之后做优化。顺序不能反反了会在排查问题时无从下手。虚拟线程和 SSE 确实是目前 JavaAI 组合里很舒服的一对搭档。同步代码风格不变却能享受到接近响应式架构的高并发红利。不过要记住虚拟线程解决的是 IO 等待不是 CPU 计算选型时把这层账面算清楚后面才不会踩坑。
返回列表