ARTICLE DETAIL

资讯详情

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

Java实现SSE流式输出:从Servlet到虚拟线程的三种方案

Java实现SSE流式输出:从Servlet到虚拟线程的三种方案 1. 为什么SSE在JavaAI场景里突然成了刚需先说说我自己的经历。去年下半年开始我陆续参与了好几个把大模型能力接入Java后端系统的项目最开始的方案很朴素——前端发一个请求后端调模型API等模型把整段回答生成完再一次性返回给前端。这个方案在demo阶段完全没问题但只要一上真实场景就露馅用户问一个稍微复杂点的问题模型要跑十几秒甚至几十秒前端页面就干巴巴地转圈用户根本不知道后台到底是在干活还是卡死了。后来我们把方案换成了SSEServer-Sent Events服务端推送事件体验立刻不一样了——模型每吐出一个token前端就能实时渲染出来用户看到文字一个个蹦出来心理等待时间大幅缩短。这就是SSE在AI场景里最核心的价值它让流式输出这件事变得极其简单。SSE本质上是一种基于HTTP的单向服务端推送技术。客户端发起一个普通HTTP请求服务端保持连接不关闭持续往客户端写数据数据格式是纯文本的text/event-stream。相比WebSocketSSE的优势在于协议简单、天然支持断线重连、走标准HTTP端口不用额外协商、浏览器原生EventSourceAPI直接支持。对于AI对话这种一问一答、服务端持续推、客户端只负责收的场景SSE几乎是量身定做的。这篇文章我打算把Java里实现SSE的三种典型姿势讲透最原始的显式调用方式、基于Spring的隐式封装方式以及JDK21虚拟线程加持下的性能飞跃。每一层我都会给出可运行的代码、参数选择的理由、以及我在真实项目里踩过的坑。如果你正在做AI应用的后端或者单纯想搞清楚SSE在Java里到底该怎么落地这篇应该能帮你省下不少查资料的时间。2. SSE核心原理与Java实现的三种姿势拆解2.1 SSE协议到底长什么样很多人第一次接触SSE会觉得神秘其实它的协议简单到令人发指。服务端返回的响应头里必须包含Content-Type: text/event-stream Cache-Control: no-cache Connection: keep-alive然后响应体就是一条条以\n\n分隔的消息块每个消息块由若干字段组成data: 你好 data: 我是AI助手 data: 今天想聊点什么 event: done data: [DONE]字段就那么几个data是消息内容可以多行每行一个data:event是自定义事件类型id是消息ID用于断线重连时定位retry告诉浏览器重连间隔。只有data字段是必须的其他都是可选。这里有个新手特别容易踩的坑每条消息必须以两个换行符\n\n结尾否则浏览器不会触发onmessage回调。我见过有人写了一个data: xxx\n就发出去前端死活收不到排查半天才发现少了一个换行。2.2 为什么AI场景偏爱SSE而不是WebSocket这个问题我被问过很多次。WebSocket是双向的SSE是单向的按理说WebSocket能力更强为什么AI场景反而更爱用SSE核心原因有三个。第一AI对话本质是请求-响应模式用户发一次问题服务端流式返回一次答案不需要客户端在同一个连接上持续发消息WebSocket的双向能力用不上。第二SSE天然支持断线重连浏览器EventSource内置了重连机制配合Last-Event-ID头可以实现断点续传而WebSocket要自己实现心跳和重连逻辑。第三SSE走标准HTTP能直接复用现有的鉴权、网关、负载均衡、日志体系运维成本几乎为零WebSocket在很多网关和LB上需要额外配置协议升级。当然SSE也有短板浏览器对同一域名的SSE连接数有限制HTTP/1.1下通常是6个而且只能服务端推客户端。但在AI场景里一个用户通常只开一两个对话窗口这个限制基本无感。2.3 Java实现SSE的三层演进路线我把Java里实现SSE的方式按抽象层次分成三层这也是我实际项目里走过的路径层次实现方式核心API适用场景痛点显式调用直接操作HttpServletResponsePrintWriter学习原理、极简场景手动管理连接、易出错隐式封装Spring MVC/WebFluxSseEmitter/FluxServerSentEvent生产项目主流线程模型需理解虚拟线程JDK21 阻塞式写法Thread.ofVirtual()高并发流式场景需JDK21下面我逐个拆解。3. 显式调用手写SSE响应把原理摸透3.1 最原始的Servlet写法先看最裸的写法用原生Servlet直接往响应里写。这段代码我建议每个做SSE的人都手敲一遍敲完你就彻底理解SSE了WebServlet(/sse/raw) public class RawSseServlet extends HttpServlet { Override protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException { resp.setContentType(text/event-stream); resp.setCharacterEncoding(UTF-8); resp.setHeader(Cache-Control, no-cache); resp.setHeader(Connection, keep-alive); resp.setHeader(X-Accel-Buffering, no); // 关键禁用Nginx缓冲 PrintWriter writer resp.getWriter(); try { for (int i 0; i 10; i) { writer.write(data: 第 i 条消息\n\n); writer.flush(); // 必须flush否则数据攒在缓冲区 Thread.sleep(500); } writer.write(event: done\ndata: [DONE]\n\n); writer.flush(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }这段代码有几个必须注意的点我一个个说。X-Accel-Buffering: no这个响应头是给Nginx看的。Nginx默认会缓冲上游响应导致你服务端明明flush了客户端还是要等一大坨才收到。加上这个头Nginx就会透传。如果你用的是其他反向代理也要查一下对应的缓冲配置。writer.flush()必须每次写完就调。Servlet的PrintWriter内部有缓冲区不flush的话数据会攒着流式效果就没了。我早期有个项目就是忘了flush测试时发现前端要等十几秒才一次性显示全部内容排查了半天。Thread.sleep这里只是模拟模型生成token的耗时。真实场景里这个位置应该是调用大模型API、逐块读取响应流、再逐块转发给前端。3.2 连接生命周期管理显式写法最大的问题是连接管理全靠自己。客户端断开连接时服务端怎么知道答案是往writer写数据时会抛IOException或者resp.isCommitted()配合检查。但更稳妥的做法是注册一个AsyncListener或者用req.getAsyncContext()配合超时。我实际项目里更推荐用异步Servlet因为同步Servlet会一直占着Tomcat的工作线程。假设你有200个并发SSE连接Tomcat默认200个工作线程就全被占满了其他普通请求直接排队。这是个非常隐蔽的性能陷阱。WebServlet(urlPatterns /sse/async, asyncSupported true) public class AsyncSseServlet extends HttpServlet { private final ExecutorService executor Executors.newVirtualThreadPerTaskExecutor(); Override protected void doGet(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException { resp.setContentType(text/event-stream); resp.setCharacterEncoding(UTF-8); resp.setHeader(Cache-Control, no-cache); resp.setHeader(X-Accel-Buffering, no); AsyncContext asyncContext req.startAsync(); asyncContext.setTimeout(0); // 不超时由业务控制 executor.submit(() - { PrintWriter writer null; try { writer asyncContext.getResponse().getWriter(); for (int i 0; i 10; i) { writer.write(data: chunk- i \n\n); writer.flush(); Thread.sleep(500); } writer.write(event: done\ndata: [DONE]\n\n); writer.flush(); } catch (Exception e) { // 客户端断开正常结束 } finally { asyncContext.complete(); } }); } }注意这里我用了Executors.newVirtualThreadPerTaskExecutor()这是JDK21的新特性后面第5节会详细讲。用异步Servlet 虚拟线程200个并发连接只占200个虚拟线程对操作系统线程几乎零压力。3.3 显式写法的适用边界说实话生产项目里我不建议直接用原生Servlet写SSE除非你有非常特殊的定制需求。原因很简单连接管理、异常处理、超时控制、心跳保活这些脏活累活全要自己干代码量大且容易出bug。显式写法的价值在于帮你理解SSE的底层机制理解了之后就该往上走用框架封装好的能力。4. 隐式封装Spring生态下的SSE最佳实践4.1 SseEmitterSpring MVC的经典方案Spring MVC从4.2开始提供了SseEmitter把SSE的连接管理、超时、完成回调都封装好了。这是目前Java后端做SSE最主流的方案。先看一个典型用法RestController RequestMapping(/api/chat) public class ChatController { GetMapping(value /stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream(RequestParam String question) { SseEmitter emitter new SseEmitter(0L); // 0表示不超时 emitter.onCompletion(() - log.info(SSE完成: {}, question)); emitter.onTimeout(() - log.warn(SSE超时: {}, question)); emitter.onError(e - log.error(SSE异常, e)); // 交给业务线程池处理避免阻塞Tomcat线程 chatExecutor.submit(() - { try { // 调用大模型逐块返回 llmClient.streamChat(question, chunk - { try { emitter.send(SseEmitter.event() .data(chunk, MediaType.TEXT_PLAIN)); } catch (IOException e) { emitter.completeWithError(e); } }); emitter.send(SseEmitter.event().name(done).data([DONE])); emitter.complete(); } catch (Exception e) { emitter.completeWithError(e); } }); return emitter; } }SseEmitter的构造函数参数是超时时间毫秒传0L表示永不超时。这里有个坑如果你不传或者传了默认值Spring会用容器默认的超时Tomcat默认30秒30秒后连接自动断开前端会收到一个错误。AI对话动辄几十秒所以一定要显式设置超时。emitter.send()的data方法有多个重载可以指定媒体类型。对于纯文本token用MediaType.TEXT_PLAIN就行。如果你要发JSON对象Spring会自动序列化但要注意序列化后的JSON里不能有换行符否则会破坏SSE的消息格式。我遇到过有人直接发一个带格式化的JSON字符串结果前端解析全乱套。4.2 线程模型SseEmitter最大的坑SseEmitter看起来简单但它背后藏着一个线程模型的陷阱这是我在生产环境被坑得最惨的一次。SseEmitter.send()是阻塞的。如果客户端网络慢或者客户端已经断开但服务端还没感知到send()会阻塞住当前线程。如果你在Tomcat的工作线程里直接调send()一个慢客户端就能拖垮一个工作线程。200个工作线程被200个慢客户端占满整个服务就挂了。所以必须把SSE的发送逻辑放到独立的线程池里就像上面代码里的chatExecutor。这个线程池的大小要根据你的并发SSE连接数来定。我一般会用一个专门的ThreadPoolTaskExecutor核心线程数设成预期并发数的1.2倍队列用有界队列防止内存溢出。但这里又引出一个新问题如果并发SSE连接是1000个你就需要1000个平台线程每个线程默认1MB栈空间光栈内存就1GB。这就是传统线程模型的瓶颈也是JDK21虚拟线程要解决的核心问题。4.3 WebFlux的响应式方案如果你的项目用的是Spring WebFlux那SSE的实现会更优雅因为WebFlux天生就是异步非阻塞的GetMapping(value /stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxServerSentEventString stream(RequestParam String question) { return llmClient.streamChat(question) .map(chunk - ServerSentEvent.Stringbuilder() .data(chunk) .build()) .concatWith(Flux.just(ServerSentEvent.Stringbuilder() .event(done) .data([DONE]) .build())); }WebFlux方案的优势是不需要额外线程池整个链路是非阻塞的一个EventLoop线程能处理成千上万个连接。但代价是学习曲线陡峭而且如果你的下游比如大模型SDK是阻塞式的你还得用subscribeOn切到弹性线程池反而更复杂。我个人的选择标准是新项目、团队熟悉响应式用WebFlux老项目、团队以阻塞式编程为主用SseEmitter 独立线程池。不要为了技术而技术能稳定跑起来才是第一位的。4.4 心跳保活与断线重连SSE连接长时间没数据中间的网络设备防火墙、LB可能会主动断开。所以生产环境必须加心跳。做法很简单起一个定时任务每隔15-30秒往所有活跃的emitter发一个注释行emitter.send(SseEmitter.event().comment(heartbeat));注释行以:开头浏览器会忽略它但能保持连接活跃。我一般设15秒因为很多LB的空闲超时是60秒15秒足够安全。断线重连方面浏览器EventSource会自动重连重连时会带上Last-Event-ID头。如果你想支持断点续传需要在发送消息时带上id服务端根据这个id决定从哪继续。AI对话场景里断线重连后通常直接重新生成更简单所以这个能力用得不多但知道有这回事。5. 虚拟线程JDK21带来的性能飞跃5.1 虚拟线程到底解决了什么问题前面反复提到传统SSE方案的瓶颈是线程数量。每个SSE连接需要一个线程来阻塞等待和发送数据1000个连接就是1000个平台线程。平台线程是操作系统线程创建成本高约1MB栈、上下文切换贵、数量有上限。JDK21正式引入的虚拟线程Virtual Threads彻底改变了这个局面。虚拟线程是JVM管理的轻量级线程栈空间按需增长初始只有几百字节一个JVM里可以轻松创建百万级虚拟线程。当虚拟线程遇到阻塞操作比如Thread.sleep、IO等待时JVM会自动把它挂起让底层的载体线程Carrier Thread去执行其他虚拟线程。用生活类比平台线程就像公司里的正式员工招一个成本高、数量有限虚拟线程就像外包团队需要的时候随时调用完就释放成本极低。对于SSE这种大量连接、每个连接大部分时间在等待的场景虚拟线程简直是天作之合。5.2 用虚拟线程重写SSEJDK21里创建虚拟线程有两种方式。一种是Thread.ofVirtual().start(runnable)另一种是用Executors.newVirtualThreadPerTaskExecutor()。后者更适合SSE场景因为每个连接一个任务Configuration public class VirtualThreadConfig { Bean(sseExecutor) public ExecutorService sseExecutor() { return Executors.newVirtualThreadPerTaskExecutor(); } }然后在Controller里注入这个executorGetMapping(value /stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream(RequestParam String question) { SseEmitter emitter new SseEmitter(0L); sseExecutor.submit(() - { try { llmClient.streamChat(question, chunk - { try { emitter.send(SseEmitter.event().data(chunk, MediaType.TEXT_PLAIN)); } catch (IOException e) { throw new UncheckedIOException(e); } }); emitter.send(SseEmitter.event().name(done).data([DONE])); emitter.complete(); } catch (Exception e) { emitter.completeWithError(e); } }); return emitter; }代码看起来和之前几乎一样但底层线程模型完全不同。现在1000个并发SSE连接只占1000个虚拟线程底层可能只有几个载体线程在真正干活。内存占用从1GB降到几十MB上下文切换开销几乎可以忽略。5.3 虚拟线程的注意事项与踩坑虚拟线程虽好但有几个坑必须知道。第一不要池化虚拟线程。虚拟线程的设计初衷就是用完即弃创建成本极低。如果你用ThreadPoolExecutor去池化虚拟线程反而会失去它的优势。newVirtualThreadPerTaskExecutor()每次任务都新建一个虚拟线程这才是正确用法。第二注意synchronized的钉住问题。在JDK21早期版本里虚拟线程在synchronized块里阻塞时会钉住pin载体线程导致载体线程无法执行其他虚拟线程。JDK21后续版本和JDK24已经大幅改善了这个问题但如果你用的是JDK21早期版本尽量用ReentrantLock替代synchronized。我实测过一个场景把synchronized换成ReentrantLock后吞吐量提升了近3倍。第三ThreadLocal要慎用。虚拟线程数量巨大如果每个虚拟线程都持有ThreadLocal的大对象内存会爆炸。JDK21引入了ScopedValue作为ThreadLocal的替代但目前还是预览特性。在SSE场景里尽量把上下文通过方法参数传递少用ThreadLocal。第四下游阻塞式SDK要确认兼容性。虚拟线程对java.net.Socket、java.nio等JDK内置IO是友好的会自动挂起。但如果你用的HTTP客户端是某些老版本的Apache HttpClient它内部用了synchronized或者native方法可能会钉住载体线程。建议用JDK11的java.net.http.HttpClient它对虚拟线程支持最好。5.4 性能对比实测数据我在一台4核8G的测试机上做了个简单压测模拟SSE连接持续推送数据30秒对比三种方案方案并发连接数平均响应延迟内存占用CPU使用率同步Servlet200稳定约400MB30%SseEmitter平台线程池1000稳定约1.2GB55%SseEmitter虚拟线程5000稳定约350MB40%数据很直观虚拟线程方案在并发能力上提升了5倍内存占用反而更低。当然这是理想环境下的数据真实场景还要考虑下游模型API的限流、网络带宽等因素。但趋势是明确的虚拟线程让Java在高并发流式场景下终于有了和Go、Node.js掰手腕的资本。6. 常见问题排查与实战避坑指南6.1 SSE连接建立后收不到数据这是最高频的问题。排查顺序我总结成一张表现象可能原因排查方法解决方案前端onopen触发但无onmessage消息格式错误抓包看响应体确保每条消息以\n\n结尾数据攒着一次性到达缓冲区未刷新看服务端是否flush每次send后flush加X-Accel-Buffering: no连接几秒后断开超时设置过短看服务端超时配置SseEmitter构造传0LNginx调大超时部分浏览器收不到连接数限制换浏览器测试HTTP/2下限制放宽或减少并发连接我印象最深的一次排查前端一直收不到数据抓包发现服务端明明发了但浏览器就是没反应。最后发现是响应头里少了Content-Type: text/event-stream被Spring的某个拦截器改成了application/json。所以如果你用了自定义拦截器或过滤器一定要确认它们没有覆盖这个头。6.2 客户端断开后服务端资源泄漏客户端关闭页面时SSE连接会断开但服务端如果不感知会继续往一个死连接写数据线程和内存都泄漏了。SseEmitter提供了onCompletion、onTimeout、onError三个回调必须都注册上在里面做资源清理emitter.onCompletion(() - { activeEmitters.remove(emitter); log.info(连接正常关闭当前活跃连接数: {}, activeEmitters.size()); }); emitter.onError(e - { activeEmitters.remove(emitter); log.warn(连接异常关闭: {}, e.getMessage()); });另外emitter.send()在客户端断开后会抛IOException一定要捕获并调用emitter.completeWithError()否则连接状态会一直挂着。6.3 大模型API流式响应的对接细节对接大模型API时不同厂商的流式协议格式不一样。有的返回SSE格式有的返回JSON Lines有的返回自定义分隔符。我一般会写一个适配层把各种格式统一转成ConsumerString回调public void streamChat(String question, ConsumerString onChunk) { HttpResponseInputStream response httpClient.send(request, HttpResponse.BodyHandlers.ofInputStream()); try (BufferedReader reader new BufferedReader( new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; while ((line reader.readLine()) ! null) { if (line.startsWith(data: )) { String data line.substring(6); if ([DONE].equals(data)) break; String content parseContent(data); if (content ! null !content.isEmpty()) { onChunk.accept(content); } } } } }这里有个细节大模型返回的token可能是半个字。比如中文的一个字可能被拆成两个token单独发出去前端会显示乱码。解决办法是在服务端做缓冲凑够一个完整字符再发。判断方法是用Character.isHighSurrogate()检查如果是高代理项就暂存等下一个低代理项拼起来再发。6.4 生产环境部署清单上线前对照这份清单检查一遍能避开80%的坑反向代理Nginx等关闭响应缓冲调大proxy_read_timeout负载均衡开启会话保持或确保SSE连接不跨节点服务端设置合理的最大连接数防止被恶意连接打满心跳间隔小于LB空闲超时建议15秒监控活跃SSE连接数、平均连接时长、异常断开率日志里记录每个连接的建立、完成、异常方便排查客户端实现重连逻辑EventSource自带但要有兜底特别提醒如果你的服务部署在容器里容器的健康检查可能会因为SSE长连接而误判。建议把SSE接口排除在健康检查路径之外或者用独立的端口。7. 我个人的一些实战体会从显式Servlet到SseEmitter再到虚拟线程这条路我走了差不多一年。最大的感受是技术选型要跟着业务阶段走不要一上来就追求最先进的方案。项目初期用SseEmitter 普通线程池完全够用等并发真的上来了再切虚拟线程也不迟。虚拟线程的迁移成本极低基本就是把Executors.newFixedThreadPool换成Executors.newVirtualThreadPerTaskExecutor代码几乎不用改。另一个体会是SSE的调试比想象中麻烦。因为它是一个长连接普通的接口测试工具不一定支持。我常用的组合是curl -N看原始流、浏览器DevTools的EventStream面板看解析后的消息、再加服务端日志。三管齐下基本没有排查不了的问题。最后分享一个小技巧如果你在本地开发时发现SSE数据总是攒着一起到先检查是不是IDE的控制台或者某个中间件在缓冲。我有一次排查了两个小时最后发现是IDEA的HTTP Client在缓冲响应。换成curl立刻就正常了。这种环境问题最耗时间但也最容易被忽略。
返回列表