SSE上生产必踩的坑:断线重连、心跳超时与SseEmitter实战
2026/9/16 14:58:02 网站建设 项目流程

你写的SSE连本地都跑得通,一上生产就断——问题出在哪

先看一段挺常见的Java服务端代码,SseEmitter往外吐数据,本地联调好好的,丢到生产环境跑几分钟,前端就收不到消息了。打开服务端日志,经常能看到这么一句:

before completion: idle timeout waiting for sse

这句话我印象太深了。第一次遇到的时候,我以为是前端断开了,查了很久才发现根本不是。这是服务端在说:连接还开着,但是太长时间没有数据往上面写了,我自己把它关了。

这个场景,做Java后端的朋友应该都不陌生。SSE(Server-Sent Events,服务端推送事件)这种东西,看着特别简单——比WebSocket轻量,走普通HTTP,不用额外握手协议,前端一个EventSource就接上了。但只要上过生产,你很快就会撞上三个问题:连接被超时掐断、断线后没人管、推送多了把服务拖垮。标题我说90%的Java人写的SSE上不了生产,一点都不夸张,因为这三个问题不解决,任何一个都能让你的服务在下一次发版后无声无息地挂掉。

这篇内容,不聊SSE的基础概念,直接讲生产环境怎么把SSE写到能用的程度。包含断线重连的完整链路、超时降级方案、服务端Session的管理和排查手段,后面还会附一个可以抄的代码骨架。无论你是准备面试(这个方案在面试里确实很能打),还是手头有个项目正被SSE问题折磨,都值得看完。

1. 90%的SSE死在生产环境的三个典型症状

先说结论:SSE不好用,问题多半不在SSE协议本身,而在实现SSE的时候把人家的生命周期和连接模型理解错了。我见过太多项目,SSE服务写出来就三个功能:SseEmitter创建、send()发消息、onCompletion()打日志。台词我都背下来了:“SSE嘛,就是个单向长连接,简单得很。”

结果上线后会以各种姿势翻车。我总结下来,90%的翻车场景逃不出下面三种。

1.1 第一种死法:连接静默被服务端掐断

这就是前面说的idle timeout问题。SSE本质是HTTP长连接。主流的Java容器——Tomcat也好,Jetty也好——对HTTP连接都有一个空闲超时限制。意思是你这个连接如果在一段时间内没有数据写入,容器出于资源保护的目的,会主动把连接断开。

问题在于:SSE这种长连接,业务上往往不是每时每刻都有数据要推。比如一个订单状态通知服务,可能用户下单后10分钟商家才接单,中间这段时间连接上没有任何字节流动。如果你没意识到容器的空闲超时机制存在,就会出现在推送完第一条数据后,连接在某个时刻被容器默默关掉。前端感知到连接关闭,EventSource会自动重连(这是浏览器行为),但服务端这边的SseEmitter已经走到onCompletion回调了,新连接重新连上来又创建一个新的SseEmitter。看起来“好像能工作”,其实你的推送链路已经断过好几次了。

更隐蔽的是,前端重连时如果不带Last-Event-ID头,服务端就不知道你断到了哪里,重连成功后只能从当前时刻开始推,断线期间的消息直接丢了。用户端的表现就是:界面上的状态突然卡住不动,刷新页面又好了。

1.2 第二种死法:断开后没有任何收敛策略

前端页面关了,或者用户网络切了一下,TCP连接其实已经断了。但服务端不一定能立刻感知到。TCP连接在没有数据传输的情况下,断开是察觉不到的——这在网络术语里叫“半开连接”。

如果推送量不大,且你写代码的时候小心处理了onErroronCompletion,这种半开连接还不会立刻引起问题。但如果你的SSE接口在给大量客户端推送高频数据,比如行情推送、协同编辑里的光标同步,那半开连接就会逐渐累积。每个SseEmitter实例都占着内存,有的还注册了定时器线程,量一大,GC压力就上来了。最恐怖的是每次调send()的时候对已断开的连接发送数据,会抛异常,异常处理不得当的话,推送主线程会被拖垮。

1.3 第三种死法:推送模型设计成散弹枪

很多团队的SSE推送接口,是把“用户标识”和“连接”直接硬编码在业务代码里的。业务A里要推送,直接塞一个SseEmitter到某个静态Map;业务B里又要推,又搞了另一套Map。推送时全量遍历,管你对不对,每个都发一遍。

你说这不也能跑吗?能跑,但问题在扩容和故障转移的时候全暴露出来。你有两台机器,用户在机器1上建了SSE连接,请求打到机器2想推数据时,从机器2的Map里找不到连接,推送失败。然后你会怎么处理?很多人的解法是:让机器1和机器2的Map做个同步,或者用Redis发布订阅广播,“反正量也不大,全广播也没事”。结果就是随着连接数上涨,消费SSE连接的那台机器成了单点,并且广播风暴把整个集群拖死。

这三种死法,本质上是同一个问题:把SSE当成HTTP接口里的一个返回类型,而不是当成一个有状态的长连接会话去对待。SSE要上生产,首先得从心态上完成这个转变。

2. 从SseEmitter的一次性生命周期说起:源头参数藏了太多坑

Spring的SseEmitter封装了SSE服务端的绝大部分逻辑。你用起来感觉很简单,但它的构造参数、回调时机、线程模型,每一样都要弄明白。你创建的每一个SseEmitter实例,背后都绑定了一条真实的TCP连接,而这个连接的生命周期并不由你的业务代码完全掌控。

2.1 三个构造参数到底动了谁的蛋糕

SseEmitter有三个构造参数:

  • timeout:连接超时时间,默认DeliveryConstants里给的是30秒?不,Spring默认给了30秒还是0?这里很多文章写错了。准确说是SseEmitter默认的timeout是0L,表示不超时。但坑就在这——你设了0,只是告诉Spring“嘿,我不想让它超时”,但如果底层容器(Tomcat的keepAliveTimeout、Jetty的idleTimeout)有自己的半关闭策略,两者就会发生冲突。就像你和同事约好“这个会议我全程不出声”,但你没法保证会议室的门不被管理员锁上。

实际生产环境中,我强烈建议给timeout设一个明确的值,比如30秒或60秒,而不是依赖默认行为。设明确值,配合心跳机制,等于你把连接的管理权从容器手里抢回来了,什么时候断什么时候续,你说了算。

  • reconnectTime:这个是服务端建议的重连间隔,单位毫秒。它会映射到SSE协议里的retry:字段。前端EventSource在收到这个字段后,断开重连时就会按这个间隔来。注意,只是“建议”,浏览器不一定完全遵守,但设置总比不设置好。

  • initializationTimeout:这个参数在最新版本里才有意义,控制SseEmitter初始化阶段的超时。早期版本调试时你可能会碰到“Initialization timeout”的错,就是它引起的。生产环境一般不用动,但出现初始化超时异常时要知道是这个参数在管。

这三个参数,说穿了就一件事:让服务端控制连接的生命周期。想明白这点,你就不会问“SseEmitter的timeout到底设多少好”这种问题了——它会为你后面设计心跳策略服务。

2.2 回调函数里,最不该干的事情是“在onTimeout里发消息”

onCompletiononTimeoutonError这三个回调,是新手最容易被蛊惑的地方。

onCompletion:连接正常关闭(客户端断开、服务端主动complete())后触发。注意,它是“最终一定会触发”的,但触发晚不代表连接还活着。你在这个回调里清理资源没问题,但如果试图判断“用户是否真的断开了”?做不到,因为客户端断开和服务端主动完成走的是同一条路径。

onTimeout:服务端因超时触发连接关闭时调用。很多第一次写SSE的同学,发现连接超时关闭后,前端会重新连上来,就想着在onTimeout里“续命”——比如重新创建一个SseEmitter塞回Map里,假装连接没断过。这种写法绝对是大忌。onTimeout触发时,底层连接已经进入关闭流程,你塞回去的新实例,将来能不能收到数据全靠运气,而且在极端情况下会引起连接泄漏。

onError:链路出现异常时触发。这个回调最大的问题是,异常类型千奇百怪,IOExceptionIllegalStateException都有可能出现,但因为回调里拿不到具体的上下文信息,你能做的只是把日志打全,然后等onCompletion

我见过比较规范的写法是:在回调里做两件事,一是从管理Map里移除当前这个SseEmitter,二是把“用户xxxx已断开”的日志打出来。就这么简单。任何创建新连接、发补偿消息的逻辑,都必须挪到更上层的“断线重连机制”里去处理,不能堆在回调里。

2.3 连接都是有状态的,别当一次性接口

有的同学喜欢这样写:

@GetMapping("/events") public SseEmitter events() { long userId = CurrentUser.get(); SseEmitter emitter = new SseEmitter(0L); // 每次都新建,推完就扔 return emitter; }

这在Demo里完全没问题,但生产环境你就等着被坑。SSE连接是“有状态”的,你不仅需要一个地方把它存起来(通常是ConcurrentHashMap),还要约定好什么时候删、什么时候换、什么时候补偿。否则,你连“这个用户现在到底连着没有”都不知道,就遑论什么断线重连和超时降级了。

3. 断线重连:服务端要做的,远不止“让前端重连”

SSE协议本身自带断线重连能力。浏览器端的EventSource在连接断开后会自动发起重连,而且会把服务端返回的retry:字段作为重连间隔。很多人觉得“这机制不是现成的嘛,还让我写什么?”关键就出在:自动重连只是把物理连接恢复给你,但你的业务推送从哪里续上,服务端完全不知道。

这就要说到SSE协议里一个非常重要的字段:Last-Event-ID

3.1 用Last-Event-ID实现断点续传

SSE规范里,客户端重连时可以在请求头里带上Last-Event-ID,告诉服务端“我最后收到的事件ID是多少,从这里之后的都给我补发”。这玩意儿就是你做断线补偿的基石。

流程是这样的:

  1. 服务端推送每条消息时,给消息带上id字段,比如订单ID + 自增序号
  2. 前端在onmessage回调里记录最近收到的事件ID。
  3. 连接断开后,前端重连时在请求头里带上Last-Event-ID
  4. 服务端解析这个Header,从事件源(数据库、MQ、Redis)里把ID之后的消息捞出来补推。

Spring MVC里获取这个Header的方式,跟普通接口没区别:

@GetMapping(value = "/subscribe", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe(@RequestHeader(value = HttpHeaders.LAST_EVENT_ID, required = false) String lastEventId) { // lastEventId就是前端带上来的断点标识 // 后续业务侧根据它做消息补偿 return sseService.createEmitter(userId, lastEventId); }

很多生产事故就是栽在这一步:没存ID,或者存了没往前传,导致重连后服务端接着从当前时刻推,断线窗口的数据永久性丢失。你看,断线重连不是“前端会自动重连”就完事了的,服务端必须把断点续传的能力实现好。

3.2 心跳机制:不是可选项,是必选项

前面说的容器空闲超时,就是心跳机制要解决的核心问题。SSE连接长时间没有数据流动,会被底层的NIO超时机制误杀。要规避这个问题,最简单有效的方式就是服务端每隔固定时间发送一个注释行。

心跳包在SSE协议里可以是一行以冒号开头的注释:

: heartbeat

也可以是显式的事件:

event: heartbeat data: ping

两者的区别在于:注释行触发不了前端的事件监听器,对业务逻辑没有干扰,是纯“保活”;而显式事件会触发前端的onmessage,如果前端没做过滤,业务层面可能会误处理。所以生产环境下,我推荐直接用注释行做心跳。

Spring这边,推荐用SseEmitter配合一个定时任务:

private ScheduledExecutorService scheduler = Executors.newSingleThreadScheduledExecutor(); private static final long HEARTBEAT_INTERVAL = 20L; // 每20秒发一次心跳 public SseEmitter createEmitter(Long userId, String lastEventId) { SseEmitter emitter = new SseEmitter(60_000L); // 保存连接 SSE_CONNECTIONS.put(userId, emitter); // 启动心跳任务 scheduler.scheduleAtFixedRate(() -> { try { emitter.send(SseEmitter.event().comment("heartbeat")); } catch (IOException e) { // 这条连接已经断了,取消心跳,交给清理流程 emitter.completeWithError(e); SSE_CONNECTIONS.remove(userId); } }, HEARTBEAT_INTERVAL, HEARTBEAT_INTERVAL, TimeUnit.SECONDS); // 确保心跳任务在连接关闭后停止 emitter.onCompletion(() -> { SSE_CONNECTIONS.remove(userId); // 注意:这里拿不到定时任务的引用,需要在实现中用更优雅的方式管理 }); return emitter; }

这段代码有一个隐藏问题:定时任务本身的停摆管理。上面这个写法里,定时任务没有绑定到具体的SseEmitter生命周期,一旦连接关闭但任务还在跑,就会不断往已关闭的emitter上发数据,产生IOException。生产级实现需要把每个连接的定时任务单独管理起来,连接关闭时取消任务。这是我在很多团队Code Review时反复强调的点。

心跳间隔怎么定?一个稳妥的参考公式是:心跳间隔 < 容器的空闲超时时间的一半。比如Jetty默认idleTimeout是30秒,那你心跳至少得15秒一次;如果Tomcat的keepAliveTimeout比你预想的长,心跳设20秒也行。原则只有一个:宁可多打几次,不能断。

3.3 重连风暴:当所有客户端同时断线

想象这个场景:服务器重启了,或者网络抖动持续了1分钟,你的生产环境有2万个在线客户端。恢复后,2万个客户端几乎同时触发重连,全部向服务端发起HTTP请求。服务端瞬间被打满,CPU飙高,然后又挂一轮。这就是重连风暴。

EventSource虽然会按retry:字段指定的间隔重连,但并不会做全网同步的“错峰”。而且如果服务端自己不设置retry:字段,浏览器默认的重连间隔在1~3秒之间,很容易造成同步冲击。

应对策略有两个层面:

  1. 服务端下发一个合理的retry:时间,比如10秒或15秒。这样能减小客户端同时重连的概率。
  2. 客户端在收到断线信号后,不要立刻重连,加一个随机抖动(jitter)。
// 前端伪代码 function connect() { const es = new EventSource('/subscribe'); es.onerror = () => { es.close(); const delay = 1000 + Math.random() * 5000; // 1~6秒随机 setTimeout(connect, delay); }; } connect();

服务端也可以做防护:在创建SseEmitter的接口上加一个“连接速率限制”,比如每个IP每分钟最多创建N条连接,超出直接返回429。这不是针对正常用户的,而是防止客户端代码写错导致死循环重连(我确实见过有人在前端把onerror里重连写成无限递归的)。

4. 超时降级:这把瑞士军刀,保住你的用户体验

超时降级往往是SSE方案里最容易被忽视的一环。很多人想的是:“SSE就是实时推送,延时低,要我降级到轮询?那还要SSE干嘛?”这话说的没毛病,但你没搞懂超时降级的本质。

超时降级不是让你放弃实时性,而是让你在连接不可用的时候,用户不至于什么都拿不到。

4.1 降级方案一:服务端缓存最近事件 + 重建后补偿

这是对SSE最自然的降级。上文提到的Last-Event-ID,本质就是服务端缓存(或可查询的持久化存储)+ 重建补偿。当连接断开时间过长,客户端重连后,服务端一看lastEventId落后太多——比如落后了10分钟的消息量——立即从消息表里把这段时间的数据全部捞出来推给前端。

这个方案的关键设计在于“消息表”和“分页快照”:

  • 每次推送的时候,除了发实时数据,还要顺带把这批数据写入一张recent_events表(或者Redis里的一个定长List)。
  • 重连时,根据lastEventId去查表。

这个方案能兜住绝大多数断线场景。真正高频情况下产生的大量消息,可以设置一个“最大补偿条数”,超出部分让前端走一次全量刷新,别试图在网络抖动的情况下传几千条数据过去。

4.2 降级方案二:短轮询兜底

如果服务端判断某台机器已经过载,或者某个用户的SSE连接反复建立又断开(比如弱网环境),最理智的做法不是一遍遍重建连接,而是主动降级成普通轮询。

具体操作给两个入口:

  1. 服务端返回特定的HTTP状态码,比如503,前端收到后自动切换成短轮询模式。
  2. 前端根据自身的重连失败次数,本地决定切换轮询。
let failCount = 0; const MAX_FAIL_COUNT = 3; function connect() { const es = new EventSource('/subscribe'); es.onerror = () => { failCount++; es.close(); if (failCount >= MAX_FAIL_COUNT) { // 切换为短轮询 startPolling(); } else { setTimeout(connect, 1000); } }; }

你可能会问:既然轮询也能拿数据,为什么不直接用轮询?这个问题问到点子上了。轮询的问题是浪费,才给了SSE存在的价值。所以这里的设计是“降级”,大部分时候走SSE,只有连接不稳定才短暂切到轮询,等网络恢复后,再平滑升级回SSE。

4.3 降级方案三:WebSocket?不是不行,但要谨慎提

面试的时候我经常会被追问:“既然断线重连这么麻烦,为什么不用WebSocket?”这是个好问题。答案分两层。

WebSocket确实是全双工的,双向通信能力碾压SSE,但代价是:

  1. 服务端实现复杂度上升一个量级。
  2. 需要额外处理WebSocket的握手、心跳(WebSocket本身也没有标准心跳,应用层自己实现)、连接恢复(WebSocket在较老标准里没有原生的断线重连字段,需要完全自研)。
  3. 在云原生环境里,WebSocket对网关的亲和性要求比SSE高得多。

所以我的建议是:如果你的场景只是服务端单方面推送、客户端被动接收,SSE是更轻的正确选择。不要为了炫技引入WebSocket,WebSocket只有在真的需要双向通信时才是正确的工具。

5. 还有几个容易被面试官追到死的细节

有一些细节,代码写起来就一两行,但面试和线上排查时总能炸出问题来。这里我挑四个高频的讲,每一个都是我在真实项目中踩过或者看别人踩过的。

5.1 连接数管理:一个用户一条连接,够了吗

很多方案里,连接Map都长这样:

private final ConcurrentHashMap<Long, SseEmitter> connections = new ConcurrentHashMap<>();

key是用户ID,value是SseEmitter。这个设计有个天然缺陷:一个用户多开几个标签页,后建立的连接会覆盖先建立的连接,导致前一个标签页收不到推送。

更好一点的设计是用连接ID作为key,用户ID作为value的属性之一:

public class SseConnection { private String connectionId; // 全局唯一 private Long userId; private SseEmitter emitter; private long lastHeartbeatTime; private String lastEventId; }

这样,一个用户可以有多个连接,推送时通过用户ID找到该用户的所有连接,批量发送。运维层面也能按连接维度看数据。

再往深一层,连接Map不能无限膨胀。每个连接需要记录lastHeartbeatTime,一个后台扫地线程定期清理掉超过N分钟没有心跳的连接。这在生产环境是大规模SSE连接保持健康的关键。

5.2 线程池隔离:推送别和业务线程抢资源

刚开始写SSE推送时,很多人直接用业务线程去send()。如果IO出问题,业务接口也跟着遭殃。我踩过一个挺典型的坑:一个订单状态的推送接口,send()时连接刚好被容器超时断掉,抛了异常,那个异常直接污染了业务方法的事务,导致订单回滚。排查时半天没找到原因,因为日志只能看到“订单更新失败”,谁会想到是推送的问题。

后来我把所有推送逻辑全部放到一个独立线程池里:

private final ThreadPoolExecutor sseSendPool = new ThreadPoolExecutor( 8, 16, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(1000), new ThreadPoolExecutor.CallerRunsPolicy() ); public void sendToUser(Long userId, Object message) { List<SseConnection> conns = connections.getByUserId(userId); for (SseConnection conn : conns) { sseSendPool.execute(() -> { try { conn.getEmitter().send(SseEmitter.event().data(message)); } catch (IOException e) { // 推送失败,标记这条连接待清理 markForRemoval(conn.getConnectionId()); } }); } }

线程池的拒绝策略我选的是CallerRunsPolicy。这个策略在推送量大的时候,会让调用线程自己执行发送逻辑,等于是“你排不上队就自己跑”,能起到背压的效果,避免任务无限堆积。当然,具体策略得看你们的业务容忍度,AbortPolicy会直接抛异常,DiscardOldestPolicy则可能丢弃最新的推送任务。

5.3 SSE的鉴权:别把连接暴露成一个裸接口

SSE接口如果直接挂在网关后面,普遍会遇到两个问题:

  1. 网关的超时配置可能会掐断SSE,因为反向代理默认对长时间没响应的连接是有限制的。
  2. SSE连接建立后,没法用普通的Header动态鉴权,所以得在建立连接的时候把鉴权做完。

我的标准做法是:客户端先调一个普通的登录/鉴权接口,获取一个短期的SSE token,再拿这个token去建立SSE连接:

@GetMapping(value = "/subscribe", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe(@RequestParam("token") String token) { // token校验 SseAuthInfo authInfo = sseTokenService.validate(token); if (authInfo == null) { throw new ResponseStatusException(HttpStatus.UNAUTHORIZED, "invalid sse token"); } return sseService.createEmitter(authInfo.getUserId(), authInfo.getLastEventId()); }

这里有个细节:token不仅要校验有效性,还要绑定用户ID和最后事件ID。否则即使你校验通过了,拿到的userId和客户端自称对不上,照样有越权风险。我见过有团队把userId直接放在路径参数里,然后只校验了token存在,登进去之后能推送别人的消息,这就非常危险了。

5.4 容器与网关:Tomcat/Jetty的配置决定了你的一半命运

前文反复提到idleTimeout,这里是真正需要检查的地方。以Spring Boot内嵌Tomcat为例,需要关注两个配置:

  • server.tomcat.keep-alive-timeout:Tomcat对keep-alive连接的保活时长,默认可能比你想的小。
  • server.tomcat.connection-timeout:建立连接时的超时。

内嵌Jetty的话,是server.jetty.idle-timeout

但是——重点来了——很多生产环境是双层的,前面还挡了一层Nginx。Nginx的proxy_read_timeout默认60秒。如果你的SSE心跳间隔是20秒,Nginx这层没问题;但如果你没做心跳,连接超过60秒没有数据,Nginx直接把连接断开。这就是为什么我前面说“心跳间隔要设计好,保证任意一个中间层的超时都不会被触发”。

Nginx的配置里,这几个值要调:

proxy_set_header Connection ''; proxy_http_version 1.1; proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; proxy_send_timeout 3600s;

proxy_buffering一定要关。否则Nginx可能会先把SSE响应缓冲起来,导致前端收到的是攒了一阵子才发出来的数据,实时性变成“伪实时”。

5.5 多实例部署时,怎么保证推送能找到连接

前面提到多实例部署时,连接只存在某一台机器的内存里,别台机器推不到。解决这个问题业界有两个主流思路:

思路一:Redis发布订阅。每台机器都订阅同一个频道,某台机器想推送时,把消息发布到Redis,所有实例收到后各自查自己的本地连接Map,谁有连接谁推送。

// 消息发布端 redisTemplate.convertAndSend("sse:push", jsonMessage); // 消息订阅端(每台机器) container.addMessageListener((message, pattern) -> { PushMessage msg = JSON.parseObject(message.getBody(), PushMessage.class); List<SseConnection> conns = localConnections.getByUserId(msg.getUserId()); for (SseConnection conn : conns) { sendFromPool(conn, msg.getData()); } }, topic);

这个方案的弊端是:如果在线用户全挂在机器1上,你发一条广播,机器2有Redis订阅也会消费,但查不到本地连接,就白白浪费一次处理。不过它胜在简单可靠。

思路二:网关层做连接路由。所有SSE连接都打到一台状态独立的Gateway节点,由Gateway记录每个用户的连接落在哪台业务机。推送请求先到Gateway,查路由表后再转发到对应节点。这就引入了新的组件复杂度,适合连接规模很大的团队。

对于中小团队,我建议先用Redis发布订阅,把成本控制住,真到了要优化推送路径了再考虑路由网关。别一上来就搞微服务治理级别的架构。

6. 手写一个可复用的生产级SSE推送服务骨架

前面说了大量理论,这一节直接给出一套能在Spring Boot 2.7+/3.x环境跑通的生产级代码骨架。这个骨架摈弃了花哨的抽象,核心目标是:缓解断线重连的补偿、防止超时断连、方便运维排查。

整体结构分三个类:主服务类、连接管理器、推送消息DTO。

6.1 主服务类:SsePushService

package com.example.sse.demo; import org.springframework.http.MediaType; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Service; import java.io.IOException; import java.util.List; import java.util.Map; import java.util.concurrent.*; @Service public class SsePushService { private static final Logger log = LoggerFactory.getLogger(SsePushService.class); // 连接管理器 private final SseConnectionManager connectionManager = new SseConnectionManager(); // 推送专用线程池 private final ThreadPoolExecutor sendPool = new ThreadPoolExecutor( 8, 16, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(2000), new ThreadPoolExecutor.CallerRunsPolicy() ); // 心跳定时线程池,公共使用 private final ScheduledExecutorService heartBeatScheduler = Executors.newScheduledThreadPool(2); // 心跳间隔(秒) private static final long HEARTBEAT_INTERVAL_SECONDS = 20L; public SseEmitter createEmitter(Long userId, String lastEventId) { SseEmitter emitter = new SseEmitter(60_000L); String connectionId = userId + "-" + System.currentTimeMillis() + "-" + ThreadLocalRandom.current().nextLong(1000, 9999); SseConnection conn = SseConnection.builder() .connectionId(connectionId) .userId(userId) .emitter(emitter) .lastEventId(lastEventId) .lastHeartbeatTime(System.currentTimeMillis()) .build(); // 注册回调 emitter.onCompletion(() -> { log.info("[SSE] connection completed, connectionId={}", connectionId); connectionManager.remove(connectionId, conn); }); emitter.onTimeout(() -> { log.warn("[SSE] connection timeout, connectionId={}", connectionId); emitter.complete(); }); emitter.onError(ex -> { log.error("[SSE] connection error, connectionId={}, error", connectionId, ex); connectionManager.remove(connectionId, conn); }); // 保存连接 connectionManager.add(conn); // 启动心跳任务(注意:这里用scheduleAtFixedRate容易堆积,如果机器卡顿,建议用scheduleWithFixedDelay) String heartbeatTaskKey = connectionId; heartBeatScheduler.scheduleWithFixedDelay(() -> sendHeartbeat(conn), HEARTBEAT_INTERVAL_SECONDS, HEARTBEAT_INTERVAL_SECONDS, TimeUnit.SECONDS); log.info("[SSE] connection created, userId={}, connectionId={}", userId, connectionId); return emitter; } private void sendHeartbeat(SseConnection conn) { if (conn.isClosed()) { return; } try { conn.getEmitter().send(SseEmitter.event().comment("heartbeat")); conn.setLastHeartbeatTime(System.currentTimeMillis()); } catch (IOException e) { // 连接已经失效,主动关闭并清理 log.warn("[SSE] heartbeat failed, closing connection, connectionId={}", conn.getConnectionId()); try { conn.getEmitter().completeWithError(e); } catch (Exception ignore) { // completeWithError内部可能还会抛 } connectionManager.remove(conn.getConnectionId(), conn); } } // 推送给指定用户的所有连接 public void pushToUser(Long userId, Object data) { List<SseConnection> connections = connectionManager.getByUserId(userId); if (connections == null || connections.isEmpty()) { log.debug("[SSE] no active connection for userId={}", userId); return; } for (SseConnection conn : connections) { sendPool.execute(() -> { try { conn.getEmitter().send(SseEmitter.event().data(data)); } catch (IOException e) { log.warn("[SSE] send failed, connectionId={}, msg={}", conn.getConnectionId(), e.getMessage()); connectionManager.remove(conn.getConnectionId(), conn); } }); } } // 推送给所有连接(谨慎使用) public void pushToAll(Object data) { connectionManager.getAll().forEach(conn -> { sendPool.execute(() -> { try { conn.getEmitter().send(SseEmitter.event().data(data)); } catch (IOException e) { connectionManager.remove(conn.getConnectionId(), conn); } }); }); } // 根据userId主动断开连接 public void disconnectUser(Long userId) { List<SseConnection> connections = connectionManager.getByUserId(userId); connections.forEach(conn -> { try { conn.getEmitter().complete(); } catch (Exception ignore) { } connectionManager.remove(conn.getConnectionId(), conn); }); } }

6.2 连接管理器:SseConnectionManager

package com.example.sse.demo; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.stereotype.Component; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CopyOnWriteArrayList; import java.util.stream.Collectors; @Component public class SseConnectionManager { private static final Logger log = LoggerFactory.getLogger(SseConnectionManager.class); // 核心存储:连接ID -> 连接对象 private final Map<String, SseConnection> connectionMap = new ConcurrentHashMap<>(); // 辅助索引:用户ID -> 连接ID列表(便于按用户批量推送) private final Map<Long, List<String>> userIndex = new ConcurrentHashMap<>(); public void add(SseConnection conn) { connectionMap.put(conn.getConnectionId(), conn); userIndex.computeIfAbsent(conn.getUserId(), k -> new CopyOnWriteArrayList<>()).add(conn.getConnectionId()); } public void remove(String connectionId, SseConnection conn) { SseConnection removed = connectionMap.remove(connectionId); if (removed != null) { List<String> ids = userIndex.get(removed.getUserId()); if (ids != null) { ids.remove(connectionId); if (ids.isEmpty()) { userIndex.remove(removed.getUserId()); } } log.info("[SSE] connection removed, connectionId={}, userId={}", connectionId, removed.getUserId()); } } public List<SseConnection> getByUserId(Long userId) { List<String> ids = userIndex.get(userId); if (ids == null || ids.isEmpty()) { return List.of(); } return ids.stream() .map(connectionMap::get) .filter(c -> c != null && !c.isClosed()) .collect(Collectors.toList()); } public List<SseConnection> getAll() { return connectionMap.values().stream() .filter(c -> !c.isClosed()) .collect(Collectors.toList()); } public int size() { return connectionMap.size(); } }

6.3 连接对象和消息DTO

package com.example.sse.demo; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; public class SseConnection { private final String connectionId; private final Long userId; private final SseEmitter emitter; private volatile String lastEventId; private volatile long lastHeartbeatTime; private volatile boolean closed; private SseConnection(Builder builder) { this.connectionId = builder.connectionId; this.userId = builder.userId; this.emitter = builder.emitter; this.lastEventId = builder.lastEventId; this.lastHeartbeatTime = builder.lastHeartbeatTime; } public static Builder builder() { return new Builder(); } public String getConnectionId() { return connectionId; } public Long getUserId() { return userId; } public SseEmitter getEmitter() { return emitter; } public String getLastEventId() { return lastEventId; } public void setLastEventId(String lastEventId) { this.lastEventId = lastEventId; } public long getLastHeartbeatTime() { return lastHeartbeatTime; } public void setLastHeartbeatTime(long lastHeartbeatTime) { this.lastHeartbeatTime = lastHeartbeatTime; } public synchronized boolean isClosed() { return closed; } public synchronized void markClosed() { this.closed = true; } public static class Builder { private String connectionId; private Long userId; private SseEmitter emitter; private String lastEventId; private long lastHeartbeatTime; public Builder connectionId(String connectionId) { this.connectionId = connectionId; return this; } public Builder userId(Long userId) { this.userId = userId; return this; } public Builder emitter(SseEmitter emitter) { this.emitter = emitter; return this; } public Builder lastEventId(String lastEventId) { this.lastEventId = lastEventId; return this; } public Builder lastHeartbeatTime(long lastHeartbeatTime) { this.lastHeartbeatTime = lastHeartbeatTime; return this; } public SseConnection build() { return new SseConnection(this); } } }

消息结构建议直接使用JSON字符串,便于前端解析和后续加字段:

{ "eventId": "order-1001-1", "timestamp": 1712563200000, "eventType": "ORDER_STATUS_CHANGE", "data": { "orderId": 1001, "status": "PAID" } }

eventId是断线补偿的核心凭据,每次推送由业务方生成,保证在同一业务会话内唯一即可。

6.4 控制器接入

package com.example.sse.demo; import org.springframework.http.MediaType; import org.springframework.web.bind.annotation.*; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; @RestController @RequestMapping("/api/sse") public class SseController { private final SsePushService ssePushService; public SseController(SsePushService ssePushService) { this.ssePushService = ssePushService; } @GetMapping(value = "/subscribe", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe(@RequestParam("token") String token) { // 伪代码:校验token,获得userId和lastEventId SseAuthInfo auth = sseAuthService.validate(token); if (auth == null) { throw new ResponseStatusException(org.springframework.http.HttpStatus.UNAUTHORIZED); } return ssePushService.createEmitter(auth.getUserId(), auth.getLastEventId()); } // 内部测试用:主动推一条消息给用户 @PostMapping("/push/{userId}") public String push(@PathVariable Long userId, @RequestBody Object data) { ssePushService.pushToUser(userId, data); return "ok"; } }

6.5 这套骨架的边界和不足

说实话,这套骨架能扛住中小规模的业务——几千到几万在线连接级别。如果到了几十万连接,还有几个问题需要进一步优化:

  1. 心跳任务是公共定时线程池,虽然scheduleWithFixedDelay不会堆积,但每个连接每20秒触发一次任务,GC压力是有的。优化思路是用Netty的EventLoop或者更轻量的环形队列定时器。
  2. 连接Map是单机内存存储,多实例部署时还需要配Redis发布订阅。
  3. SseEmitter的send是同步的,如果消费者消费速度跟不上生产速度,CallerRunsPolicy会用业务线程兜底,此时要警惕任务队列积压变长。

先知道这套方案的边界,你才知道它什么时候该退役。这是架构层面很重要的一件事。

7. 面试时如何把SSE方案讲得“有架构感”

前面这堆实战内容,面试的时候就是你的弹药。但很多同学问题在于:有详细的细节,却没有清晰的表达框架,讲着讲着就被面试官带偏了。

我自己当面试官时,一般会这么考察候选人:

第一步,先问:SSE和WebSocket的区别是什么? 第二步,再问:如果让你给一万个用户做订单状态推送,怎么设计? 第三步,追问:连接断开了怎么办?服务端怎么感知断开?推送失败了要不要补偿?

如果你的回答能形成下面这条链路,面试官基本会满意:

“我用SSE主要是因为它走HTTP,天然穿透各种代理和网关,前端接入成本低。服务端拿到SseEmitter后,把它封装成一个有状态的连接对象,按userId维度维护连接Map。推送时通过独立的线程池发送,避免影响业务线程。

连接保活方面,我设置了比容器idleTimeout更小的心跳间隔,定时往连接里写注释行。同时记录每个连接的lastHeartbeatTime,由一个清理任务把死连接移除。

断线恢复方面,客户端通过EventSource自动重连,重连时带上Last-Event-ID,服务端根据事件ID从事件表里做补偿推送,保证消息不丢。

如果重连失败次数达到阈值,前端会切换成短轮询降级,保证用户最终能看到数据。多实例部署时用Redis发布订阅广播推送,每台机器只推送本机持有的连接。

整体设计就是:连接有状态、心跳有兜底、推送有补偿、流量有降级。”

这一段讲下来,架构感、业务感、细节感就都到位了。面试官要是再往深挖“心跳定时器怎么管理”“Redis广播的重复消费怎么处理”“线程池参数怎么定”,那就看你今天看的这篇文章到底消化了多少。

8. 我在生产环境踩过的三个真实坑,希望你绕开

最后分享三个真实的踩坑经历。这三个坑不在框架代码里,而是在“你以为你懂了,但其实没有”的细节里。

坑一:超时时间设置成0,以为是“永不超时”,结果反而被坑。最早我把SseEmitter的timeout设成0L,想着永远不超时最省事。结果某次流量波动时,连接大量堆积,GC停顿导致所有线程都卡住,连接全被底层强制回收。后来我改成显式设置60秒,配上20秒心跳,世界就清净了。0确实是不超时,但底层容器的资源管理器不认可你的“不超时”,它会按自己的节奏办事。主动设成你想要的超时时间,让行为可预测,才是生产环境该有的姿态。

坑二:心跳发送用scheduleAtFixedRate,在GC停顿后出现“心跳风暴”。scheduleAtFixedRate是严格按照时间基准来调度的。如果某次GC停顿了2秒,原本该在T+20秒执行的任务,会在GC结束后立刻执行;如果期间累积了好几次没执行,它会连追好几个周期,一下子发出两三条心跳。这本身不算灾难,但如果同一时间上万连接都在追心跳,就是一次微型的CPU尖峰。改成scheduleWithFixedDelay后,心跳始终是上一次执行完再延迟固定时间,不会追赶补发,行为温和很多。

坑三:Nginx代理没关缓冲,SSE变成“假实时”。第一次上SSE时压测,发现前端总是延迟5~10秒才收到消息,还以为是网络问题。排查到最后发现是Nginx的proxy_buffering默认开启,把SSE的数据缓冲起来了,缓冲到一定量才一次性转发。就这一行配置,直接改变了SSE的实时性。把proxy_buffering off打开后,数据一条发一条,前端秒收。这种问题,你在本地直连服务端是永远测不出来的,因为本地没有Nginx这层。所以生产环境一定要从用户的链路去排查,不能只看服务端通不通。

这三个坑放到一起,可以看出一个共性:SSE上生产的难点不在于SSE协议本身,而在于它跑在一个“正常HTTP请求”模型并不适用的环境里。连接、缓冲、超时、并发——这些隐形的约束才是决定成败的地方。你把这些约束挨个摸透了,SSE就像一把顺手的手术刀,指哪打哪;摸不透,它就是一台随时会崩的跑步机。

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

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

立即咨询