新闻详情

新闻详情

首页 / 资讯中心 / 详情

Java虚拟线程+SSE实现大模型流式输出:从显式调用到高并发封装

发布时间:2026/9/29 8:46:29来源:尧图网络
Java虚拟线程+SSE实现大模型流式输出:从显式调用到高并发封装
1. 为什么大模型流式输出绕不开 SSE做过 AI 对话类产品的后端同学应该都有体会用户点下发送按钮之后如果界面要等模型把整段回答全部生成完才一次性渲染那种体验基本没法用。大模型生成一段几百字的回答耗时可能从两三秒到十几秒不等用户盯着一个转圈的加载图标干等心理上早就判定这个产品卡死了。所以流式输出不是锦上添花的功能而是 AI 交互产品的及格线。流式输出这件事落到技术选型上绕不开SSEServer-Sent Events服务器推送事件。它本质上是基于 HTTP 的一种单向、长连接、文本协议服务端可以持续往客户端推送text/event-stream格式的数据块客户端用浏览器原生的EventSource或者fetch的流式读取就能接。相比 WebSocketSSE 的优势在于走标准 HTTP不需要协议升级握手天然兼容现有的网关、鉴权、日志体系服务端实现也简单得多。对于用户提问、模型逐字回答这种典型的单向推送场景SSE 几乎是性价比最高的选择。但真到了 Java 后端落地的时候问题就来了。Java 传统的 Servlet 线程模型是一个请求占一个线程而 SSE 连接是长连接一个用户开着对话页面这条连接可能挂几分钟甚至更久。如果每个 SSE 连接都独占一个平台线程几百个并发在线用户就能把线程池打满后面的请求全部排队。这就是为什么JDK21引入的虚拟线程Virtual Threads在这个场景下特别香——它让一个连接一个线程这种最符合直觉的写法重新变得可行而不用被迫去写复杂的响应式代码。这篇内容我打算把这条链路完整走一遍从最原始的显式 SSE 调用到把 AI 交互逻辑封装成隐式调用再到用虚拟线程把并发能力拉起来。中间会讲清楚每一步为什么这么设计、踩过哪些坑、参数怎么算。适合已经写过 Java Web、想上手 AI 流式接口的同学也适合正在做 AI 应用但被并发问题卡住的老手。2. 显式 SSE 调用先把最原始的链路跑通2.1 显式调用的本质是什么所谓显式调用指的是开发者手动去管理 SSE 的整个生命周期手动设置响应头、手动拿输出流、手动按 SSE 协议格式拼数据、手动 flush、手动处理连接断开。这种方式最贴近协议本身也最能暴露问题所以我建议任何人在用框架封装之前都先手写一遍。在 Spring MVC 里一个最朴素的 SSE 接口大概长这样GetMapping(value /chat/stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public void stream(RequestParam String question, HttpServletResponse response) throws IOException { response.setContentType(text/event-stream); response.setCharacterEncoding(UTF-8); response.setHeader(Cache-Control, no-cache); response.setHeader(Connection, keep-alive); PrintWriter writer response.getWriter(); // 调用大模型拿到流式返回 modelClient.stream(question, chunk - { writer.write(data: chunk \n\n); writer.flush(); }); writer.write(data: [DONE]\n\n); writer.flush(); }这段代码看着简单但每一行都有讲究。produces TEXT_EVENT_STREAM_VALUE决定了响应类型Cache-Control: no-cache是防止中间代理缓存住流式内容导致用户看不到实时输出Connection: keep-alive是告诉连接别急着关。数据格式上SSE 规定每条消息以data:开头以两个换行\n\n结尾这是协议硬性要求少一个换行客户端就解析不出来。2.2 显式调用最容易踩的三个坑第一个坑是缓冲区。很多人写完上面那段代码发现明明服务端在 flush前端却要等好几秒才一次性收到全部内容。原因通常是中间有反向代理或者 Servlet 容器自带的缓冲。Nginx 需要加proxy_buffering off;Tomcat 层面要确认响应没有被Content-Length提前锁定。我一般会在响应头里再加一个X-Accel-Buffering: no专门用来关掉 Nginx 的缓冲这个头是 Nginx 认的实测很管用。第二个坑是连接断开后的资源释放。用户关掉页面、切走标签页SSE 连接就断了但服务端可能还在傻乎乎地调模型、往一个已经死掉的流里写数据。这时候writer.write会抛IOException如果不捕获线程就带着异常结束了模型调用那边的资源也没释放。正确做法是监听连接状态一旦写失败就主动中断模型调用。这里就引出了后面要讲的abort机制。第三个坑是超时。SSE 长连接如果长时间没有数据推送中间的网络设备负载均衡、网关可能会判定连接空闲然后掐断前端就会看到stream disconnected before completion: idle timeout waiting for sse这类报错。解决办法是定期发送心跳注释行比如每隔 15 秒发一个: heartbeat\n\n冒号开头的行在 SSE 协议里是注释客户端会忽略但能保持连接活跃。提示心跳间隔不要设太短15 到 30 秒比较合适。太短会浪费带宽太长又起不到保活作用。具体值要看你链路上最急性子的那个设备的空闲超时配置。2.3 显式调用的价值与局限显式调用最大的价值是透明。你能清楚看到数据从模型出来、经过你的代码、写到 HTTP 响应的全过程出问题的时候排查路径非常清晰。但它的局限也很明显每个接口都要重复写一遍响应头设置、格式拼接、异常处理、心跳逻辑代码重复度高而且一旦要换模型供应商改动面很大。这就自然过渡到了下一个阶段——把这些重复的、和业务无关的 SSE 细节封装起来让业务代码只关心我要推什么内容而不是怎么推。3. 隐式封装让 AI 交互逻辑对业务透明3.1 封装的核心思路隐式封装的目标是业务层写起来就像调用一个普通的返回字符串的方法但底层自动变成流式推送。要做到这一点核心是把 SSE 的推抽象成一个回调或者一个流式返回值业务代码只管往里塞数据。我比较推荐的封装方式是定义一个StreamSink接口业务侧只依赖它public interface StreamSink { void send(String chunk); void complete(); void error(Throwable t); boolean isClosed(); }然后框架层负责把StreamSink和真实的 HTTP 响应绑定起来处理格式拼接、flush、心跳、异常。业务代码变成这样public void handle(String question, StreamSink sink) { modelClient.stream(question, chunk - { if (sink.isClosed()) return; // 连接已断别再推了 sink.send(chunk); }); sink.complete(); }这么一改业务代码里再也看不到data:、\n\n、flush这些协议细节。换模型供应商的时候只要新的客户端也支持流式回调业务代码一行不用动。3.2 基于什么技术栈封装 AI 交互逻辑热词里有个问题问得很实在基于什么技术栈封装 AI 交互逻辑。我的答案是分层的Web 层用 Spring MVC 或者 WebFlux模型调用层用各家 SDK 或者统一的 HTTP 客户端中间加一层适配器把不同模型的流式协议统一成StreamSink。为什么不直接用 WebFlux 的Flux一路到底因为响应式编程的学习曲线陡团队里只要有人不熟维护成本就上去了。而且很多模型 SDK 本身是阻塞式的硬塞进响应式链路里反而别扭。用虚拟线程之后阻塞式写法 高并发这个组合重新成立封装的复杂度能降一大截。适配器这一层是关键。不同模型供应商的流式返回格式不一样有的返回 JSON 行有的返回纯文本块有的带[DONE]结束标记有的靠连接关闭表示结束。适配器的职责就是把这些差异吃掉对外只暴露统一的send/complete/error。我一般会为每个供应商写一个ModelStreamAdapter注册到一个工厂里根据配置选择。3.3 abort 机制用户点了停止按钮之后AI 对话产品几乎都有停止生成按钮。用户点下去前端会中断 SSE 连接后端必须感知到并且立刻停止模型调用否则你还在为一次用户已经放弃的请求付费这在按 token 计费的场景下是真金白银的浪费。感知连接断开的方式在 Servlet 里可以监听AsyncListener的onError和onTimeout或者更简单粗暴地在每次send时检查sink.isClosed()。但检查是被动的如果模型那边半天不吐一个字你也没机会检查。所以更可靠的做法是把模型调用的Future或者取消句柄存起来一旦检测到连接断开就主动cancel。AtomicReferenceCancellable current new AtomicReference(); sink.onClose(() - { Cancellable c current.get(); if (c ! null) c.cancel(); }); Cancellable c modelClient.stream(question, chunk - sink.send(chunk)); current.set(c);这里有个细节onClose的回调触发时机和模型调用的启动时机之间可能有竞态。如果连接在模型调用启动之前就断了current还是空的取消就落空了。所以要么用锁保护要么在设置current之后立刻再检查一次连接状态。这个坑我在实际项目里踩过表现是偶尔有请求明明用户已经关了页面后端还在跑查了半天才发现是竞态。3.4 封装带来的可测试性提升封装还有一个容易被忽略的好处可测试。显式调用的时候你要测一个 SSE 接口得启动整个 Web 容器模拟 HTTP 连接很重。封装成StreamSink之后你可以写一个CollectingSink把推送的内容收集到列表里然后直接对业务方法做单元测试断言推送的内容和顺序。模型客户端也可以 mock 成一个按预设序列吐 chunk 的假实现。这样流式逻辑的正确性就能在毫秒级跑完不用等真实模型。4. 虚拟线程让一连接一线程重新可行4.1 平台线程模型在 SSE 场景下的死结传统 Java 里一个请求进来Tomcat 从线程池里分配一个平台线程这个线程从头到尾跟着请求走。SSE 是长连接意味着这个线程要被占用很久。Tomcat 默认最大线程数 200也就是说理论上最多 200 个并发 SSE 连接超出的请求全部排队。而 AI 对话场景下用户开着页面不关连接可能挂几分钟200 个并发对稍微有点规模的产品来说完全不够看。有人会说那就把线程池调大。但平台线程是操作系统线程一个线程占 1MB 左右的栈空间几千个线程就是几个 G 的内存而且线程上下文切换的开销也会急剧上升。这条路走不通。另一条路是响应式用少量线程处理大量连接。但前面说了响应式代码写起来累和阻塞式的模型 SDK 配合也别扭。4.2 虚拟线程为什么能解这个结虚拟线程是 JDK21 正式转正的特性JEP 444。它的核心思想是把线程的调度从操作系统层面搬到 JVM 层面。一个虚拟线程在遇到阻塞操作比如 IO 等待时JVM 会把它挂起让底层的载体线程carrier thread去跑别的虚拟线程。等阻塞结束再把虚拟线程恢复挂到某个载体线程上继续跑。关键在于虚拟线程的栈是存在堆里的可以按需增长创建成本极低。你可以轻松创建几十万个虚拟线程内存占用远小于同等数量的平台线程。对于 SSE 这种大部分时间在等待的场景虚拟线程简直是量身定做——每个连接一个虚拟线程等待模型返回的时候虚拟线程被挂起载体线程去服务其他连接资源利用率拉满。在 JDK21 里启用虚拟线程非常简单Spring Boot 3.2 之后可以直接配置Bean public TomcatProtocolHandlerCustomizer? protocolHandlerCustomizer() { return protocolHandler - { protocolHandler.setExecutor(Executors.newVirtualThreadPerTaskExecutor()); }; }或者更简单在配置文件里加一行spring.threads.virtual.enabledtrue。这一行下去Tomcat 处理请求的线程就全变成虚拟线程了。4.3 虚拟线程不是银弹两个必须注意的点第一个注意点是synchronized 的 pinning 问题。在 JDK21 早期版本里如果一个虚拟线程在synchronized块里阻塞它会把载体线程一起钉住pinning导致载体线程没法去跑别的虚拟线程虚拟线程的优势就没了。JDK21 里这个问题在大部分场景已经缓解但如果你在 SSE 链路里大量用synchronized保护共享状态还是要留个心眼。替代方案是用ReentrantLock它和虚拟线程配合得更好。第二个注意点是线程本地变量ThreadLocal。虚拟线程数量巨大如果每个虚拟线程都塞一堆 ThreadLocal 数据内存会涨得很快。SSE 链路里常见的 ThreadLocal 用途是传递用户上下文、traceId 之类建议改用 JDK21 的ScopedValue预览特性或者显式传参别滥用 ThreadLocal。注意虚拟线程适合 IO 密集型任务不适合 CPU 密集型。如果你的 SSE 链路里夹着大量 JSON 序列化、加解密这类 CPU 操作虚拟线程帮不上忙该优化算法还得优化算法。4.4 实测并发能力对比我在一台 4 核 8G 的测试机上做过对比。同样的 SSE 接口模拟客户端持续接收流式数据每个连接保持 60 秒。线程模型最大稳定并发连接内存占用备注平台线程池200约 200约 1.2G超过后请求排队平台线程池1000约 900约 3.5G上下文切换开销明显虚拟线程5000约 1.5G未测到上限受限于文件描述符这个数据不是绝对的跟具体业务逻辑、模型响应速度都有关系但趋势很清楚虚拟线程把并发上限从线程池大小这个硬约束里解放出来了。真正需要关注的上限变成了文件描述符数量、网络带宽这些更底层的资源。5. 完整实操从零搭一个虚拟线程 SSE 服务5.1 环境准备与 JDK21 安装要点先把 JDK21 装上。Linux 服务器上我一般用包管理器或者直接下压缩包。下压缩包的话解压到/usr/local/jdk-21然后配置环境变量export JAVA_HOME/usr/local/jdk-21 export PATH$JAVA_HOME/bin:$PATH写到/etc/profile.d/jdk21.sh里source一下生效。验证用java -version看到 21 就对了。这里提醒一句如果你服务器上还跑着依赖老版本 JDK 的应用别直接把系统默认 JDK 换掉用update-alternatives管理多版本更稳妥。Spring Boot 版本要选 3.2 及以上它才正式支持spring.threads.virtual.enabled。Maven 里把 parent 版本提上去Java 编译版本设成 21。5.2 核心代码结构整个服务我分成四层Controller 层负责接收请求和建立 SSE 连接Service 层是业务逻辑Adapter 层对接模型Sink 层处理 SSE 输出。目录结构大概这样com.example.chat ├── controller/ChatController.java ├── service/ChatService.java ├── adapter/ModelStreamAdapter.java ├── adapter/impl/MockModelAdapter.java └── stream/StreamSink.javaController 层GetMapping(value /chat, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter chat(RequestParam String q) { SseEmitter emitter new SseEmitter(0L); // 0 表示不超时 chatService.handle(q, new SseEmitterSink(emitter)); return emitter; }这里用SseEmitter是 Spring 提供的封装比手写HttpServletResponse省事它内部帮你处理了格式和 flush。超时设成 0 表示永不超时具体要不要设超时看业务一般对话场景我会设个 5 分钟兜底。SseEmitterSink把SseEmitter适配成我们的StreamSinkpublic class SseEmitterSink implements StreamSink { private final SseEmitter emitter; private final AtomicBoolean closed new AtomicBoolean(false); public SseEmitterSink(SseEmitter emitter) { this.emitter emitter; emitter.onCompletion(() - closed.set(true)); emitter.onTimeout(() - closed.set(true)); emitter.onError(e - closed.set(true)); } Override public void send(String chunk) { if (closed.get()) return; try { emitter.send(SseEmitter.event().data(chunk)); } catch (IOException e) { closed.set(true); } } // complete / error / isClosed 略 }5.3 心跳与超时的参数计算心跳间隔怎么定我的经验公式是心跳间隔 链路上最短空闲超时 / 3。比如 Nginx 默认proxy_read_timeout是 60 秒那心跳设 20 秒比较安全。如果链路上有负载均衡设了 30 秒空闲超时那心跳就得压到 10 秒以内。心跳的实现用一个定时任务每隔固定时间往 sink 里发一个空注释。注意心跳也要检查isClosed连接断了就别发了。ScheduledExecutorService heartbeat Executors.newSingleThreadScheduledExecutor(); heartbeat.scheduleAtFixedRate(() - { if (!sink.isClosed()) sink.send(); // 空数据客户端会忽略 }, 15, 15, TimeUnit.SECONDS);任务结束的时候记得把这个定时任务关掉否则会泄漏。我一般把心跳任务的生命周期和 sink 绑定sink 关闭时一起取消。5.4 前端如何消费这条流前端用fetch而不是EventSource因为EventSource只支持 GET而且没法自定义请求头比如带 token。用fetch拿到response.body的ReadableStream配合TextDecoder逐块解析const resp await fetch(/chat?q encodeURIComponent(question)); const reader resp.body.getReader(); const decoder new TextDecoder(); let buffer ; while (true) { const { done, value } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const lines buffer.split(\n\n); buffer lines.pop(); for (const line of lines) { if (line.startsWith(data:)) { const text line.slice(5).trim(); appendToUI(text); } } }这里buffer的处理是关键。网络传输是分块的一个 SSE 消息可能被拆到两个 chunk 里所以不能假设每次read拿到的都是完整消息。用buffer缓存按\n\n切分最后一段留在 buffer 里等下一块这样才不会丢数据或者解析出半截内容。这个细节很多教程不讲但实际项目里不处理必然出 bug。6. 常见问题与排查技巧实录6.1 流式输出相关的典型故障做这类服务问题基本集中在几个地方。我把踩过的和帮别人排查过的整理成一张表方便对照。现象可能原因排查方向解决前端等很久才一次性收到全部内容中间层缓冲检查 Nginxproxy_buffering、Tomcat 缓冲关缓冲加X-Accel-Buffering: no连接中途断开报 idle timeout空闲超时看网关/负载均衡超时配置加心跳调大超时用户停止后后端还在跑未处理 abort检查连接断开监听监听 close主动 cancel 模型调用并发上不去请求排队线程池打满看线程池活跃数启用虚拟线程虚拟线程启用后没效果pinning检查 synchronized 块换 ReentrantLock偶发数据丢失或乱码分块解析错误检查前端 buffer 逻辑按\n\n切分并缓存残段6.2 一个隐蔽的竞态问题前面提过 abort 的竞态这里展开说。场景是这样的用户发起请求Controller 创建 sink然后异步启动模型调用。如果用户手速很快在模型调用还没启动的时候就点了停止连接断开事件先触发此时current里的取消句柄还是 null取消落空。等模型调用真正启动后它不知道连接已经断了继续跑。解决办法是在设置取消句柄之后再检查一次 sink 状态Cancellable c modelClient.stream(question, chunk - sink.send(chunk)); current.set(c); if (sink.isClosed()) { c.cancel(); // 补一刀 }这个设置后再检查的模式在并发编程里很常见本质上是把检查和设置这两个操作在时间上做双向覆盖确保无论谁先谁后都能被感知到。6.3 模型 SDK 阻塞导致的虚拟线程失效有些模型 SDK 内部用的是自己的线程池或者用了synchronized保护的连接池。这种情况下即使你的 Web 层用了虚拟线程真正的阻塞发生在 SDK 内部虚拟线程的优势发挥不出来。排查方法是打线程 dump看载体线程是不是被大量占用。如果确认是 SDK 的问题两个方向一是换用支持异步的 SDK 版本二是把 SDK 调用包一层用CompletableFuture配合虚拟线程执行器把阻塞点显式暴露出来。我一般倾向于后者改动小可控。6.4 关于数据一致性的提醒热词里有人问java 怎么保证数据一致性放到 SSE 场景下主要指的是推送顺序和最终状态。SSE 本身是单连接内的有序流同一个连接里消息顺序是有保证的。但如果你在服务端用了多个线程往同一个 sink 里写顺序就乱了。所以要么保证单线程写要么在 sink 内部加锁。另外流式推送的内容和最终落库的内容要一致别出现推给用户的和存下来的不一样这种问题在带工具调用的 AI Agent 场景里特别容易出。7. 一些实战心得虚拟线程这个特性我最大的感受是它把简单还给了 Java 并发。以前为了扛并发被迫学响应式、学各种操作符代码写得像天书。现在用虚拟线程你可以继续用最朴素的阻塞式写法性能还更好。对于 AI 应用这种 IO 密集、逻辑又不算特别复杂的场景这个组合几乎是当前的最优解。封装这件事我的建议是先显式跑通再隐式封装。别一上来就追求优雅的抽象先把最原始的 SSE 链路手写一遍把缓冲、心跳、abort 这些坑都踩一遍你才知道封装的时候该封什么、不该封什么。我见过太多人直接套框架结果出了问题连数据从哪来都说不清。最后分享一个小技巧调试 SSE 的时候用curl -N命令最直观。-N是禁用缓冲你能在终端里实时看到数据一块块冒出来比在浏览器里看 Network 面板清楚多了。命令大概是这样curl -N http://localhost:8080/chat?q你好如果终端里数据是实时冒的说明服务端没问题问题在前端或者中间层如果终端里也是憋半天才一次性出来那问题就在服务端或者它前面的代理。这一招能帮你快速定位问题出在哪一段链路上省下大量瞎猜的时间。
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

Claude Code 主创放弃写 Prompt 了:他改写循环,TaoToken 统一 Key 通道怎么接? 2026/9/29 9:36:40

Claude Code 主创放弃写 Prompt 了:他改写循环,TaoToken 统一 Key 通道怎么接?

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
模型部署与推理优化:INT8量化原理与实践全解析 2026/9/29 9:36:39

模型部署与推理优化:INT8量化原理与实践全解析

模型部署与推理优化02:量化原理与实践(INT8 矩阵乘、校准、QAT 与 LLM 量化)量化这个事,我这两年接触得越多越觉得它被低估了。很多人提到INT8量化,第一反应是“降精度换速度”,好像就是个简单粗暴的取舍。…

阅读更多 →
Verilog Testbench 从入门到实战:仿真验证环境搭建全攻略 2026/9/29 9:36:32

Verilog Testbench 从入门到实战:仿真验证环境搭建全攻略

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
广东电网识别挑战赛赛道三亚军方案:电力缺陷检测全流程解析 2026/9/29 9:36:26

广东电网识别挑战赛赛道三亚军方案:电力缺陷检测全流程解析

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
嵌入式Linux父进程与子进程:从fork到僵尸进程的实战解析 2026/9/29 9:36:25

嵌入式Linux父进程与子进程:从fork到僵尸进程的实战解析

最近在折腾飞凌嵌入式ElfBoard这块板子,跑Linux调进程时碰到一堆和父进程、子进程相关的乱七八糟问题。嵌入式开发里进程这个概念太关键了,尤其当你开始用fork创建子进程,或者发现系统里多了个僵尸进程却不知道是谁搞出来的时候,真…

阅读更多 →
算法表达三维法:流程图、伪代码与N-S图协同建模 2026/9/29 9:36:18

算法表达三维法:流程图、伪代码与N-S图协同建模

1. 算法不是代码,而是“可执行的思维蓝图”很多人学C语言时一上来就写for循环、调printf,结果调试三天搞不定一个冒泡排序——不是语法错了,是脑子里压根没形成算法的“形状”。我带过三十多届嵌入式方向的实习生,发现一个铁律&am…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉