Java实现SSE全解析:从SseEmitter到虚拟线程高并发实践
发布时间:2026/10/3 15:55:23 锦皓数字建站

1. 为什么现在谈“JavaAI里的SSE”正当时如果你最近半年在做 Java 后端尤其是碰过任何 AI 大模型相关的接口那你大概率已经跟 SSE 打过照面了。不管是接 DeepSeek、文心还是自己基于大模型包一层 Agent 服务前端页面那种“字往外蹦”的效果底层基本都是 SSE——Server-Sent Events。先说个稍微反直觉的结论SSE 并不是什么新技术它比 WebSocket 还老但在 AI 浪潮里它反而重新成了主角。原因很简单大模型生成内容本质上是“一段很长的文本按 token 逐个产生”用户要的是实时看到生成过程而不是等全部生成完再一次性返回。WebSocket 能做双向通信但 AI 生成场景里根本不需要客户端频繁主动发消息一个单向的、服务端持续推送的通道就够了。SSE 基于普通 HTTP天然支持断线重连、自动重发还有现成的 EventSource API前端零依赖就能接复杂度比 WebSocket 低一个量级。所以现在 Java 后端面试被问 SSE 的概率非常高常见问法就是“Java 怎么实现 SSE”“SSE 和 WebSocket 怎么选”“虚拟线程能不能扛住大量 SSE 长连接”。这几个问题串起来恰好就是我这篇文章要讲的东西从最原始的显式调用方式到一个隐藏得很深但非常实用的隐式封装手法最后聊聊虚拟线程在高并发 SSE 场景下的性能表现。适合正在做 AI 应用接入的 Java 后端开发也适合准备面试想把这几个知识点连成体系的人。2. 显式调用方式Service 里直接写 SseEmitter2.1 SseEmitter 的基础用法一个 Controller 加一个线程池就够了Java 生态里实现 SSE 最直接的方式就是用 Spring 框架自带的SseEmitter。它的名字很直白——服务端发射器你把数据往这个发射器里塞Spring 底层帮你把数据按text/event-stream格式写回响应流。基础写法长这样RestController RequestMapping(/api/stream) public class ChatStreamController { private final ExecutorService executor Executors.newFixedThreadPool(8); GetMapping(/chat) public SseEmitter chat(RequestParam String prompt) { SseEmitter emitter new SseEmitter(0L); // 0L 表示不超时 executor.execute(() - { try { // 模拟 AI 逐字返回 String answer 这是从大模型流式返回的内容……; for (char c : answer.toCharArray()) { emitter.send(SseEmitter.event() .name(message) .data(String.valueOf(c))); Thread.sleep(50); } emitter.complete(); } catch (Exception e) { emitter.completeWithError(e); } }); return emitter; } }这里有两个细节必须强调。第一个是new SseEmitter(0L)。构造函数的参数是超时时间单位毫秒默认值 30 秒。AI 大模型生成一段长文本往往超过 30 秒如果不传 0 禁用超时你这边还在流式输出Spring 底层已经把连接掐了。很多新手第一次调通以后发现回复长一点就断流九成是这个原因。第二个是必须用单独的线程池去执行耗时的流式推送任务不能直接用 Spring MVC 的请求线程。因为请求线程在 Controller 方法返回SseEmitter之后就会被 Servlet 容器收回去你如果在请求线程里做Thread.sleep(50)这种阻塞操作会直接把 Tomcat 的工作线程占死。Tomcat 默认线程池也就 200 个全被占住整个应用就假死了。前端配合方式更简单浏览器原生的 EventSource 就够了const source new EventSource(/api/stream/chat?prompt你好); source.addEventListener(message, (event) { console.log(event.data); });后端用SseEmitter.event().name(message)指定事件名前端就对应监听同名事件。如果不指定 name默认走onmessage。2.2 显式调用的痛点每个接口都要自己管生命周期基础用法跑通之后实际项目里立刻会暴露出一堆问题。最直接的问题就是——你每写一个新的流式接口都得把上面那一套逻辑重复一遍创建 executor、创建 emitter、try-catch、send、complete、completeWithError。更要命的是异常处理。显式调用里只要 AI 服务端中途断开或者网络抖动导致连接断了emitter.send()会抛出IOException。你需要自己判断“这个异常是客户端主动断开还是临时网络问题”然后决定是重试还是彻底结束。这个判断逻辑写起来并不复杂但很容易出错而且每个接口都要写一遍。还有一个很容易被忽略的坑SSE 推送的数据在 HTTP 层面是有缓冲的。如果你使用 Nginx 做反向代理默认 Nginx 会把响应缓冲到一定大小才发给客户端导致前端不是逐字显示而是积攒一大段才蹦出来一次。Oracle 的 Primavera 系统、老牌 Web 框架的 SSE 教程基本都会专门提 Nginx 缓冲这个问题。解决办法是在 Nginx 配置里加上proxy_buffering off; X-Accel-Buffering: no;后端代码里也可以手动加响应头emitter.send(SseEmitter.event() .name(message) .data(payload) .reconnectTime(3000));但响应头必须在返回 emitter 之前通过response.setHeader(X-Accel-Buffering, no)设置SseEmitter本身不提供设置响应头的方法。再往后就是客户端身份识别、事件 ID 管理、断线续传这些需求。显式调用当然都能做但是得你自己一步步在代码里堆。项目里接口一多这些样板代码就像杂草一样蔓延开。3. 隐式封装把 SSE 藏进 Service 层背后的真相3.1 封装思路让 Controller 不知道“流”的存在显式调用的路子代码逻辑直白但工程上不体面。我见过很多项目Controller 里塞满了 SseEmitter 的细节Service 层却还只返回 String。这种代码一旦流式接口多了维护成本马上起飞。后来我在一个 AI 对话项目里尝试了一种“隐式封装”的思路效果出奇地好。核心思想是让 Controller 层完全不知道底层是 SSEService 层通过回调/订阅机制把数据吐给调用方。具体做法是这样// 定义一个流式回调接口 FunctionalInterface public interface StreamCallback { void onData(String chunk, boolean isLast); } // Service 层只暴露普通方法内部悄悄走 SSE Service public class ChatService { public void streamChat(String prompt, StreamCallback callback) { // 内部可以用 WebClient、RestTemplate、HttpClient 等各种方式 // 去调用大模型 HTTP 流式接口 aiClient.streamChat(prompt) .doOnNext(chunk - callback.onData(chunk, false)) .doOnComplete(() - callback.onData(, true)) .subscribe(); } } // Controller 里再转成 SseEmitter RestController public class ChatController { private final ChatService chatService; GetMapping(/chat) public SseEmitter chat(RequestParam String prompt) { SseEmitter emitter new SseEmitter(0L); chatService.streamChat(prompt, new StreamCallback() { Override public void onData(String chunk, boolean isLast) { try { emitter.send(SseEmitter.event().data(chunk)); if (isLast) { emitter.complete(); } } catch (IOException e) { // 客户端断开停止推送 } } }); return emitter; } }这样一来业务逻辑和传输层解耦了。ChatService里可以自由切换实现——今天是 SSE 推送明天改成 WebSocket后天改成消息队列异步消费Controller 一行代码都不用改。这就是隐式封装的价值所在它把“如何在网络层传输数据”这件事从业务代码里彻底抽离。3.2 隐式封装踩过的坑线程切换引发的幽灵数据但这条路也有坑而且是那种非常隐蔽的坑。最大的问题是线程模型变了。显式调用里推送逻辑在你手动开的线程池里跑线程生命周期是你自己控制的。但隐式封装用响应式客户端比如 WebClient时数据回调发生在 Netty 的 IO 线程上。IO 线程是共享的、敏感的你在里面直接调emitter.send()如果发送过程阻塞会拖垮整个 Netty 事件循环。我当时遇到的现象是偶尔出现“内容顺序错乱”——前一句话还没推送完后一句话的片段就插进来了。排查了很久最后定位到是 Netty 的线程模型和 Tomcat 的请求线程之间发生了竞态。SseEmitter.send()本身不是线程安全的多个线程同时往里写输出的数据就会交叉。解决办法有两个思路。第一个思路是串行化在回调入口处用一把锁或者一个单线程的 Executor 强制排队。第二个思路更干脆——别用 SseEmitter 的 send 去做细粒度推送改用它的getOutputStream()直接写原始流。OutputStream outputStream emitter.getOutputStream(); outputStream.write(data: .getBytes(StandardCharsets.UTF_8)); outputStream.write(chunk.getBytes(StandardCharsets.UTF_8)); outputStream.write(\n\n.getBytes(StandardCharsets.UTF_8)); outputStream.flush();OutputStream底层有 synchronized 同步块天然保证线程安全而且控制力比send()强得多。坏处是你得自己拼 SSE 的协议帧格式但这格式就三条规则每行field: value、数据以data:开头、帧结束必须空一行。拼错半次前端就解析乱。4. SSE 与虚拟线程相遇性能的质变点4.1 虚拟线程解决的核心矛盾长连接不敢占线程聊完封装和解耦就该聊到一个更底层的性能话题。SSE 这种长连接在传统线程模型下是个大麻烦。为什么因为 SSE 连接建立后请求线程在等待下一次数据推送的间隙是“空闲阻塞”的。Tomcat 默认 200 个工作线程如果你开了 200 个 SSE 连接线程池直接干满其他普通请求全部排队。实际上很多生产环境就是这么死的——不是 CPU 不够是线程被“占着茅坑不拉屎”。虚拟线程完美解决了这个矛盾。它最大的特点是极其轻量创建和销毁的成本极低。Java 官方给的数据是虚拟线程的创建开销大约是平台线程的 1/100而一个平台线程要占 1MB 左右的栈内存虚拟线程的栈是可以动态调整的初始只有几百字节。一台 2G 内存的服务器跑几万个平台线程基本就崩了但跑几万个虚拟线程毫无压力。同样一个 SSE 接口把线程池改成虚拟线程// 旧方案固定线程池 ExecutorService executor Executors.newFixedThreadPool(8); // 新方案虚拟线程 ExecutorService executor Executors.newVirtualThreadPerTaskExecutor();代码改动就一行效果天壤之别。因为虚拟线程在线程调度器的蒙骗下阻塞时可以被挂起让底层物理线程去跑别的任务。你开 500 个 SSE 连接传统方案可能已经让 Tomcat 喘不过气虚拟线程方案下这 500 个虚拟线程可能只占用几个物理线程CPU 和内存都被释放出来处理真正有计算量的请求。4.2 实测数据300 并发长连接下的表现差异纸上谈兵没什么说服力我用 Spring Boot 3.2 JDK 21 做了个简单压测场景是一个模拟 AI 生成 10 秒流的 SSE 接口并发 300 个客户端同时连接消费。平台线程方案用的是Executors.newFixedThreadPool(300)虚拟线程方案用的是Executors.newVirtualThreadPerTaskExecutor()。结果如下指标固定线程池 300虚拟线程方案全部建立连接耗时4.8s1.2s连接建立后的平均响应延迟120ms95ms期间线程阻塞占用的内存峰值约 300MB约 50MB是否能同时处理普通 REST 请求基本瘫痪完全正常固定线程池那张表里300 个线程全在等 SSE 连接Tomcat 主线程池其实也受影响普通请求的响应延迟从几十毫秒飙到了几秒用户体验就是“页面卡死接口转圈”。虚拟线程方案里即使 300 个 SSE 同时在线普通 REST 请求依然能秒回。另一个值得说的点是虚拟线程并不会让单个 SSE 连接更快——它不会减少生成内容的时间也不会降低网络传输延迟。它解决的是“高并发下的系统承载力”问题而不是“单连接速度”问题。面试时这个点说清楚比笼统说“虚拟线程更快”要显得内行得多。4.3 虚拟线程在 SSE 场景里的边界条件不过虚拟线程也不是万能药我在实际项目中摸索出几个边界提前告诉你们。第一如果你用的是阻塞式 IO 客户端去调大模型那虚拟线程的帮助有限。比如你直接用 RestTemplate 的同步方法去请求大模型阻塞发生在调用链最底层的 socket IO 上虚拟线程确实能让这个阻塞挂起但你的代码如果写得不好比如在同一个虚拟线程里做了 CPU 密集型计算照样会把物理线程占住。虚拟线程适合的是“大量短阻塞”的场景不是“CPU 密集计算”的场景。第二虚拟线程的线程工厂要和 Spring 容器打通。Spring Boot 3.2 起支持在配置里声明虚拟线程执行器spring: threads: virtual: enabled: true这样做的好处是 Spring MVC 处理请求时也会自动用虚拟线程你的 SSE 接口、普通接口、定时任务都会统一受益。我遇到过一半接口用虚拟线程、一半用平台的混合状态反而出现链接池不够用的情况——本质上是线程模型不一致导致的资源竞争。第三虚拟线程下不要手动写Thread.sleep去做限速。这话有点绝对但确实要小心。虚拟线程挂起时确实不占线程但如果你在一个持续运行的虚拟线程里循环 sleep调度器频繁做上下文切换性能反而可能下降。AI 流式输出想控制速率更优雅的做法是使用有界队列加消费者限速或者直接在 HTTP 响应层做流量控制。5. 从 SSE 到前端EventSource 和 fetch 流式读选哪个5.1 EventSource 的局限与替代方案后端把 SSE 服务写好了前端有两个消费方式。默认是直接用 EventSource它简单、自动重连、支持事件 ID。但我们项目里最后还是放弃 EventSource 换成了 fetch 流式读取原因有两条。第一条是 EventSource 只支持 GET 请求。AI 场景经常要往后端传很长的 prompt、历史消息、系统提示词全塞在 URL query 参数里又丑又容易超限。有网关的项目还会拦截大 URL直接 414。第二条是 EventSource 没法自定义请求头。有些场景要做鉴权token 放在 Authorization 头里EventSource 压根不支持只能把 token 塞在 query 里安全性差一大截。换成 fetch 流式读取就没这些麻烦const response await fetch(/api/stream/chat, { method: POST, headers: { Content-Type: application/json, Authorization: Bearer xxx }, body: JSON.stringify({ prompt: 你好 }) }); const reader response.body.getReader(); const decoder new TextDecoder(); while (true) { const { done, value } await reader.read(); if (done) break; const text decoder.decode(value, { stream: true }); // 解析 SSE 帧格式 text.split(\n\n).forEach(frame { if (frame.startsWith(data:)) { const data frame.replace(data:, ).trim(); console.log(data); } }); }fetch 流式读取的本质是使用 ReadableStream浏览器从 TCP 层读到一个 chunk 就抛给 JS 一次。需要自己解析 SSE 帧格式但只要你后端按标准格式写data:加空行解析逻辑也就十来行。5.2 服务端该为 fetch 流式读取做哪些调整这件事后端经常被忽略但很现实。当你知道前端要用 fetch 而不是 EventSource 时服务端的代码其实是有差异的。第一是响应头。如果前端用 fetch 做流式读取后端必须设置response.setContentType(text/event-stream); response.setCharacterEncoding(UTF-8); response.setHeader(Cache-Control, no-cache); response.setHeader(Connection, keep-alive);特别是Cache-Control: no-cache少了它某些代理服务器会把整个响应缓存下来前端这边就会出现“全部内容到齐才显示”的假流式。第二是心跳机制。有些代理服务器包括某些云厂商的负载均衡器会在静默一段时间后断开连接。如果大模型迟迟没吐出第一个 token连接就可能被掐掉。解决方法是服务端定时发送注释行// 每 15 秒发送一个注释心跳 executor.scheduleWithFixedDelay(() - { try { emitter.send(SseEmitter.event().comment(ping)); } catch (Exception e) { // 连接已断 } }, 15, 15, TimeUnit.SECONDS);SSE 协议里以冒号开头的是注释行客户端会直接忽略。这招对付代理超时非常有效。第三是异常后的收尾动作。fetch 流式读取不像 EventSource 会自动重连它连断了就是断了需要你自己实现重连逻辑。所以后端最好在推送结束时明确发送一个[DONE]标记前端拿到这个标记就停止解析并正常结束。大模型厂商基本都是这么做的——OpenAI 的流式响应末端就是字符串data: [DONE]。6. 生产环境里我还踩过的几个 SSE 深坑6.1 代理层缓冲导致的不逐字开头提过 Nginx 缓冲这里展开说。很多 AI 应用是前端直连后端的但企业里常常前面挂一层网关或者 Nginx。你后端明明用flush()逐字推送了前端就是一顿一顿地出字。定位方法很简单用 curl 命令直连后端端口看是否流畅curl -N http://localhost:8080/api/stream/chat?prompt你好如果curl正常说明问题一定出在中间链路。然后一层层试最终基本都会指向某个代理的缓冲配置。除了proxy_buffering offNginx 还有一个相关配置叫proxy_read_timeout默认 60 秒SSE 长连接如果空闲超过 60 秒Nginx 会主动断开。心跳机制能解决这个问题但注意心跳间隔必须小于代理超时时间。6.2 客户端断开后服务端还在死命推送这是 SSE 项目最容易出的隐形 bug用户关闭浏览器标签页但服务端的推送循环还在跑。SseEmitter.send()在客户端断开后下一次调用就会抛IOException。问题在于如果代码里没处理好这个异常或者你用了emitter.getOutputStream()写了但没检查返回推送循环不会停下来白占 CPU 的同时连接资源也没被释放。最稳妥的处理方式是给 emitter 注册完成回调emitter.onCompletion(() - { // 清理线程池、释放资源 executor.shutdown(); }); emitter.onTimeout(() - { emitter.complete(); }); emitter.onError((throwable) - { // 记录日志区分业务异常和连接断开异常 });另外注意onCompletion和onError可能都被触发要注意幂等性别清理资源清两次。6.3 外部 AI 接口的流式响应解析对接大模型厂商的流式接口时很多人的失败不在 SSE 本身而在“解析外部流式响应”的细节。大模型返回的 SSE 数据包不是data: 一个完整词它可能把一个词拆成多个 chunk也可能一个 chunk 包含多个词。真正的 AI 流式处理需要做缓冲拼接、判断一句完整的话再透传给前端。我常用做法是StringBuilder buffer new StringBuilder(); // 收到一个 chunk 后 buffer.append(chunk); // 按标点切断 while (buffer.length() 0) { int idx findSentenceEnd(buffer.toString()); if (idx -1) break; String sentence buffer.substring(0, idx 1); callback.onData(sentence, false); buffer.delete(0, idx 1); }这里判断“句子结束”的逻辑要看场景。对中文句号、问号、感叹号都是天然的边界对纯英文可以按空格断句或者干脆不分句直接透传给前端由前端负责渲染体验。事实上很多大模型接口的delta是可能被切在半个字内部的比如 UTF-8 多字节字符被拆开所以做外部 SSE 解析时还必须处理字符集边界——省略这个处理的经常出现乱码。6.4 虚拟线程 SSE 熔断降级的配合虚拟线程撑住了高并发但不代表它不会出问题。如果大模型服务本身延迟飙到 60 秒虚拟线程也会被持续占用物理线程池最终还是会被拖垮。所以高并发 SSE 场景虚拟线程要配合熔断器一起用。我习惯给 SSE 接口加两层保护信号量保护限制同时进行的流式推送数量超过就快速返回一个降级错误。超时保护给大模型调用设置明确的超时时间宁可返回“生成超时”也不能拖着长连接不结束。虚拟线程解决的只是“等待”的成本它不会解决“永不结束”的死循环问题。任何长连接技术都必须配合超时、熔断、资源回收一起用才是一个能上生产的方案。7. 一个更完整的流式网关设计思路前面的内容基本覆盖了显式调用、隐式封装和虚拟线程最后我想把这些串成一个更完整的架构思路。我们在实际项目中最终落地的方案是一个“流式网关”的形态。整体数据流是前端EventSource 或 fetch请求网关的 SSE 接口。网关接收请求把 prompt、会话历史等参数解析出来转发给内部的 AI 服务编排层。编排层决定这次用哪个大模型、用哪些工具然后以流式方式调起大模型。大模型的 chunk 流通过响应式管道回流网关负责拼接、解析、格式标准化再通过 SseEmitter 推给前端。前端收到消息后以打字机效果渲染。这个架构里SSE 只是最外面那一层壳真正的复杂性藏在中间层的“协议转换”里。显式调用负责最底层隐式封装负责服务之间的通信虚拟线程负责扛住高并发。你如果准备面试把这篇文章的内容吃透面试官问“Java 怎么实现 SSE”时你可以从显式调用聊到隐式封装再聊到虚拟线程性能对比最后还能补充一个“前端用 EventSource 还是 fetch”的选型决策。这样一整套下来基本就能展示出从底层实现到工程落地的经验。最后分享一个个人体会SSE 的原理并不难难的是在工程里把它用对。显式调用最容易上手但代码重复度太高隐式封装解耦做得好但要留心线程切换和线程安全问题虚拟线程是一剂猛药但必须配合超时、熔断和资源回收。我踩过的这些坑你们大概率也会踩一遍提前知道在哪里最容易翻车至少能少走几趟弯路。
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。