从线上告警到代码修复,我花了两天时间才彻底搞懂SpringBoot里SseEmitter断开连接时的正确回收姿势。这个问题表面上看只是一个"客户端关闭了连接,服务端没收到通知"的小事,实际牵涉到HTTP协议层、Servlet异步容器、Spring的事件回调机制,还有SpringBoot不同版本之间的行为差异。如果你正在用SSE做实时推送,或者准备在生产环境上SseEmitter,这篇文章值得你花十分钟看完。
我当时遇到的问题很典型:部署的推送服务运行几天后,连接数持续上涨,内存占用缓慢爬升,最后触发告警。重启之后恢复正常,但过两天又复发。最初怀疑是代码里有连接没关,后来抓了堆栈、开了JMX监控、翻了Spring和Tomcat源码,才定位到根因——客户端断开时SseEmitter的回收没有我预想得那么"自动"。下面我把完整的排查链路、原理分析和最终方案写出来,希望能帮你少踩几个坑。
1. 线上告警:连接数只增不减背后的SseEmitter泄漏案
1.1 事故现场的异常特征
那是一个很普通的下午,监控平台突然弹出推送服务的活跃连接数告警。我打开Grafana看了一眼,曲线呈现出非常规律的爬升趋势——每增加一台客户端设备,连接数就稳定地增加一条,几乎没有回落。这就排除了流量突增的可能,更像是某种"使用一次就泄漏一次"的资源管理缺陷。
紧接着JVM内存告警也来了,老年代使用率一路攀升,GC频率明显变高。用jstack抓了一次线程快照,发现大量处于WAITING状态的线程堆积在Tomcat的NIO事件循环里,看起来像是在等待某个永远不会来的IO事件。因为服务本身流量不大,CPU没有异常,所以一开始我怀疑是某个业务线程卡住了,但仔细看堆栈,发现真正的问题出在SseEmitter上——每个连接都持有一份响应资源和对应的异步上下文,这些对象没有被释放。
更让我困惑的是,服务端的日志里并没有出现任何异常。代码里明明给每个SseEmitter注册了onCompletion回调,按理说客户端断开时应该触发这个回调,连接应该被及时清理。但日志一片安静,仿佛这些连接根本不知道客户端已经走了。
1.2 初步排查:线程、内存、连接三层证据
我先从操作系统层面查了TCP连接状态,执行netstat命令后看到大量连接处于ESTABLISHED状态,RemoteAddress指向的端口已经不是客户端的真实端口了——这意味着TCP四次挥手已经完成,或者说经过了NAT设备的转换,服务端实际已经无法感知这条连接是否活着。
接着我查看了Spring的事件监听日志,确认SseEmitter的onCompletion、onTimeout、onError三个回调都没有被触发。这就很奇怪了,因为按照官方文档的说明,无论连接是正常完成、出错还是超时,都应该有对应的回调被调用。
为了进一步确认,我写了几个模拟客户端,用不同的方式断开连接:有的直接关闭浏览器标签页,有的用curl请求后立即杀掉进程,有的干脆拔掉网线。结果发现:只有调用response.getWriter().close()这种正常结束的连接能触发onCompletion,其他异常断开方式大概率什么回调都不触发。这时我意识到,SseEmitter的回收机制并没有我原来以为的那么完善,它依赖底层容器对连接断开的感知能力,而这种感知在真实网络环境下非常不可靠。
2. SseEmitter的完整生命周期:协议层、容器层与应用层三层视角
2.1 从HTTP响应到异步释放:SSE连接的底层原理
要理解SseEmitter为什么会在客户端断开时"装死",得先搞清楚它底层是怎么工作的。SSE(Server-Sent Events)本质上是HTTP协议的一种长连接用法——客户端发起一个普通的GET请求,服务端不立即结束响应,而是把响应流保持打开,通过chunked编码持续向客户端推送数据。SseEmitter正是Spring对这个模式的高层封装。
当我们调用new SseEmitter(0L)并把它返回给Spring MVC时,Spring会把当前请求从Tomcat的请求处理线程中解绑,切换到异步模式。这个过程中,Tomcat会为这个异步请求创建一个AsyncContext,底层通过NIO的Selector监听这个连接上的读写事件。数据推送时,服务端通过响应输出流写数据;客户端断开时,理论上Tomcat的NIO通道会触发一个读事件,告知连接已关闭。
这段逻辑本身没有大问题,问题出在"客户端断开"的时机和方式上。TCP协议是流式的,服务端只有在尝试往一个已关闭的连接上写数据时,才会收到RST或者Broken Pipe错误。如果客户端只是静默关闭了socket,而服务端又没有主动写数据,那么这个连接在操作系统层面可能很长时间内都不会有任何异常表现。说白了,服务端没办法用一种绝对可靠的方式"主动感知"对端断开,除非对端主动通知,或者服务端自己主动探测。
2.2 Spring提供的三个回调:onCompletion、onTimeout、onError的真实触发场景
SseEmitter提供了三个回调方法,字面上看覆盖了所有结束场景:
- onCompletion:响应完成的回调,官方说法是"请求已经完成,无论是因为正常结束还是因为出错"。
- onTimeout:异步请求超时的回调。
- onError:发生异常时的回调。
但真实的触发行为比官方说法复杂得多。我翻阅了Spring Framework的源码,这三者的触发都依赖容器层对异步请求状态的通知。具体来说,Tomcat的AsyncContext在请求真正结束时会调用监听器,Spring再把这个事件转发给SseEmitter的回调。如果连接在操作系统层面处于半开状态(比如客户端断电、网线断开、NAT会话超时),Tomcat根本不会收到任何事件,自然也就不会触发任何回调。
另外还需要注意,onCompletion一旦触发,SseEmitter就进入完成状态,之后调用send方法会抛出IllegalStateException。所以很多人的回收逻辑是"在onCompletion里移除连接",这本身是对的,但它依赖的前提是onCompletion必须被触发。如果回调本身没有机会执行,那这段清理代码就是摆设。
2.3 为什么"客户端断开"并不等于"服务端感知"
我再用一个生活化的类比解释这个问题。你把电话听筒拿起来没挂断,然后人离开了房间。理论上对方能听到环境音,但如果房间隔音很好,对方完全察觉不到你已经离开,只会一直"喂喂喂"地等待回应。TCP连接也是这样,只要没有数据流动,连接两端的任何一方都不知道对方是否还活着。
更麻烦的是,中间还可能有NAT网关、负载均衡器和代理服务器。客户端应用层主动断开时,TCP FIN包会一路传回来,服务端大概率能感知。但如果客户端是移动设备,网络切换或者系统休眠,TCP连接可能既不发FIN也不发RST,就这么默默消失。这时候服务端不仅不知道连接断了,还会一直维护这条连接上的AsyncContext、SseEmitter对象、响应缓冲区,直到操作系统最终清理TCP连接(这个时间可能长达几十分钟甚至更久)。
这就是SseEmitter"泄漏"的本质:不是代码里忘了调complete,而是连接断开这个事件根本没有被可靠地传递到应用层。解决思路也就清晰了——服务端要主动建立一套"探测-超时-回收"机制,而不是被动依赖底层事件。
3. 根因定位:SpringBoot版本与容器对断连检测的影响
3.1 SpringBoot版本差异:从内嵌Tomcat到异步超时行为
排查过程中有一个非常关键的发现:同样的代码在SpringBoot 2.3和SpringBoot 2.7上行为完全不同。在低版本上,默认的内嵌Tomcat对异步请求有一个默认超时时间,即使你不主动设置,连接也会在一段时间后因为超时被容器回收,从而触发onTimeout回调。而高版本的SpringBoot(尤其是SpringBoot 3.x基于Spring Framework 6的版本)调整了异步请求超时策略,配合Tomcat 10+的NIO实现,行为更加保守——如果没有配置超时,某些版本下异步请求会被长期保持。
这里就引出了热搜词里那个经典报错:stream disconnected before completion: idle timeout waiting for sse。这个报错信息其实是Tomcat通过外部网络层抛出的,意思是等待SSE连接时发生了空闲超时。它出现在客户端配置了读取超时且超时时间小于服务端两次心跳间隔的情况下——客户端主动放弃了连接,但服务端不一定收到通知。
所以如果你用的是较新的SpringBoot版本,new SseEmitter(0L)这种"永不超时"的设置反而可能是隐患的源头。0L确实会关闭Spring层面的异步超时,但如果你不写这一行,不同版本的默认超时逻辑又不一样。我的建议是:不要依赖版本默认行为,显式配置spring.mvc.async.request-timeout或者在创建SseEmitter时给出一个明确的值。
3.2 客户端正常关闭与异常断网的两种路径
我把断连场景拆成两类,分别说明服务端会经历什么。
第一类是"正常关闭"。客户端调用EventSource的close()方法,或者浏览器标签页正常关闭,TCP连接会走完整的四次挥手流程。FIN包到达服务端后,Tomcat的NIO层会在这个socket上检测到EOF,Spring有机会触发onCompletion回调。但这里还有个小坑:如果服务端此时还有缓冲数据未发送,可能先触发的是传输异常而不是正常完成,顺序不固定。
第二类是"异常断网"。手机突然进电梯、电脑休眠、断网,TCP连接处于半开状态。服务端不会收到任何协议层通知。这种情况下,SseEmitter会一直存活,直到下一次服务端尝试向它写数据。如果业务上没有心跳推送,或者心跳推了但消息量特别小,服务端可能很长时间不会触发SocketException,因为少量的写入操作可能被操作系统缓冲,不会立即暴露连接问题。这类连接,如果没有主动回收机制,就会变成真正的"僵尸连接"。
3.3 关键证据:日志中没有Completion回调
我在代码里加了一行日志,打印onCompletion、onError、onTimeout三个回调的触发情况。实测了超过一百次异常断连,结果让我脊背发凉:onCompletion触发率不足15%,onTimeout触发率取决于SseEmitter是否设置了超时,onError几乎不触发。而服务端日志里唯一出现的,是客户端重连时服务端尝试向旧连接写数据抛出的IOException——但这已经是事后了,而且IOException只触发那一次,旧连接对象仍然没有被清理。
这个现象说明什么?说明之前代码里"在onCompletion里remove连接"的方案只覆盖了理想情况。真实生产环境里,客户端断连的形态千奇百怪,回调机制根本兜不住。我们必须把回收逻辑从"事件驱动"改为"事件驱动+定时巡检兜底",也就是我下一章要写的方案。
4. 正确回收SseEmitter:心跳、超时、注册表三件套
4.1 方案设计:把"事件驱动"降级为"轮询兜底"
经过上面的分析,我最终确定了一套可靠性优先的回收方案:核心思想是在SseEmitter的注册表里维护每个连接的最后活跃时间,通过定时任务周期性扫描,把超过阈值的连接主动complete,同时配合心跳来探测连接的存活状态。这套方案的关键在于三个组件:
- 一个ConcurrentHashMap作为SseEmitter注册表,key是客户端ID,value是封装了SseEmitter和最后活跃时间的Holder对象。
- 一个固定频率的心跳任务,向所有注册的SseEmitter发送ping消息,既能保持连接活跃,又能及时暴露已经断开的连接。
- 一个固定频率的清理任务,扫描注册表,把超过N秒没有活跃的SseEmitter主动complete并移除。
为什么要用ConcurrentHashMap而不是普通HashMap?因为SseEmitter的创建和销毁是并发进行的,客户端连接的建立和断开没有固定顺序,如果用一个非线程安全的Map,在高并发下很容易出现ConcurrentModificationException。用ConcurrentHashMap虽然不能完全避免遍历时的一致性偏差,但至少不会在并发读写时抛异常。
还有一个设计细节:心跳消息必须使用SseEmitter.event().name("heartbeat").data("ping")这种格式,而不能只发送一条普通数据。带事件名的心跳消息在客户端可以单独监听,比如EventSource.addEventListener("heartbeat", ...),这样业务数据和心跳数据互不干扰。
4.2 服务端定时探测与资源回收代码
先看注册表和连接管理器的最简实现:
@Component public class SseConnectionManager { private static final long HEARTBEAT_INTERVAL_MS = 15_000L; private static final long IDLE_TIMEOUT_MS = 60_000L; private final ConcurrentHashMap<String, ClientConnection> clients = new ConcurrentHashMap<>(); private static class ClientConnection { SseEmitter emitter; volatile long lastActiveTime; volatile boolean completed; } public SseEmitter createConnection(String clientId) { ClientConnection conn = new ClientConnection(); conn.emitter = new SseEmitter(0L); conn.lastActiveTime = System.currentTimeMillis(); clients.put(clientId, conn); conn.emitter.onCompletion(() -> { conn.completed = true; clients.remove(clientId); }); conn.emitter.onTimeout(() -> { conn.completed = true; conn.emitter.complete(); clients.remove(clientId); }); conn.emitter.onError(e -> { conn.completed = true; conn.emitter.complete(); clients.remove(clientId); }); return conn.emitter; } @Scheduled(fixedRate = HEARTBEAT_INTERVAL_MS) public void sendHeartbeat() { long now = System.currentTimeMillis(); clients.forEach((clientId, conn) -> { try { conn.emitter.send(SseEmitter.event().name("heartbeat").data("ping")); conn.lastActiveTime = now; } catch (IOException e) { safeComplete(conn); } catch (IllegalStateException e) { safeComplete(conn); } }); } @Scheduled(fixedRate = 30_000L) public void recycleIdleConnections() { long now = System.currentTimeMillis(); clients.forEach((clientId, conn) -> { if (conn.completed) { clients.remove(clientId); return; } if (now - conn.lastActiveTime > IDLE_TIMEOUT_MS) { safeComplete(conn); clients.remove(clientId); } }); } private void safeComplete(ClientConnection conn) { try { conn.completed = true; conn.emitter.complete(); } catch (Exception ignored) { } } }这段代码有几个地方需要仔细解释。第一,new SseEmitter(0L)传0表示禁用Spring层面的异步超时,因为超时由我们的定时任务统一管理。如果你传了一个有限值,比如60秒,那么无论客户端是否活着,60秒后onTimeout都会触发,可能造成正常连接被误杀。两种方式没有绝对的对错,就看你想让哪一层来主导回收。
第二,心跳任务和清理任务都通过@Scheduled注解实现。如果你服务的定时任务比较多,建议给SseConnectionManager单独配置一个线程池,避免心跳发送被其他耗时任务阻塞。我在生产环境用的是Spring的TaskScheduler自定义线程池,corePoolSize设置为2就够用了。
третье,心跳发送时捕获了两种异常:IOException和IllegalStateException。前者是网络层面的写入失败,后者是SseEmitter已经完成后再调用send会抛出的异常。这两种异常都说明这条连接已经不可用,直接走安全清理逻辑。
4.3 业务推送与心跳的协同:避免误判连接死亡
有了心跳和清理任务之后,还需要特别注意一个边界:心跳任务刚把lastActiveTime刷新了,但业务消息推送的频率很低,比如每隔5分钟才推一条业务数据,这5分钟内只有心跳在维持连接活跃。清理任务的阈值如果设置得太小,比如30秒,那么只要心跳和清理任务的执行节奏错开一点,就可能误杀正常连接。
我的建议是:清理阈值至少是心跳间隔的4倍以上。比如心跳每15秒一次,清理阈值设置为60秒及以上。这样即使某一次心跳因为网络抖动没有成功刷新时间,还有连续多次心跳机会来纠正,不会因为单次失败就误杀。
另外,如果你在业务推送方法里也需要更新lastActiveTime,可以在发送业务数据成功后顺手刷新,如下所示:
public void sendBusinessMessage(String clientId, Object data) { ClientConnection conn = clients.get(clientId); if (conn != null && !conn.completed) { conn.emitter.send(SseEmitter.event().name("message").data(data)); conn.lastActiveTime = System.currentTimeMillis(); } }这个细节能让"活跃时间"更贴近真实情况,避免出现"心跳没断但业务数据已经推不进去"的伪健康状态。
5. 验证与压测:确认回收逻辑经得起真实断连
5.1 模拟客户端断开:curl、脚本、拔网线三种场景
回收逻辑写完后,我在本地和测试环境各做了一轮验证。模拟断开的工具很简单,核心要覆盖三类场景:正常关闭、进程被杀、网络断开。
第一类场景用curl模拟:
curl -N http://localhost:8080/sse/stream/1001 & sleep 3 kill %1curl的-N参数表示禁用缓冲,让响应内容即时输出。这里kill掉curl进程后,TCP连接会正常关闭,服务端应该能通过心跳任务在下一个周期感知到。
第二类场景用一段Python脚本模拟:
import requests, time, os resp = requests.get('http://localhost:8080/sse/stream/1002', stream=True) print("connected", resp.status_code) time.sleep(3) os._exit(0) # 不发送FIN包直接退出,模拟崩溃进程直接退出时操作系统会关闭socket,但不会主动通知服务端。这个场景更接近小程序端突然崩溃的情况。
第三类场景是拔网线。我在有网线的机器上跑了个测试客户端,建立连接后直接拔掉网线。这是最严酷的场景,模拟的是设备失联,一般心跳机制就是为此设计的。
5.2 验证结果与监控指标:如何确认回收真的生效
验证过程中我主要盯着三个指标:SseEmitter注册表的大小、Tomcat活动连接数、JVM堆内存变化。
我给SseConnectionManager写了一个Actuator端点,可以直接查看当前注册表里有多少条连接以及各自的存活状态。实测结果也很直观:
- curl正常关闭后,下一次心跳发送时立刻抛出IOException,注册表在心跳周期内(15秒内)就会移除该连接。
- Python进程直接退出后,服务端不会立即可感知,但清理任务会在30秒一次的巡检中发现这个连接已经有超过60秒没有活跃,然后主动complete并移除。实测大约在断连后45-75秒内完成回收。
- 拔网线的场景下,由于发出的心跳迟迟没有异常,只能等清理任务的阈值触发。实测回收时间在60秒左右,满足我的预期。
Tomcat活动连接数的曲线也验证了机制的有效性:之前持续爬升的连接数曲线开始变成锯齿形,每一轮都有一批旧连接被回收,一批新连接被创建。JVM老年代内存的爬升趋势也大大缓解,说明SseEmitter对象能被随GC回收。
5.3 压测的意外发现:心跳频率过高会放大问题
我在压测阶段还发现了一个反直觉的现象:心跳频率设置得越密集,客户端断连时服务端收到IOException的速度越快,但如果在并发连接数非常大的情况下(比如上万条),每15秒一次的心跳就变成了一次对注册表的全量遍历和发送,反而增加了服务端的CPU和内存压力。
我的优化方式是:心跳发送不再一个连接一个连接地串行发送,而是把每个连接的最后活跃时间先检查一遍,如果距离上次发送不足一个周期就跳过;同时在心跳代码里直接捕获IllegalStateException而不再往里写日志,因为异常数量在断连频繁时会爆炸式增长,日志刷屏会拖垮磁盘和日志系统。对确实需要记录的连接异常,我只记录clientId到一个统计计数里,定期汇总上报。
6. 避坑清单扩展:生产环境SseEmitter其他常见问题
6.1 Nginx代理导致的空闲连接假死
如果你在应用前面挂了Nginx做反向代理,SSE默认会遇到一个非常坑的问题:Nginx的proxy_read_timeout默认60秒,60秒内上游服务器没有返回任何数据,Nginx就会主动断开客户端与上游的连接,客户端表现为收到的响应被中断。而且Nginx断开时,它向Spring服务端发送的也是一个正常的关闭操作,服务端可能感知不到这是"代理超时"还是"客户端自己关闭"。
解决方法是给SSE的location单独配置一个足够长的读超时时间,同时关闭proxy_buffering,让数据实时转发:
location /sse/ { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Connection ''; proxy_buffering off; proxy_read_timeout 3600s; proxy_send_timeout 3600s; }proxy_buffering off尤其重要,如果开着,Nginx会等上游攒够缓冲才发给客户端,SSE的实时性就完全失去了。Connection头设为空,是为了避免Nginx默认启用KeepAlive导致连接语义混乱。
6.2 客户端重连时SseEmitter实例管理:同ID新旧连接覆盖问题
实际生产中,客户端的断线重连是家常便饭。客户端断网后,它会自动发起一个新的SSE连接,此时如果你的注册表以clientId为key,那么新连接put进来时,旧连接的对象还残留在Map里,就造成了新老连接互相覆盖的问题。
我的推荐做法是:创建新连接前,先检查同一个clientId是否已有旧连接,如果有,先把旧连接安全complete,再put新连接:
public SseEmitter createConnection(String clientId) { ClientConnection old = clients.get(clientId); if (old != null && !old.completed) { old.emitter.complete(); } ClientConnection conn = new ClientConnection(); // ... 省略创建逻辑 }这样能保证同一时刻一个客户端在服务端只对应一条活的SSE连接,避免重复推送和资源浪费。如果你不想用同步锁,也可以用ConcurrentHashMap的compute方法实现原子替换,我为了可读性在示例里用了简单方式,生产建议加上并发细节。
还有一个隐蔽问题:客户端的EventSource会自动携带Last-Event-ID请求头,但如果你的SSE是JSON数据流,建议在推送数据时显式带上id字段,客户端重连时才能正确断点续传。这个和SseEmitter回收关系不大,但属于SSE落地时容易被忽略的点。
6.3 线程池阻塞与异步任务泄漏
SseEmitter本身不需要占用线程,它属于异步请求,但如果你在推送数据时使用了阻塞式的IO或者同步的业务处理,这些任务会在线程池里排队。如果某个连接长期不回收,线程池里堆积的任务也会随之越来越多。
我在线上就踩过一次坑:SseEmitter的心跳和业务推送共用了同一个线程池,结果业务推送偶尔出现慢请求,把线程池占满了,心跳任务排队等待,导致所有连接的lastActiveTime无法更新,清理任务随即误判所有连接都"超时",数据库里的活跃连接数一瞬间清零,把所有在线客户端全部踢下线。
这个问题非常典型。解决方案很简单:给心跳、清理这类系统任务分配一个独立的、固定大小的线程池,让它和业务推送任务完全隔离。即使业务线程池被占满,心跳线程依然能正常工作,保证系统判断连接状态的机制不被业务拖垮。
另外,在分布式环境下,如果同一个服务有多台实例,客户端通过负载均衡连接到不同的实例,那么每台实例的SseEmitter注册表都是独立的。简单场景下可以让客户端根据clientId做一致性哈希路由,保证同一个客户端的连接始终落在同一台实例上;复杂场景则需要引入Redis发布订阅或者消息队列来广播推送消息。这个扩展比较大,这里先提一笔,等之后有时间可以单独写一篇。
最后再分享一个我踩过泥坑后总结的小技巧:即使你的SSE业务非常简单,也一定要在SseEmitter创建后的第一分钟里主动向客户端发送一条初始消息。这不仅能让客户端快速确认连接建立成功,还能让SseEmitter在底层立即执行一次真实的IO写入,把潜在的连接问题提前暴露出来。如果这条初始消息发送就失败了,说明连接从一开始就是坏的,直接回收即可,不要让它占用注册表名额。这个简单的"首消息探活"加上心跳巡检,基本能覆盖我遇到过的绝大多数SseEmitter回收难题。