资讯动态

Java后端SSE实战:从SseEmitter到虚拟线程的完整指南

发布时间:2026/9/29 7:25:56 来源:尧图企业网站定制
这几年做Java后端最明显的一个变化就是“AI交互”这件事把很多老技术重新激活了SSE就是其中之一。明明SSE协议早在2011年就有了过去几年一直不温不火结果大模型一出来所有AI应用的前端都想要“打字机效果”SSE直接从冷门协议变成了面试题、架构选型、性能优化三条线都会碰到的常客。我也是从最早手动写HttpServletResponse输出流到后来封装成统一的AI交互模块再折腾虚拟线程解决长连接并发问题中间踩了不少坑。今天把这条完整的路线梳理一遍希望能给正在做JavaAI项目的朋友一些参考。1. 为什么AI大模型场景绕不开SSE1.1 从“全部生成完再返回”到“边生成边推送”传统接口的逻辑很简单请求进来业务计算完成返回一个完整结果。但大模型生成几百个字可能要好几分钟接口如果非得等全部生成完才返回用户盯着空白页面转圈感知上就是“系统卡死了”。流式输出则完全不同它把生成过程拆成一段一段推给前端前端收到一块就渲染一块视觉上就是逐字蹦出来的打字机效果。我在做AI辅助工具的时候对比过同样一段回答走非流式接口首字延迟可能在3到10秒之间用户很容易中途刷新或者放弃走SSE之后300毫秒内就能看到第一个字用户会觉得系统“有反应了”。很多大模型厂商的官方接口默认都建议开启stream参数也是这个原因。对Java后端来说这个变化意味着接口不再是一个方法调用一次返回而是要处理连接生命周期、分段消息、异常中断这些东西。说白了你的代码从“算完返回”变成了“边算边发”数据交付从一次性的变成了持续性的这正好是SSE的主场。1.2 为什么是SSE而不是WebSocket或者轮询很多第一次接触流式推送的人会问直接用WebSocket不是更全能吗轮询也不是不行吧我的经验是选型不能看谁功能多要看场景是否匹配。维度SSEWebSocket轮询通信方向服务端单向推送给客户端全双工客户端主动多次请求底层协议HTTP独立握手协议帧格式HTTP连接成本一条HTTP长连接需要处理握手、心跳、帧分片每次都新建连接自动重连浏览器EventSource自带要自己实现自己实现自定义Header原生EventSource不支持需要fetch方案握手时支持请求里支持典型场景大模型流式回答、进度推送聊天室、协同编辑、在线游戏低频简单轮询AI回答本质上就是服务端向单个用户单向推送文本片段客户端不需要反向发送随机帧最多就是发一个“停止生成”的指令。这种情况用SSE最合适。拿WebSocket去实现相当于开卡车送快递不是不能送但握手、心跳、断线重连全都得自己写复杂度凭空高了一截。轮询就更不匹配了。要做到“打字机”效果要么每几百毫秒拉一次全量文本对接口、数据库、带宽都是灾难要么拉增量但服务端还得维护一个让客户端请求的游标复杂度很快就上来了。而且轮询的实时性天然差一截用户看到的内容永远是“上一次请求的结果”。顺带说一句如果业务里确实需要客户端随时发消息给服务端比如AI助手不仅要流式回复还要支持用户实时控制那也不一定非要上WebSocketSSE加普通的POST请求混用也能解决而且更简单。1.3 SSE协议细节别只盯着text/event-streamSSE虽然走HTTP但有几个点非常容易踩坑。响应头里最重要的三个设置Content-Type必须是text/event-streamCache-Control建议no-cacheConnection保持keep-alive。消息体格式也有讲究每条消息以空行分隔字段包括data、event、id、retry。data后面跟实际内容多行data会被拼成一条带换行的消息event用来区分消息类型id配合客户端的Last-Event-ID实现断线续传retry告诉客户端重连间隔多少毫秒。实际对接大模型的时候很多人把OpenAI流式接口里的[DONE]当成协议的一部分其实它只是业务层的约定。服务端看到[DONE]就知道流结束了此时应该调用complete()或者结束响应。还要注意编码和行尾问题消息行必须以\n结尾如果只发\r\n很多解析器会直接出问题中文内容务必保证UTF-8否则前端渲染出来是一堆乱码。还有一个很隐蔽的坑SSE是用来做服务端到客户端持续推送的但很多中间环节默认会缓冲响应或主动掐断连接。Nginx默认开启proxy_buffering会把你的流式输出攒起来一次性发给客户端打字机效果直接没了网关层还有各种read_timeout长时间没有数据就断连。这些到后面排查章节我再细说。2. 显式调用Java实现SSE的原生玩法2.1 Spring Boot里的SseEmitter从哪开始Java里做SSE有几种路线。最底层的是Servlet的AsyncContext配合HttpServletResponse手动写流生产环境这么干的人很少它把所有细节都暴露给你写起来非常容易出错。Spring 4.2以后提供了SseEmitter这是对ResponseBodyEmitter的扩展专门处理text/event-stream格式。Spring Boot项目里你只需要在Controller方法里返回SseEmitter框架会把HTTP响应切换成长连接模式你往里面写内容、最后complete()就行。用SseEmitter的时候有三个回调要注册onCompletion处理正常结束onTimeout处理超时onError处理异常。构造函数可以传入超时毫秒数不传就用容器默认值。很多初学者以为这个超时是“最多能连接多久”其实它指的是空闲超时也就是说连接一直在发数据就不会断一旦空闲超过设置的时间才会触发onTimeout回调。2.2 一个能跑起来的流式接口示例先给一段典型代码这是显式调用阶段最常见的写法。它虽然能跑但问题也很多后面我再逐个展开。RestController RequestMapping(/api/ai) public class ChatController { private final ExecutorService executor Executors.newFixedThreadPool(8); GetMapping(value /stream, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream(RequestParam String prompt) { SseEmitter emitter new SseEmitter(60_000L); emitter.onCompletion(() - log.info(SSE completed)); emitter.onTimeout(() - emitter.complete()); emitter.onError(e - log.error(SSE error, e)); executor.execute(() - { try { ChatClient client ChatClient.create(); StringBuilder buffer new StringBuilder(); client.streamChat(prompt, delta - { buffer.append(delta); emitter.send(SseEmitter.event().name(message).data(delta)); }); emitter.send(SseEmitter.event().name(done).data(buffer.toString())); emitter.complete(); } catch (Exception e) { emitter.completeWithError(e); } }); return emitter; } }这段代码里能看到几个典型问题线程池用固定8线程并发一高就会排队SseEmitter的60秒是空闲超时大模型思考时间稍微长一点连接就被掐断没有处理客户端主动断开的情况断连后emitter可能一直留在内存里。初版这么写没问题线上跑起来你会发现全是坑。2.3 显式方案的三个坑超时、线程池、资源回收第一个坑是超时。SseEmitter构造参数里的60秒指的是“连接空闲多久会超时”。如果大模型在排队或者网络慢中间有一段时间没有新数据推给前端连接可能提前断掉。这时候前端EventSource会自动重连但重连后服务端又会重新发起一轮大模型调用重复生成还重复计费。解法是心跳在等待期间定时发送一个空注释行或者event为heartbeat的数据帧告诉中间链路“连接还是活的”。第二个坑是线程池。上面代码用固定8线程池同时20个用户并发时任务全在队列里排队前端半天没反应。很多初版项目就是把线程池参数随手一写一压测就崩。这里要理解一点每条SSE连接背后占着一个工作线程线程池大小和在线连接数直接相关这跟普通接口那种“算完就释放”的模型完全不同。后面讲到虚拟线程这个问题的解法会彻底改变但显式阶段你必须先意识到这个矛盾。第三个坑是资源回收。客户端断网、浏览器刷新、用户点了停止生成服务端的emitter不会自动清理。你必须监听前端的abort或者断开通知主动调用emitter.complete()释放连接。漏掉这个的话很快会积累一堆僵尸连接把Tomcat连接数打满。Spring在处理这种断开时通常抛ClientAbortException很多人在onError里只知道打日志打完就完了连接还是没释放这是非常典型的线上事故。3. 隐式封装把AI交互逻辑收敛成一种可复用的结构3.1 封装前先想清楚边界当项目里只有一两个AI流式接口时直接写SseEmitter当然没问题。但很快你会发现每个接口都要重复做相同的事创建emitter、注册回调、拼接消息、处理异常、发心跳、结束时通知前端“回答完了请刷新业务状态”。这些代码属于基础设施逻辑不应该散落在Controller里和业务参数搅在一起。这就是从显式到隐式封装的动机把SSE交互逻辑收敛成一个业务无关的模块。业务方只需要声明“我要一个流式回答内容是XXX”剩下的连接、推送、异常都由框架完成。封装好了之后再来一个新的AI能力只需要写业务回调不用关心底层是哪个模型厂商、消息怎么格式化、连接怎么管理。但也要提醒一句别过度封装。如果封装层把SseEmitter完全隐藏只在最后返回一个完整字符串那等于把流式退回了非流式用户又得等半天。真正的封装是保留事件流语义把繁琐细节拿走把能力留下来。3.2 核心封装组件连接管理、事件格式化、厂商适配我实际落地的封装模块大概有这几个部分。连接管理StreamSessionManager负责维护每个会话的SseEmitter引用、会话状态连接中/已完成/已取消、心跳任务。内部用ConcurrentHashMap存储sessionId到emitter的映射再配一个全局定时任务统一发心跳避免每个emitter各自起一个线程。事件格式化StreamMessageFormatter统一消息结构。比如业务消息用eventmessage结束消息用eventdone错误消息用eventerror并携带错误码。前端只要解析这几种事件类型不用关心底层是哪个模型厂商。厂商适配ModelClientAdapter屏蔽不同大模型厂商流式接口的差异。底座是同一个接口比如StreamModelClient各个厂商分别实现输出统一往StreamSessionManager里写数据。这样以后换模型厂商只要改适配器上层业务和前端都不用动。这里有个细节封装层内不要用裸的String拼接JSON消息很容易出格式错误。可以定义统一的StreamEvent对象通过Jackson序列化序列化失败时走error事件而不是直接把连接弄断。错误事件本身也要标准化前端才能统一展示或者弹出重试按钮。3.3 对接大模型流式API与Abort处理封装过程中最难的是“转发”这件事。大模型SDK大多数都支持流式回调比如onDelta、onComplete、onError你要把这些回调翻译成SSE事件流。伪代码如下public void handleStreamRequest(String requestId, String prompt, StreamSession session) { try { aiClient.streamChat(prompt) .doOnNext(delta - { session.getBuffer().append(delta); session.send(SseEmitter.event().name(message).data(delta)); }) .doOnComplete(() - { session.send(SseEmitter.event().name(done) .data(Map.of(fullText, session.getBuffer().toString()))); session.complete(); }) .subscribe(); } catch (Exception e) { session.sendError(e); } }注意buffer一定要用StringBuilder不要用字符串加号拼接。流式token非常多的时候字符串的不可变特性会导致频繁创建中间对象性能差距很明显。再就是Abort处理。用户点了停止生成或者浏览器端连接断了服务端必须感知。前端用AbortController取消fetch的时候后端的Servlet响应通常会抛出ClientAbortException此时要触发一个onAbort回调把这个requestId下的大模型流式请求也取消掉避免云端还在继续生成、token继续计费。这也是从显式到隐式封装后做得更顺的一件事封装层统一处理“用户取消”这类生命周期事件而不是每个接口自己catch一遍。3.4 流式结束之后落库、文档生成与数据一致性SSE只是传输通道流式结束后往往还有一堆业务动作比如把完整回答落库、更新任务状态、生成报告文档。这些动作不能放在SSE线程里同步做否则会拖慢连接释放还容易阻塞下一条消息。我现在的做法是流式完成后先把“完整答案”写入消息队列或者本地任务表由单独的任务处理器去落库和生成文件。这里会涉及数据一致性问题。比如用户看到的回答已经完整渲染了但数据库里还没写入用户刷新页面可能找不到记录。解法通常靠两点一是用任务表和幂等字段保证后端任务一定执行且只执行一次二是对“AI回答已生成”这件事设计好事务边界先写内容表再更新任务状态最后发送异步通知事件。如果你还要用POI Word生成带图表的报告流程就要更谨慎。POI里的XWPFDocument确实可以创建图表用XWPFChart设置类别和系列数据就能生成线图、柱状图之类的但图表数据要提前绑定到单元格区域不然跑出来的图是空的。业务上建议先保存完整内容和业务快照必要的时候做深拷贝再让后台任务读取快照生成文档避免生成过程中数据被其他并发修改影响。快照的好处是即使后续业务字段变化了文档里记录的仍然是生成那一刻的状态。4. 虚拟线程让长连接不再吃满线程池4.1 传统线程模型下的SSE困境前面提到过显式方案里每条SSE连接都要占一个工作线程去等大模型返回。传统Tomcat默认工作线程数不过200可以通过maxThreads调大但也不是无脑调每个平台线程都有栈空间和切换代价线程太多反而会把CPU消耗在上下文切换上。这意味着当AI服务同时有200个流式连接Tomcat的线程池就被占满了其他普通接口全都在排队。实际项目里我遇到过类似事故一个AI助手服务上线后普通查询接口突然变慢排查下来就是流式连接把线程池拖死的。核心矛盾在于SSE连接处于“等待IO”的状态时线程什么都不干但它仍然占用一个宝贵的执行槽。Java线程是直接映射到操作系统线程的阻塞式模型下线程数总是会被“等待”吃光而CPU利用率却上不去。4.2 虚拟线程到底解决了什么问题JDK 21正式引入了虚拟线程它是JVM在用户态创建的线程不直接映射到操作系统线程。可以把平台线程理解成“工位”虚拟线程理解成“工牌”。传统并发模型是一个员工必须占一个工位直到干完活虚拟线程则是一个员工拿到工牌干活的时候找空工位坐下一旦阻塞在IO等待上JVM立刻把工位让给别的虚拟线程等数据到了再回来接着干。对SSE来说这个特性是天然匹配的每个流式连接需要“一个线程”去等待、回调、发送但大多数时间线程都在阻塞读。虚拟线程可以让大量连接同时存在而底层只需要几十个平台线程在轮转调度。所以JDK 21发布后Java做AI流式接口的并发上限从“几百连接”直接跳到“几万连接”而且代码模型没有根本性改变还是熟悉的Thread、Executor那一套。4.3 Spring Boot与Tomcat中启用虚拟线程启用方式并不复杂。Spring Boot 3.2开始支持虚拟线程你只需要在配置里打开开关spring.threads.virtual.enabledtrue这个配置会把Tomcat的请求处理线程池切换成虚拟线程机制每个请求进入时不再从固定线程池取平台线程而是创建一个虚拟线程来执行。同时Spring的Async、Scheduled等如果配置了虚拟线程执行器也会一并生效。要注意硬性前提JDK必须21及以上Spring Boot版本3.2及以上。如果还在Java 8或者Spring Boot 2.x只能通过外部线程包装或者引入第三方库做类似的事情别指望一个配置开关就能解决。另外如果你使用了WebFlux响应式栈其实并不缺线程虚拟线程主要解决的是阻塞式Servlet栈的问题两者路线不同不需要混在一起比较。4.4 虚拟线程场景实测与注意事项我做过一次不严谨的对比压测同样的流式接口模拟大模型每100毫秒返回一块内容连接保持60秒。平台线程池配置200线程时大概350个并发就开始出现队列等待和超时切换虚拟线程后同样的服务器配置下3000并发依然能平滑响应首字延迟几乎没变化。这组数据不能当正式基准看但趋势非常明显虚拟线程对SSE长连接场景的优化是实打实的。不过虚拟线程也不是无脑好用几个点要特别注意。别在synchronized块里做阻塞IO。JDK 21之前这叫固定pinning问题虚拟线程在synchronized里阻塞时会把底层平台线程也卡住导致性能退化。尽量用ReentrantLock、Semaphore这些基于AQS的并发工具AQS会让线程进入挂起队列等待唤醒这种机制和虚拟线程的调度方式配合得更好。数据库连接池、HTTP连接池的大小不会因为虚拟线程变大。虚拟线程让你能开很多“脑洞”但底层连接池如果只有10个那最终瓶颈还是在连接池上相当于把压力转移到了基础设施资源上。要一起调大连接池或者考虑让连接池也具备更高的并发能力。ThreadLocal要慎用。虚拟线程可以挂载到不同的平台线程上执行如果你在ThreadLocal里存了和某个连接绑定的上下文可能被错误复用。Java 21里更推荐用ScopedValue这样的方案或者显式传上下文对象别依赖线程局部变量。5. 常见问题排查与避坑清单5.1 流式断连idle timeout waiting for SSE这个报错在社区里很常见stream disconnected before completion: idle timeout waiting for SSE。意思是某个环节在等待SSE数据时超过了空闲超时把连接掐断了。排查顺序可以按三层走。先确认是不是服务端自己没数据。如果大模型接口慢或者业务处理卡住不产生任何消息连接自然空转。对策是心跳定时发送一个注释行或者eventping的数据帧让链路知道连接还活着。再查中间链路。Nginx的proxy_read_timeout默认60秒如果流式回答中间有停顿超过这个值连接会被Nginx断开。需要关闭代理缓冲并把超时调大proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; proxy_send_timeout 3600s;最后查客户端。浏览器原生EventSource、各种HTTP库、移动端的流式解析库各自都有空闲超时配置需要确认全链路的值是一致的。很多时候服务端没问题客户端提前把连接关了反而更容易让人摸不着头脑。5.2 前端能不能直接读SSEEventSource与fetch原生EventSource代码非常简单几行就能订阅但它有一个很大限制不能自定义请求头。AI应用里经常要在请求头带Token或者业务标识这时候原生EventSource就没法直接用了。两种绕法一种是在服务端把鉴权信息放到URL参数或者Cookie里但URL参数会出现在日志里安全性一般另一种是前端用fetch读取流式响应手动解析SSE格式。用fetch读SSE时需要按行解析只处理以data:开头的行同时关注event字段和[DONE]结束标记。React项目里可以封装一个useSSE hook内部用AbortController管理生命周期组件卸载、用户停止生成时都调用abort让请求真正取消。我见过一些团队在SSE和WebSocket之间来回纠结其实如果只是“服务端有东西就推、没有就等”SSE加fetch是最省事的一条路轮询文件变化这类场景也够用。5.3 连接不会自动断、乱码、重复重连怎么办几个高频问题一起说。连接不自动断服务端没调用complete()客户端也没主动关闭连接一直挂着。排查时看服务端连接数是不是只增不减然后检查所有分支是否都调用了complete或者completeWithError。尤其要在catch和finally里处理别只盯着正常路径。乱码响应头里charset没设对或者中间代理把Content-Type里的charset改掉了。统一用text/event-stream; charsetutf-8服务端写出的字节保证UTF-8。还有前面说的行尾必须是\n不然解析器会出错。重复重连客户端EventSource会自动重连如果服务端返回了错误码而不是正常关闭客户端会按retry字段疯狂重连。正常结束时一定要返回200并且关闭响应不要给客户端留下错误状态。AI回答这种场景断线后一般没必要恢复中间过程重新发起一次请求即可但业务上要做好幂等避免重复生成、重复扣费。我在实际项目里最大的体会是SSE这件事技术含量不算高但特别考验工程细节。从显式调用到隐式封装提升的不只是代码整洁度更重要的是把连接生命周期、取消机制、异常处理这些最容易出错的地方统一管起来虚拟线程则是把长连接的资源模型问题从根本上解掉了。如果你正在做JavaAI项目建议先把第一版用显式方式跑通亲身体会一下哪里痛再做封装最后再上虚拟线程。不要一上来就追求最先进的方案先让每一步的问题暴露出来你才会真正理解为什么需要这些优化。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑