☰
Java SSE 实战:从 SseEmitter 到虚拟线程与 Spring AI 流式对话
2026/10/1 13:37:51 网站建设 项目流程

SSE 这个东西,最早我在做后台任务进度推送的时候用过,那时候还是 Servlet 3.0 的AsyncContext手搓,代码写得跟意大利面一样。后来 Spring 出了SseEmitter,封装了一层,用起来舒服多了。再后来 JDK 21 把虚拟线程正式落地,Spring Boot 3.2 直接一行配置就能开启,SSE 这种"一个请求挂很久"的场景,突然就从"不敢多用"变成了"随便开"。这篇文章就把我这几年来在 Java 里折腾 SSE 的几条路径捋一遍——从最原始的显式写法,到 Spring AI 里的隐式封装,再到虚拟线程带来的吞吐量变化,顺带把踩过的坑和实测数据都摊开讲。

1. 为什么 SSE 在 Java 里一直是个"尴尬"的存在

1.1 SSE 到底解决了什么问题

先说清楚 SSE(Server-Sent Events)是什么。它本质上就是一个长连接的 HTTP 响应,服务端把Content-Type设成text/event-stream,然后保持连接不关闭,持续往客户端写数据。客户端用浏览器原生的EventSource对象接收,收到一条解析一条。

它和 WebSocket 最大的区别在于:SSE 是单向的,只能服务端推、客户端收。但恰恰因为单向,它比 WebSocket 简单太多——不需要协议升级、不需要握手协商、不需要心跳保活(HTTP 层自带)、走的就是普通 HTTP 端口,Nginx 配一下proxy_buffering off就能用。

我见过太多团队一提到"实时推送"就上 WebSocket,结果发现业务根本不需要双向通信,白白背上了连接管理、心跳、重连、鉴权这一堆复杂度。SSE 在这些场景下是更务实的选择:

  • 大模型流式对话输出(这个现在最火)
  • 后台任务进度条
  • 日志实时 tail
  • 股票/行情推送
  • 消息通知

1.2 传统 Servlet 模型下的致命伤

问题来了。SSE 的核心特征是"连接长时间挂着",而传统 Java Web 的线程模型是一个请求占一个线程(thread-per-request)。Tomcat 默认 200 个线程,意味着最多同时挂 200 个 SSE 连接,第 201 个请求就得排队。

你可以把 Tomcat 线程池调大,比如调到 2000。但每个线程在 JVM 里默认占 1MB 栈空间(-Xss默认值),2000 个线程就是 2GB 内存,还没算上下文切换的开销。而且这些线程 99% 的时间都在wait,纯纯的资源浪费。

这就是为什么早些年大家对 SSE 又爱又恨——功能好用,但并发上不去。一个用户开一个 SSE 连接,一万个在线用户就是一万个线程,服务器直接跪。

这里有个常见的误解:很多人以为用了SseEmitter就"异步"了,就不占线程了。其实不是。SseEmitter只是把响应的写入权交出来,底层那个处理请求的 Tomcat 线程该占还是占,只是它不再阻塞在业务逻辑上,而是阻塞在等待数据上。真正的解放要等到虚拟线程或者响应式编程。

1.3 三条技术路线的分野

所以 Java 里做 SSE,实际上有三条路:

路线代表技术线程模型适用场景
显式调用SseEmitter/AsyncContext平台线程阻塞连接数少、逻辑简单
响应式WebFlux +Flux<ServerSentEvent>事件循环,无线程阻塞高并发、团队熟悉 Reactor
虚拟线程JDK 21 + Spring Boot 3.2虚拟线程阻塞高并发、想保留同步写法

WebFlux 那条路我走过,学习曲线陡,调试痛苦,一个Flux的背压没处理好就各种诡异问题。对于大多数业务团队来说,虚拟线程是性价比最高的方案——代码还是同步写法,但并发能力直接起飞。下面我按这三条路依次展开。

2. 显式调用:SseEmitter 的完整落地与那些文档没写的细节

2.1 最小可用版本长什么样

先上一个能跑的最小例子,Spring MVC 环境:

@RestController public class StreamController { @GetMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter stream() { SseEmitter emitter = new SseEmitter(0L); // 0 表示不超时 ExecutorService executor = Executors.newSingleThreadExecutor(); executor.execute(() -> { try { for (int i = 0; i < 10; i++) { emitter.send(SseEmitter.event() .id(String.valueOf(i)) .name("message") .data("chunk-" + i)); Thread.sleep(500); } emitter.complete(); } catch (Exception e) { emitter.completeWithError(e); } finally { executor.shutdown(); } }); return emitter; } }

这段代码有几个点必须说清楚:

第一,SseEmitter(0L)的超时设置。默认超时是 30 秒(取决于容器),到点会自动 complete。做长连接推送必须设成 0 或者一个很大的值,否则用户看个进度条看到一半连接就断了。但设成 0 也有风险——如果客户端异常断开而服务端没感知到,这个 emitter 会一直挂着,直到 TCP keepalive 超时。所以生产环境我一般设成 30 分钟,配合前端定时重连。

第二,produces必须显式声明。不写MediaType.TEXT_EVENT_STREAM_VALUE,浏览器拿到的是普通响应,EventSource直接报错。这个坑我见过至少三个同事踩过。

第三,emitter.send()的线程安全。SseEmitter内部对 send 做了同步,但如果你在多个线程里并发 send,顺序是不保证的。做流式输出时,务必保证单线程顺序写。

2.2 连接生命周期管理:onCompletion / onTimeout / onError

SseEmitter提供了三个回调,这是管理连接的关键:

emitter.onCompletion(() -> { log.info("SSE completed, emitterId={}", emitterId); // 清理资源、从在线列表移除 }); emitter.onTimeout(() -> { log.warn("SSE timeout, emitterId={}", emitterId); emitter.complete(); }); emitter.onError((ex) -> { log.error("SSE error, emitterId={}", emitterId, ex); // 客户端断开通常走这里 });

这里有个非常隐蔽的坑:客户端主动关闭页面时,服务端不一定立刻触发onError。因为 TCP 连接可能还处于半开状态,服务端下一次send()才会抛IOException。所以如果你依赖onError来清理资源,可能会延迟很久。我的做法是每次 send 都包一层 try-catch,一旦抛异常立即 complete 并清理:

private void safeSend(SseEmitter emitter, Object data) { try { emitter.send(data); } catch (IOException e) { emitter.completeWithError(e); throw new RuntimeException("client disconnected", e); } }

2.3 心跳:不是可选项,是必选项

SSE 连接长时间没数据,中间的任何一层(Nginx、负载均衡、防火墙)都可能把它掐掉。默认的 idle timeout 通常是 60 秒。所以必须发心跳。

心跳有两种做法:

  • 注释心跳:emitter.send(SseEmitter.event().comment("ping")),客户端EventSource会忽略注释行,但连接保持活跃。
  • 业务心跳:发一个event: heartbeat的空事件,前端可以选择性处理。

我一般用注释心跳,每 15 秒一次,用一个ScheduledExecutorService统一调度:

ScheduledExecutorService heartbeatScheduler = Executors.newScheduledThreadPool(1); heartbeatScheduler.scheduleAtFixedRate(() -> { emitters.forEach((id, emitter) -> { try { emitter.send(SseEmitter.event().comment("ping")); } catch (IOException e) { emitters.remove(id); } }); }, 15, 15, TimeUnit.SECONDS);

注意:心跳线程和业务发送线程可能并发操作同一个 emitter。虽然SseEmitter内部有锁,但为了顺序可控,我建议把心跳也走同一个发送队列。不过实践中直接发也没出过问题,Spring 内部用的是synchronized。

2.4 显式调用的天花板在哪

SseEmitter用起来直观,但它的天花板很明显:每个连接占一个 Tomcat 线程。我做过压测,单机 4C8G,Tomcat 默认配置,SSE 连接数到 180 左右就开始出现新请求排队,到 200 直接拒绝。

调大server.tomcat.threads.max到 1000,能撑到 900 多连接,但内存占用涨了将近 1GB,而且 GC 压力明显变大。这就是平台线程的物理限制——线程是稀缺资源。

所以SseEmitter适合什么场景?内部管理系统、连接数几百以内的场景、或者你根本不在乎并发。一旦面向 C 端、要扛几千上万连接,就得换路子。

3. Spring AI 里的隐式封装:流式对话是怎么把 SSE 藏起来的

3.1 ChatClient 的 stream() 背后发生了什么

Spring AI 现在做流式对话,代码简单到令人发指:

@GetMapping(value = "/chat", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public Flux<String> chat(@RequestParam String message) { return chatClient.prompt() .user(message) .stream() .content(); }

一行stream()就搞定了。但这里面藏了好几层封装,值得拆开看:

第一层,ChatClient把模型调用抽象成了StreamResponseSpec。你调.stream()返回的是一个StreamResponseSpec,再调.content()拿到Flux<String>。这个Flux就是流式返回的文本块序列。

第二层,底层模型客户端(比如 OpenAI 的)把 HTTP 的 SSE 响应解析成 Flux。模型服务端返回的本来就是 SSE 格式(data: {...}\n\n),Spring AI 用 WebClient 接收,逐行解析,把每个 chunk 转成 Flux 的元素。

第三层,Spring MVC 把Flux<String>再序列化回 SSE 写给前端。注意这里有个"SSE 进、SSE 出"的转换——模型给我 SSE,我再给前端 SSE,中间用 Flux 做管道。

所以你在前端看到的,就是标准的 SSE 流:

data: 你 data: 好 data: , data: 我 data: 是 ...

3.2 为什么这里必须用 WebFlux 而不能用 MVC

上面那段代码,如果你在纯 Spring MVC 项目里跑,会报错或者行为异常。原因是Flux这种响应式类型,只有 WebFlux 才能正确处理它的背压和异步订阅。Spring MVC 虽然从 5.0 开始支持响应式返回值,但对Flux的支持是有限的——它会尝试把整个 Flux 收集完再返回,那就失去流式意义了。

所以 Spring AI 的流式对话,底层依赖 WebFlux。哪怕你的项目主体是 MVC,引入 Spring AI 后也会带上 WebFlux 的依赖。这一点在选型时要清楚。

那有没有办法在纯 MVC 里做流式对话?有,用SseEmitter手动桥接:

@GetMapping(value = "/chat", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter chat(@RequestParam String message) { SseEmitter emitter = new SseEmitter(0L); chatClient.prompt().user(message).stream().content() .subscribe( chunk -> { try { emitter.send(chunk); } catch (IOException e) { emitter.completeWithError(e); } }, emitter::completeWithError, emitter::complete ); return emitter; }

但这么写就回到了"一个连接一个线程"的老问题。而且subscribe是在 Reactor 的线程上跑的,和 MVC 的线程模型混在一起,调试起来很别扭。

3.3 流式对话里几个容易翻车的点

第一个,超时。大模型生成慢的时候,两个 chunk 之间可能隔好几秒。如果你的网关或者 Nginx 有 60 秒 idle timeout,长回答会被截断。我遇到过用户问一个需要长推理的问题,回答到一半连接断了,前端显示"stream disconnected before completion"。解决办法是调大网关超时,或者在应用层加心跳——但 Spring AI 的 Flux 管道里插心跳比较麻烦,通常还是靠调网关配置。

第二个,错误处理。模型调用失败时,Flux 会发onError。如果你没处理,前端收到的是一个断掉的流,用户看到的是"回答到一半没了"。正确做法是在 Flux 上加onErrorResume,把错误转成一条正常的文本消息推给前端:

.stream() .content() .onErrorResume(e -> Flux.just("\n[生成中断:" + e.getMessage() + "]"))

第三个,token 统计。流式模式下,模型的 usage 信息通常在最后一个 chunk 里(如果模型支持的话)。Spring AI 的ChatResponse里有 metadata,但.content()只取文本,拿不到 usage。要统计 token,得用.chatResponse()而不是.content(),然后从 metadata 里抠。

第四个,取消。用户点了"停止生成",前端关闭 EventSource,服务端要能感知并取消模型调用。WebFlux 里靠的是订阅取消(cancel()),但模型 API 的 HTTP 请求能不能真的取消,取决于底层客户端。这块 Spring AI 还在演进,实测有时候取消不彻底,模型还在后台跑。

3.4 隐式封装的代价:你失去了什么控制权

Spring AI 把 SSE 封装得很优雅,但代价是你对底层连接的控制权变弱了。比如:

  • 你想自定义 SSE 的 event name、id、retry 字段,.content()给不了你,得用.chatResponse()然后自己组装。
  • 你想在流中间插入自定义事件(比如"思考中"状态),得在 Flux 上做mergeWith。
  • 你想控制 flush 时机,基本没戏,Reactor 有自己的调度。

所以我的经验是:简单对话用.content()图省事,复杂交互(多模态、工具调用、状态提示)就得下沉到Flux<ChatResponse>甚至自己管 SseEmitter。封装是双刃剑,用之前想清楚自己要什么。

4. 虚拟线程:让 SSE 的并发能力真正起飞

4.1 虚拟线程到底解决了什么

JDK 21 的虚拟线程(Virtual Threads,JEP 444)是这几年 Java 最实在的一个特性。它的核心思想是:把"线程"这个概念的调度从操作系统交给 JVM。

平台线程是 1:1 映射到操作系统线程的,创建成本高(MB 级栈内存),数量有限(几千个就到顶)。虚拟线程是 M:N 映射——大量虚拟线程复用在少量平台线程(称为 carrier thread)上。虚拟线程的栈是存在堆里的,可以动态伸缩,创建一个虚拟线程的成本跟创建一个普通对象差不多。

关键在于:当虚拟线程遇到阻塞操作(IO、sleep、锁等待)时,JVM 会把它从 carrier thread 上卸载下来,让 carrier 去跑别的虚拟线程。等阻塞结束再挂回去。这样一来,一个 carrier thread 可以支撑成千上万个虚拟线程。

这对 SSE 意味着什么?每个 SSE 连接占一个虚拟线程,但不再占一个平台线程。一万个连接就是一万个虚拟线程,底层可能只用几十个 carrier thread。内存占用从 GB 级降到 MB 级。

4.2 Spring Boot 3.2 开启虚拟线程有多简单

一行配置:

spring.threads.virtual.enabled=true

就这一行。Spring Boot 3.2+ 会自动把 Tomcat 的请求处理线程池换成虚拟线程执行器。你的SseEmitter代码一个字不用改,并发能力直接上一个数量级。

我实测过,同样的 4C8G 机器,同样的SseEmitter代码:

配置稳定连接数内存占用新请求延迟
平台线程(默认 200)~180正常连接满后排队
平台线程(调到 1000)~900+1GB连接满后排队
虚拟线程10000++200MB基本无感

这个对比很能说明问题。虚拟线程不是"优化",是量级上的改变。

4.3 虚拟线程下的 SSE 代码要不要改

大部分情况不用改。但有几个细节要注意:

第一,synchronized会 pin 住 carrier thread。这是虚拟线程最大的坑。如果一个虚拟线程在synchronized块里阻塞,它没法被卸载,会一直占着 carrier thread。JDK 21 里这个问题还存在(JDK 24 的 JEP 491 才解决)。所以虚拟线程环境下,能用ReentrantLock就别用synchronized。

SseEmitter内部用的是synchronized,这意味着高并发下可能有 pinning 问题。实测在几千连接时影响不大,但上万连接时能观察到 carrier thread 被占满。如果遇到吞吐上不去,用jdk.tracePinnedThreads参数排查:

java -Djdk.tracePinnedThreads=full -jar app.jar

第二,ThreadLocal 要慎用。虚拟线程数量巨大,如果每个线程都存一份 ThreadLocal,内存会爆。Spring 的RequestContextHolder默认用 ThreadLocal,在虚拟线程下每个请求一份,量大时要注意。JDK 提供了ScopedValue(预览特性)作为替代,但还没正式落地。

第三,线程池不要乱用。虚拟线程环境下,Executors.newFixedThreadPool这种反而成了瓶颈。要提交任务用Executors.newVirtualThreadPerTaskExecutor(),每个任务一个虚拟线程。

4.4 虚拟线程 + SseEmitter 的完整生产配置

把前面的东西整合一下,一个生产可用的 SSE 配置大概长这样:

@Configuration public class SseConfig { @Bean(destroyMethod = "shutdown") public ScheduledExecutorService heartbeatExecutor() { return Executors.newScheduledThreadPool(2, r -> { Thread t = new Thread(r, "sse-heartbeat"); t.setDaemon(true); return t; }); } @Bean public SseEmitterManager sseEmitterManager(ScheduledExecutorService heartbeatExecutor) { return new SseEmitterManager(heartbeatExecutor); } }
public class SseEmitterManager { private final Map<String, SseEmitter> emitters = new ConcurrentHashMap<>(); private final ScheduledExecutorService heartbeatExecutor; public SseEmitterManager(ScheduledExecutorService heartbeatExecutor) { this.heartbeatExecutor = heartbeatExecutor; heartbeatExecutor.scheduleAtFixedRate(this::heartbeat, 15, 15, TimeUnit.SECONDS); } public SseEmitter register(String clientId) { SseEmitter emitter = new SseEmitter(30 * 60 * 1000L); // 30 分钟 emitter.onCompletion(() -> emitters.remove(clientId)); emitter.onTimeout(() -> { emitter.complete(); emitters.remove(clientId); }); emitter.onError(e -> emitters.remove(clientId)); emitters.put(clientId, emitter); return emitter; } public void send(String clientId, Object data) { SseEmitter emitter = emitters.get(clientId); if (emitter == null) return; try { emitter.send(data); } catch (IOException e) { emitter.completeWithError(e); emitters.remove(clientId); } } private void heartbeat() { emitters.forEach((id, emitter) -> { try { emitter.send(SseEmitter.event().comment("ping")); } catch (IOException e) { emitters.remove(id); } }); } }

配合application.properties:

spring.threads.virtual.enabled=true server.tomcat.threads.max=200 server.tomcat.accept-count=1000

注意server.tomcat.threads.max在虚拟线程模式下其实不太重要了,因为请求处理走的是虚拟线程执行器,这个配置主要影响的是 acceptor 和 poller 线程。但设小一点没坏处。

5. 三条路线的选型决策与实测对比

5.1 一张表说清楚怎么选

维度SseEmitter(平台线程)WebFlux + Flux虚拟线程 + SseEmitter
代码复杂度低高低
并发能力低(百级)高(万级)高(万级)
调试难度低高低
团队学习成本低高低
生态兼容性好一般(很多库不支持响应式)好
内存占用中低低
适合场景内部系统纯响应式架构大多数业务场景

我的建议很直接:新项目、JDK 21+、Spring Boot 3.2+,无脑上虚拟线程 + SseEmitter。除非你的团队已经深度使用 Reactor,否则没必要为了 SSE 去啃 WebFlux。

5.2 一个真实的迁移案例

去年我把一个内部监控系统的 SSE 推送从平台线程迁到虚拟线程。原系统用SseEmitter,Tomcat 线程池调到 500,稳定支撑 400 多个连接,再多就开始丢。迁移动作只有两步:

  1. 升级 JDK 到 21,Spring Boot 到 3.2
  2. 加一行spring.threads.virtual.enabled=true

改完压测,连接数直接干到 8000 没压力,内存反而降了(因为不用维护 500 个平台线程的栈了)。整个过程代码零改动,这就是虚拟线程最爽的地方——它是透明的。

但也不是完全没坑。迁移后我们发现一个定时任务偶尔会卡住,排查发现是任务里用了synchronized锁一个共享对象,虚拟线程被 pin 住了。改成ReentrantLock后解决。所以迁移后一定要用jdk.tracePinnedThreads跑一遍,把 pinning 点找出来。

5.3 压测数据背后的原理

为什么虚拟线程能撑这么多连接?核心在于内存模型。

平台线程的栈是预分配的,默认 1MB(-Xss),即使线程啥也不干,这 1MB 也占着。10000 个平台线程 = 10GB 栈内存,直接 OOM。

虚拟线程的栈是按需分配、存在堆里的,初始只有几百字节,随着调用深度增长。一个挂着的 SSE 虚拟线程,栈深度很浅,可能就几 KB。10000 个虚拟线程 = 几十 MB。这就是量级差异的来源。

再加上 carrier thread 的数量默认等于 CPU 核数(可以用jdk.virtualThreadScheduler.parallelism调整),4 核机器就 4 个 carrier,上下文切换成本极低。

5.4 什么时候虚拟线程也救不了你

虚拟线程不是银弹。以下场景它帮不上忙:

  • CPU 密集型任务:虚拟线程的优势在 IO 阻塞时卸载,纯计算任务没有阻塞点,和平台线程没区别。
  • 大量synchronized阻塞:pinning 会让 carrier 被占死,退化成平台线程。
  • native 方法阻塞:JNI 调用里的阻塞,JVM 感知不到,没法卸载。
  • 下游是瓶颈:如果 SSE 推送的数据来自一个慢查询数据库,虚拟线程再多也堵在数据库那。

所以上虚拟线程之前,先确认你的瓶颈真的在"线程数量"上,而不是别的地方。

6. 踩过的坑与排查手册

6.1 连接莫名其妙断掉:从网关到应用逐层排查

SSE 最烦的问题就是"连接断了但不知道谁断的"。我的排查顺序是这样的:

第一步,看前端报错。EventSource的onerror里能拿到readyState。如果是CLOSED,说明连接被关闭了。但前端拿不到关闭原因,得看服务端。

第二步,看服务端日志。onError回调有没有触发?触发了说明是客户端或中间层断的。没触发说明是服务端主动 complete 的(超时或者业务逻辑)。

第三步,看网关配置。Nginx 的proxy_read_timeout默认 60 秒,proxy_buffering默认 on。SSE 必须关 buffering,否则数据会被攒着一起发:

location /stream { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; proxy_set_header Connection ''; proxy_http_version 1.1; chunked_transfer_encoding off; }

第四步,看负载均衡。有些 LB 对长连接有硬性超时,比如 5 分钟。这个只能查 LB 文档或者抓包看。

第五步,看应用超时配置。SseEmitter的 timeout、Tomcat 的connectionTimeout、Spring 的async.request-timeout,任何一个到点都会断。

6.2 "stream disconnected before completion" 的三种成因

这个报错我在用 Spring AI 时遇到过好几次,成因有三种:

成因一:模型 API 超时。模型生成慢,超过了 HTTP 客户端(WebClient)的响应超时。解决:调大spring.ai.openai.chat.options.timeout或者底层 WebClient 的 timeout。

成因二:网关 idle timeout。两个 chunk 之间间隔太久,网关认为连接空闲就掐了。解决:调大网关超时,或者让模型尽快吐第一个 token。

成因三:应用层 Flux 被取消。前端关闭了连接,Reactor 取消订阅,模型调用被中断。这个是正常行为,但日志里看起来像错误。解决:在onErrorResume里区分CancellationException,别当错误处理。

6.3 虚拟线程 pinning 的定位方法

前面提过synchronized会 pin 住 carrier thread。定位方法:

java -Djdk.tracePinnedThreads=full -jar app.jar

启动后,一旦发生 pinning,控制台会打印完整的堆栈,告诉你哪个synchronized块导致的。输出长这样:

Thread[#123,ForkJoinPool-1-worker-1,5,CarrierThreads] java.base/java.lang.VirtualThread$VThreadContinuation.onPinned(VirtualThread.java:183) ... com.example.MyService.doSomething(MyService.java:42) <== synchronized block

看到<== synchronized block那行,就是罪魁祸首。改成ReentrantLock即可。

注意:JDK 21 里jdk.tracePinnedThreads是有效的,但 JDK 24 之后这个参数被移除了,因为 JEP 491 已经解决了大部分 pinning 问题。如果你用的是 JDK 24+,基本不用担心这个。

6.4 内存泄漏:emitter 没被清理

SSE 最隐蔽的 bug 是 emitter 泄漏。表现是:连接数看起来正常,但内存一直涨,最后 OOM。

原因通常是:客户端断开了,但服务端的 emitter 还挂在 Map 里没删。onCompletion和onError不一定及时触发,尤其是客户端异常断开(比如拔网线)时。

我的防御措施有三层:

  1. 每次 send 都 try-catch,抛异常立即从 Map 移除。
  2. 心跳时顺便清理,心跳 send 失败的 emitter 直接移除。
  3. 定期扫描,用一个定时任务检查 emitter 的最后活跃时间,超过阈值没活动的强制 complete 并移除。
// 定期清理僵尸连接 heartbeatExecutor.scheduleAtFixedRate(() -> { long now = System.currentTimeMillis(); emitters.entrySet().removeIf(entry -> { if (now - entry.getValue().lastActive > 5 * 60 * 1000) { entry.getValue().complete(); return true; } return false; }); }, 1, 1, TimeUnit.MINUTES);

7. 一些零散但有用的经验

7.1 前端 EventSource 的重连策略

EventSource自带重连,默认 3 秒一次。但它的重连是"无脑重连",不管服务端是不是真的挂了。如果服务端在重启,前端会疯狂重连,打爆服务端。

我的做法是:服务端在关闭连接前,发一个event: close事件,前端收到后主动close()并停止重连。等业务需要时再手动重连。

// 服务端优雅关闭 emitter.send(SseEmitter.event().name("close").data("server shutting down")); emitter.complete();
// 前端 const es = new EventSource('/stream'); es.addEventListener('close', () => { es.close(); // 延迟后手动重连 setTimeout(connect, 5000); });

7.2 SSE 的鉴权怎么做

EventSource有个限制:不能自定义请求头。所以你不能像普通 AJAX 那样带Authorization头。解决方案有三种:

  • Cookie:最省事,浏览器自动带。但跨域时要注意withCredentials。
  • URL 参数:/stream?token=xxx,简单但 token 会出现在日志里,安全性差。
  • 先 POST 换 ticket,再用 ticket 建 SSE 连接:最安全,但多一次请求。

我一般用 Cookie,配合SameSite=Lax,够用。

7.3 数据格式的选择

SSE 的data字段只能是字符串。要传对象,得自己序列化。我一般统一用 JSON:

emitter.send(SseEmitter.event() .name("message") .data(objectMapper.writeValueAsString(payload)));

前端解析:

es.addEventListener('message', (e) => { const data = JSON.parse(e.data); // ... });

注意data里如果有换行符,SSE 协议会把它拆成多行data:,前端拼回来时换行会丢。所以 JSON 序列化时最好把换行转义掉,或者用 base64。

7.4 和 WebSocket 的取舍再强调一次

最后再说一次选型。SSE 和 WebSocket 不是替代关系,是互补:

  • 只需要服务端推:SSE。简单、走 HTTP、自动重连。
  • 需要双向通信:WebSocket。比如聊天室、协同编辑。
  • 需要传二进制:WebSocket。SSE 只能传文本。
  • 需要极低延迟:WebSocket。SSE 有 HTTP 头开销。

大模型对话这种场景,用户输入走普通 POST,模型输出走 SSE,完美契合。没必要上 WebSocket。

8. 写在最后的一点个人体会

折腾 SSE 这几年,我最大的感受是:技术选型要看"约束条件",而不是"哪个更先进"。WebFlux 比 MVC 先进,但团队不熟就是灾难。虚拟线程比平台线程先进,但 JDK 版本不够就是空谈。

虚拟线程真正改变游戏规则的地方,不是它"更快",而是它让同步写法重新变得可行。以前为了高并发,被迫学 Reactor、学响应式、学各种操作符,现在一行配置就回到熟悉的同步世界。这对大多数业务团队来说,是实实在在的减负。

Spring AI 的封装也是同理。它把 SSE 的复杂度藏起来,让你专注在业务逻辑上。但封装总有边界,遇到它解决不了的问题,你还是得下沉到Flux甚至SseEmitter。所以理解底层原理,比会用 API 重要得多。

如果你现在正在做 SSE 相关的功能,我的建议是:先用SseEmitter把功能跑通,理解连接生命周期和心跳机制;然后上虚拟线程解决并发;最后如果业务复杂到需要响应式,再考虑 WebFlux。别一上来就追求"最优解",能跑通、能维护、团队能接手的,才是好方案。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询