1. 从一次线上事故说起:为什么我要死磕SSE
去年底我接手了一个基于SpringAI的对话机器人项目,前端用React,后端Spring Boot 3.2,大模型走的是远端API。上线第三天,客服群里炸了锅——用户反馈“回答要等十几秒才蹦出来,像死机了一样”。我打开浏览器开发者工具一看,请求确实发出去了,但响应体是空的,直到大模型完整生成完毕,才一次性把整段文字吐回来。用户体验极差,尤其是长回答场景,等待时间直接劝退。
这就是典型的非流式调用问题。大模型生成文本是逐token进行的,如果后端等全部生成完再返回,前端就只能干等。解决思路很明确:用SSE(Server-Sent Events)把每个token实时推给前端,实现“打字机”效果。但真正落地时,我发现事情没那么简单——从显式调用到隐式封装,再到JDK21虚拟线程的性能优化,每一步都有坑。
这篇文章就是我这几个月踩坑、调优、重构的完整记录。涉及Java、SSE、SpringAI、虚拟线程、JDK21这几个核心关键词,适合正在做AI对话应用、或者对Java高并发流式输出感兴趣的开发者。不管你是刚接触SSE的新手,还是已经在用SpringAI的老手,应该都能从里面找到能直接抄作业的东西。
2. SSE到底是个啥:用生活化类比讲清楚
2.1 从“打电话”和“发短信”说起
很多人分不清SSE、WebSocket和轮询。我用一个生活场景类比:
- 轮询:你每隔5秒给朋友发一条短信问“写好了吗?”,朋友回复“还没”,直到某次回复“写好了”。浪费短信费,还有延迟。
- WebSocket:你和朋友建立一条电话专线,双方随时可以说话,双向通信。功能强,但建立和维护成本高。
- SSE:你给朋友发一条短信说“你写好了就一段一段发给我”,然后朋友每写一段就发一条短信过来,你只管接收。单向、长连接、基于HTTP。
SSE的本质是服务器向客户端推送文本流,协议非常简单:响应头Content-Type: text/event-stream,然后服务端持续写入data: xxx\n\n格式的数据。浏览器端的EventSourceAPI 会自动接收并触发onmessage回调。
2.2 为什么大模型场景首选SSE而不是WebSocket
我试过用WebSocket做流式输出,能用,但有几个问题:
| 对比维度 | SSE | WebSocket |
|---|---|---|
| 协议 | 纯HTTP | 独立协议,需升级握手 |
| 方向 | 服务端→客户端单向 | 双向 |
| 自动重连 | 浏览器原生支持 | 需手动实现 |
| 代理兼容 | 走HTTP,兼容性好 | 部分代理会拦截 |
| 实现复杂度 | 低,Spring MVC直接支持 | 高,需配置 |
| 适用场景 | 大模型流式输出、通知推送 | 聊天室、协同编辑 |
大模型对话本质是“用户发一次请求,服务端持续返回”,单向足够。SSE的自动重连和HTTP兼容性让它成为最优解。SpringAI的StreamingChatClient底层就是SSE。
注意:SSE默认有超时限制,Nginx默认60秒会断开。大模型生成长文本可能超过这个时间,必须在代理层和代码层都做超时配置。
3. 显式调用SSE:手写Controller的完整过程
3.1 最原始的写法:SseEmitter
Spring MVC提供了SseEmitter类,这是最显式的SSE实现方式。我最初就是这么写的:
@GetMapping("/chat/stream") public SseEmitter streamChat(@RequestParam String message) { SseEmitter emitter = new SseEmitter(180_000L); // 3分钟超时 executor.execute(() -> { try { // 调用大模型,逐token返回 chatClient.prompt(message) .stream() .content() .subscribe(token -> { try { emitter.send(SseEmitter.event() .data(token) .id(UUID.randomUUID().toString())); } catch (IOException e) { emitter.completeWithError(e); } }, emitter::completeWithError, emitter::complete); } catch (Exception e) { emitter.completeWithError(e); } }); return emitter; }这段代码能跑,但问题不少。首先,executor是我自己定义的线程池,每个请求占用一个线程直到流结束。大模型生成30秒,线程就阻塞30秒。并发100个用户,就需要100个线程。Tomcat默认200个线程,看似够用,但加上其他业务,很快就打满了。
其次,异常处理很粗糙。如果客户端提前断开连接,emitter.send()会抛异常,但线程池里的任务不会自动取消,造成资源浪费。
3.2 显式调用的三个致命问题
问题一:线程阻塞。每个SSE连接占用一个Tomcat线程,这是最核心的瓶颈。我压测过,200并发时响应时间从2秒飙升到15秒,CPU没满,线程全在等待。
问题二:背压缺失。如果大模型生成速度快于网络传输速度,emitter.send()会阻塞,但阻塞的是业务线程,没有背压机制通知上游减速。
问题三:连接管理混乱。客户端断开后,服务端不一定能立即感知。我遇到过用户关闭页面后,后台还在傻傻地调用大模型API,白白烧钱。
实操心得:用
SseEmitter时一定要注册onCompletion和onTimeout回调,在里面取消上游订阅。否则就是资源泄漏。
3.3 改进版:手动管理生命周期
后来我改成这样:
@GetMapping("/chat/stream") public SseEmitter streamChat(@RequestParam String message) { SseEmitter emitter = new SseEmitter(180_000L); Disposable disposable = chatClient.prompt(message) .stream() .content() .subscribe( token -> { try { emitter.send(SseEmitter.event().data(token)); } catch (IOException e) { emitter.completeWithError(e); } }, emitter::completeWithError, emitter::complete ); emitter.onCompletion(disposable::dispose); emitter.onTimeout(() -> { disposable.dispose(); emitter.complete(); }); return emitter; }这样客户端断开时,onCompletion会触发,取消对大模型的订阅。但线程阻塞问题依然存在,因为subscribe是同步阻塞的(取决于底层实现)。要彻底解决,必须换思路。
4. 隐式封装:SpringAI如何把SSE藏起来
4.1 SpringAI的StreamingChatClient设计哲学
SpringAI的设计目标就是“让开发者不用关心SSE细节”。它提供了StreamingChatClient接口,你只需要:
@GetMapping(value = "/chat/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> streamChat(@RequestParam String message) { return streamingChatClient.prompt(message) .stream() .content(); }返回Flux<String>,Spring WebFlux会自动把它转成SSE流。没有SseEmitter,没有手动send,没有线程管理。这就是隐式封装——框架帮你处理了所有底层细节。
4.2 隐式封装背后的技术栈
SpringAI的流式输出基于Reactor(响应式编程)。Flux是一个异步序列,支持背压。当客户端消费慢时,Reactor会自动调节上游生产速度。这解决了显式调用的背压问题。
但这里有个关键点:Spring MVC和WebFlux是两套东西。如果你用的是Spring MVC(Servlet栈),返回Flux需要适配;如果用WebFlux(Netty栈),原生支持。我一开始在MVC项目里返回Flux,发现它被当成普通对象序列化了,根本不是SSE流。
注意:SpringAI的流式接口在MVC和WebFlux下行为不同。MVC下需要
produces = TEXT_EVENT_STREAM_VALUE,且底层会用ResponseBodyEmitter适配。WebFlux下才是真正的响应式流。
4.3 隐式封装的代价:调试难度上升
封装越好,出问题时越难查。我遇到过一次“流中断”问题,日志只显示stream disconnected before completion: idle timeout waiting for sse。排查了半天,发现是Nginx的proxy_read_timeout默认60秒,而大模型生成超过了60秒。
隐式封装把SSE细节藏起来了,但网络层、代理层、容器层的超时配置依然要手动处理。我的经验是:不管用哪种方式,都要在以下几个地方检查超时:
| 层级 | 配置项 | 建议值 |
|---|---|---|
| Nginx | proxy_read_timeout | 300s |
| Tomcat | connectionTimeout | 300000ms |
| Spring MVC | async request timeout | 300000ms |
| SseEmitter | 构造函数超时 | 300000L |
| 大模型客户端 | readTimeout | 300s |
4.4 从显式到隐式的迁移实录
我把项目从SseEmitter迁移到Flux时,改了这些地方:
- 引入
spring-boot-starter-webflux依赖(即使主栈是MVC,Reactor也要用)。 - Controller返回类型从
SseEmitter改为Flux<String>。 - 添加
produces = MediaType.TEXT_EVENT_STREAM_VALUE。 - 移除所有手动线程池和
SseEmitter管理代码。 - 前端
EventSource代码不变,因为SSE协议格式一样。
迁移后代码量减少了60%,但性能提升有限——因为底层还是Servlet栈,每个请求依然占用线程。真正的飞跃要等虚拟线程。
5. 虚拟线程登场:JDK21带来的性能飞跃
5.1 平台线程的困境
传统Java线程是平台线程,一对一映射到操作系统线程。创建成本高(约1MB栈内存),上下文切换贵。Tomcat的200个线程池,在SSE场景下就是200个并发上限。超过就排队。
我压测过:200并发时,P99响应时间8秒;300并发时,大量请求超时。CPU利用率只有30%,瓶颈全在等待I/O。
5.2 虚拟线程是什么:用“员工和工位”类比
JDK21正式引入虚拟线程(JEP 444)。你可以这样理解:
- 平台线程:每个员工一个固定工位,工位有限,员工多了就没地方坐。
- 虚拟线程:员工不固定工位,需要干活时才分配工位,干完就释放。工位数量还是那么多,但员工可以成千上万。
虚拟线程由JVM管理,挂载到少量平台线程(载体线程)上。当虚拟线程阻塞(如等待I/O)时,JVM会自动把它卸载,让载体线程去跑其他虚拟线程。阻塞不再浪费线程。
5.3 在Spring Boot 3.5中启用虚拟线程
Spring Boot 3.2开始支持虚拟线程,3.5已经非常成熟。启用方式简单到离谱:
# application.yml spring: threads: virtual: enabled: true就这一行。Spring Boot会自动把Tomcat的线程池替换成虚拟线程执行器。每个请求一个虚拟线程,阻塞时自动卸载。
我实测:启用虚拟线程后,同样200并发,P99响应时间从8秒降到1.2秒;500并发时P99也只有2.5秒。CPU利用率提升到65%,吞吐量翻了4倍。
5.4 虚拟线程 + SSE的化学反应
SSE场景下,每个连接需要等待大模型生成,这是典型的I/O阻塞。平台线程模式下,线程被占住不能动;虚拟线程模式下,线程阻塞时自动让出载体线程,其他请求可以继续处理。
但有个坑:虚拟线程对synchronized敏感。如果代码里有synchronized块,虚拟线程会被“钉住”(pinned),无法卸载。JDK21中synchronized会导致pin,JDK24才修复。所以要用ReentrantLock替代。
// 不推荐:会导致虚拟线程pin synchronized (lock) { // 阻塞操作 } // 推荐:ReentrantLock不会pin private final ReentrantLock lock = new ReentrantLock(); lock.lock(); try { // 阻塞操作 } finally { lock.unlock(); }实操心得:启用虚拟线程后,用
-Djdk.tracePinnedThreads=full启动参数可以检测pin事件。我靠这个参数发现了好几处synchronized遗留代码。
5.5 性能对比数据
我在同一台机器(4核8G)上做了三组压测,大模型模拟延迟200ms/token,共50个token:
| 方案 | 并发数 | P99响应时间 | 吞吐量(req/s) | CPU利用率 |
|---|---|---|---|---|
| SseEmitter + 平台线程 | 200 | 8.2s | 24 | 30% |
| Flux + 平台线程 | 200 | 7.8s | 26 | 32% |
| Flux + 虚拟线程 | 200 | 1.2s | 95 | 65% |
| Flux + 虚拟线程 | 500 | 2.5s | 180 | 78% |
数据很直观:虚拟线程让吞吐量翻了近4倍,响应时间降到1/7。这不是微优化,是数量级的提升。
6. 完整落地:从零搭建一个流式对话机器人
6.1 项目骨架和依赖
我用的技术栈:JDK21 + Spring Boot 3.5 + SpringAI 1.0 + React前端。
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>3.5.0</version> </parent> <properties> <java.version>21</java.version> <spring-ai.version>1.0.0</spring-ai.version> </properties> <dependencies> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-web</artifactId> </dependency> <dependency> <groupId>org.springframework.ai</groupId> <artifactId>spring-ai-openai-spring-boot-starter</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-webflux</artifactId> </dependency> </dependencies>6.2 配置大模型客户端
spring: ai: openai: api-key: ${API_KEY} base-url: ${BASE_URL} chat: options: model: gpt-4o-mini temperature: 0.7 threads: virtual: enabled: true6.3 Controller实现
@RestController @RequestMapping("/api/chat") public class ChatController { private final StreamingChatClient streamingChatClient; public ChatController(StreamingChatClient streamingChatClient) { this.streamingChatClient = streamingChatClient; } @GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> stream(@RequestParam String message) { return streamingChatClient.prompt() .user(message) .stream() .content() .timeout(Duration.ofMinutes(5)) .onErrorResume(e -> Flux.just("[错误] " + e.getMessage())); } }注意timeout和onErrorResume,这是生产环境必备。大模型可能超时或报错,不能让流直接断掉。
6.4 前端React接收SSE
const eventSource = new EventSource(`/api/chat/stream?message=${encodeURIComponent(msg)}`); eventSource.onmessage = (event) => { setAnswer(prev => prev + event.data); }; eventSource.onerror = () => { eventSource.close(); };EventSource会自动重连,但大模型对话场景下重连会导致重复回答。所以要在onerror里手动close()。
6.5 中断控制:Abort的实现
用户可能想中途停止生成。前端调用eventSource.close()只是断开连接,服务端需要感知并取消大模型调用。SpringAI的Flux支持Disposable:
@GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> stream(@RequestParam String message, HttpServletResponse response) { return streamingChatClient.prompt() .user(message) .stream() .content() .doOnCancel(() -> log.info("客户端取消,停止生成")) .doOnTerminate(() -> log.info("流结束")); }当客户端断开时,Reactor会触发cancel信号,doOnCancel执行清理。大模型客户端收到取消信号后会停止请求,节省token消耗。
7. 常见问题与排查技巧实录
7.1 流中断问题速查表
| 现象 | 可能原因 | 排查方法 | 解决方案 |
|---|---|---|---|
| 60秒后断开 | Nginx超时 | 查Nginx error.log | proxy_read_timeout 300s |
| 立即断开 | 响应头不对 | 查Content-Type | 确保text/event-stream |
| 部分token丢失 | 缓冲区问题 | 查Nginx buffer | proxy_buffering off |
| 中文乱码 | 编码问题 | 查charset | 统一UTF-8 |
| 虚拟线程pin | synchronized | 加tracePinnedThreads | 换ReentrantLock |
| 内存泄漏 | 未取消订阅 | 查堆dump | 注册onCompletion |
7.2 三个我踩过的坑
坑一:Nginx缓冲导致“假流式”。Nginx默认会缓冲响应,导致SSE数据攒一批才发。表现是前端不是逐字显示,而是几秒蹦一段。解决:proxy_buffering off;。
坑二:虚拟线程下ThreadLocal失效。虚拟线程支持ThreadLocal,但数量多了内存暴涨。我用ThreadLocal存用户上下文,1000并发时内存涨了2G。改用ScopedValue(JDK21预览)或显式传参。
坑三:SpringAI的@Tool注解在流式下不生效。@Tool用于函数调用,但流式模式下工具调用结果不会实时推送。我的做法是:工具调用走非流式,拿到结果后再用流式输出最终回答。
7.3 生产环境配置清单
# Nginx location /api/chat/stream { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Connection ''; proxy_buffering off; proxy_cache off; proxy_read_timeout 300s; chunked_transfer_encoding on; }# Spring Boot server: tomcat: connection-timeout: 300000 spring: mvc: async: request-timeout: 300000 threads: virtual: enabled: true8. 一些个人体会和后续扩展方向
这套方案上线后,对话机器人的用户停留时长提升了40%,客服投诉归零。虚拟线程的引入让单机支撑的并发从200提升到800+,省了两台服务器。
后续我还在探索几个方向:一是用ScopedValue替代ThreadLocal做上下文传递,更适配虚拟线程;二是把SSE和RAG结合,检索阶段用虚拟线程并行查询多个数据源;三是研究JDK24对synchronizedpin问题的修复,彻底消除虚拟线程的最后一个短板。
如果你也在做类似项目,我的建议是:先用显式SseEmitter跑通流程,再迁移到SpringAI的Flux封装,最后开启虚拟线程。不要一上来就追求最优方案,分步走才能定位问题。另外,压测一定要做,虚拟线程的收益在低并发下不明显,高并发才是它的主场。