ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SSE连接泄漏不用怕:SseEmitter检测与回收方案全解析

SSE连接泄漏不用怕:SseEmitter检测与回收方案全解析 做实时推送这一两年SSE算是被越来越多团队接受了。原因很简单走HTTP协议不用像WebSocket那样另外搭握手和鉴权前端直接用EventSource就能接服务端在SpringBoot里开一个SseEmitter就能推。尤其是AI大模型流式回答火起来之后SSE几乎成了流式输出的默认方案。但SSE用起来爽坑也一样多。我印象最深的一次线上事故是连接数监控一直在涨明明没有那么多用户在页面Tomcat的线程占用却居高不下。结果排查下来是服务端每个主动推送的SseEmitter都没有在客户端断开时正确回收连接被一直挂着像一堆没人收的根。那之后我把SseEmitter的整个生命周期翻了个底朝天也踩过“stream disconnected before completion: idle timeout waiting for sse”这类报错。这篇文章就把我沉淀下来的检测和回收方案完整写出来希望能让还在SSE坑里的同学少走几步弯路。1. 先从一次“连接泄漏”说起SSE的三种典型死法1.1 不是只有网络断开才算断开客户端断开这个词听起来很简单但在真实网络环境里至少有三种完全不同的情况。第一种是客户端主动关闭连接。用户刷新页面、关闭浏览器标签页、从大屏页面离开或者前端调用了eventSource.close()。这类断开TCP层会正常发出FIN包服务端如果正在等待写入数据其实有机会感知到。第二种是网络链路闪断。最常见的是Wi-Fi切换、手机从4G切到5G、公司网络突然不通。这一类断开不是主动的TCP层没有标准的FIN对服务端来说这条连接可能看起来还“好好的”。第三种是客户端进程直接被杀。浏览器崩溃、电脑休眠、App被系统回收同样没有机会通知服务端。这三种情况在服务端SseEmitter眼里表现完全不同。第一种相对好处理后面两种几乎只能靠“超时加心跳”之类的方式去猜测。如果代码里压根没管这件事就会出现最典型的症状服务端SseEmitter对象还在Map里定时任务还在往这条已经没有接收方的连接上发数据连接数只增不减。1.2 SseEmitter的生命周期到底由谁控制写SseEmitter的第一印象就是它像个“推送句柄”Controller里new一个return给Spring MVC之后业务线程用它send数据。但很多新手没意识到SseEmitter并不是一个普普通通的数据对象它的本质是Servlet异步请求机制在Spring MVC里的封装。具体来说当Controller返回SseEmitter后DispatcherServlet不会立刻把响应提交回去而是把这个请求标记为异步原线程直接释放。后续你通过SseEmitter.send()写入数据时数据才会实际写到对应的HTTP响应流里。也就是说SseEmitter的存活周期不再受Controller方法控制它关联的是底层Servlet容器的异步请求上下文。这就引出一个关键点这个异步请求什么时候销毁正常结束时需要代码主动调用complete()异常时触发error回调超时时触发timeout回调。但所有这些回调前提都是“容器感知到了状态变化”。而客户端断开后容器能不能感知取决于你是否在写入时发生了异常。如果你只是把SseEmitter挂在Map里既不发数据也不主动close那它就是一条永远悬挂的异步连接谁也发现不了。2. 为什么onCompletion和onError靠不住2.1 回调不触发的真实原因官方文档和网上教程都会告诉你SseEmitter支持三个回调emitter.onCompletion(() - log.info(连接正常结束)); emitter.onTimeout(() - log.info(连接超时)); emitter.onError(ex - log.error(连接异常, ex));看起来挺全的但你实测下来就会发现客户端直接刷新页面时onCompletion大概率不会打印。为什么因为onCompletion本质上是异步请求结束时的通知而“请求结束”这件事容器只会在自己主动处理时才会触发。客户端断开到服务端发现断开中间隔着一条TCP连接的状态检测光靠回调覆盖不了所有场景。客户端断开场景服务端直观表现onCompletion能否感知主动刷新/关闭页面可能有IOException也可能毫无感知不保证网络闪断数据堆积后续send可能成功也可能失败基本感知不到进程被杀长时间无数据交互彻底静默完全感知不到最实际的说法是这三个回调只在服务端自己触发结束/超时/异常时有意义。比如你主动调complete()onCompletion会触发到了配置的超时时间onTimeout会触发send的时候抛了异常onError会触发。但如果你什么都不做客户端断开对服务端来说就像消失了一样回调根本没有机会执行。2.2 断开后才能发现的异常发送时IOException那客户端断开服务端是不是就完全没感知也不是。当你尝试往一个已经断开的连接上写数据时底层TCP会通过写失败的方式暴露出来。在Spring里SseEmitter.send()如果写不进去通常会抛IOException这就是最直接、最可靠的断开信号。我自己用下来最顺手的写法就是给每个客户端绑定一个SseEmitter实例然后在所有发送消息的地方统一做try-catch捕获IOException之后立刻清理。public void sendMessage(String clientId, String message) { SseEmitter emitter emitterMap.get(clientId); if (emitter null) { return; } try { emitter.send(SseEmitter.event().name(message).data(message)); } catch (IOException e) { log.warn(向客户端 {} 推送失败连接已断开执行回收, clientId); emitterMap.remove(clientId); emitter.complete(); } }注意清除顺序很重要先remove再complete。如果不先remove而是等complete的回调里再remove万一客户端是假死状态complete可能也会卡住Map就会一直保留这条记录。2.3 被忽略的容器异步超时idle timeout还有一个藏在暗处的角色就是Servlet容器的异步请求超时时间。SpringBoot默认内嵌的Tomcat对异步请求有一个async timeout默认大概在30秒左右不同版本有差异。如果SseEmitter在这个时间内没有任何数据写入容器会认为是空闲超时直接结束异步请求同时日志里可能就会打印stream disconnected before completion: idle timeout waiting for sse这条日志本身不是致命的但它说明一个非常重要的事SSE连接必须要有“心跳”或者需要把超时时间拉长到业务允许的范围。否则你只是建了个SseEmitter一分钟没推数据它就会被容器静默回收前端那边就会收到连接关闭然后自动重连服务端也看不到明确原因。3. 正确回收SseEmitter的三板斧3.1 发送失败就是最直接的信号上一部分已经提到捕获IOException是识别断开的第一个信号。这里有三个实践上的细节值得单独拿出来说。所有对外发送消息的地方必须走同一个出口不要散落在各个业务代码里。否则漏一个try-catch就多一个泄漏点。我见过不少项目一开始接入SSE时只在一个Controller里send后来功能多了到处直接调emitter.send结果某个角落里漏了异常处理线上连接数就开始悄悄上涨。不只是业务消息要try-catch心跳发送也要。因为假死连接最容易被心跳发不出去这件事暴露。捕获到IOException后的动作不应该仅仅是打印日志而是完整的清理从Map移除、调用complete()、释放关联资源。如果你用定时任务向一批SseEmitter广播消息那业务的try-catch更要包在for循环外面。我曾经见过一段代码for循环里其中一个emitter.send抛异常后整个方法中断后面的客户端全都没收到消息。这就把“单点断连”扩大成了“全部推送失败”。3.2 心跳机制让“假死”连接现形相信很多人在接入SSE的第二天就学会了发心跳。因为如果长时间不推业务数据容器超时或前端断连的问题马上就会出现。心跳的本质就是让连接保持“活跃”同时验证连接是否真的还通。推荐心跳间隔设置在15到20秒之间理由有两点大于容器默认的异步超时时间但小于自己设置的SseEmitter超时时间。前端如果配置了浏览器级别的超时比如经过了代理或网关心跳间隔不能太接近留足余量。心跳数据不需要额外构造复杂的JSON发一个简单事件即可。前端EventSource不会因为命名事件触发onmessage所以对业务无感知。Spring可以这样发emitter.send(SseEmitter.event().name(heartbeat).data(ping));如果客户端是假死状态TCP链路已经断了心跳发送时同样会抛IOException那正好借这个机会把这根连接清掉。所以心跳定时任务必须和清理逻辑联动不要只发不问结果。3.3 定时扫描兜底配合Map做连接清理心跳能解决大部分问题但还不够。原因很简单如果你设置了SseEmitter超时时间为0表示永不超时或者业务休息时间大于心跳间隔但小于你的清理预期假死连接仍然可能残留一段时间。更稳妥的是在Map基础上维护“最后活跃时间”再启动一个定时扫描任务把超过一定时间没有活跃过的连接统一清掉。先简单设计一个结构用SseEmitter的子类来做public class ClientSseEmitter extends SseEmitter { private volatile long lastActiveTime; public ClientSseEmitter(long timeout) { super(timeout); this.lastActiveTime System.currentTimeMillis(); } public void updateLastActiveTime() { this.lastActiveTime System.currentTimeMillis(); } public long getLastActiveTime() { return lastActiveTime; } }然后在所有send成功后调用updateLastActiveTime()。定时扫描逻辑就很简单Scheduled(fixedDelay 30_000) public void sweepDisconnectedClients() { long now System.currentTimeMillis(); emitterMap.entrySet().removeIf(entry - { ClientSseEmitter emitter entry.getValue(); if (now - emitter.getLastActiveTime() 90_000) { log.warn(客户端 {} 超过90秒未活跃执行强制回收, entry.getKey()); emitter.complete(); return true; } return false; }); }这里90秒的阈值要结合你的业务心跳间隔来定。一般心跳15秒一次90秒没活跃基本可以断定这个客户端已经不在了。这个兜底策略的好处是即便心跳和send都没发现异常最终也能把僵尸连接清理掉。3.4 一个容易踩坑的细节超时时间到底设多少SseEmitter构造方法接收一个超时时间单位毫秒。比如new SseEmitter(60_000L)表示60秒没有任何数据写入就触发onTimeout。如果你想要一个“只有主动complete才会关闭”的长期连接一般会写成new SseEmitter(0L)0表示永不超时。但这里有个大坑即使SseEmitter设置了0底层Tomcat的async timeout依然可能生效。也就是说你创建SseEmitter时说好了永不超时但容器层还有一道超时限制。这种情况在SpringBoot 2.2之后的版本出现过表现为日志里突然出现“idle timeout waiting for sse”前端也会无规律地重新连接。处理办法有两种修改SpringBoot的配置项把容器异步请求超时调大或设为-1。老老实实发心跳。我个人推荐先用心跳解决因为修改容器级别的全局配置影响面太大其他异步请求也会被波及。另外如果使用了Nginx或负载均衡做反向代理也别忘了调proxy_read_timeout。这个我放在经验篇里详细说。4. 完整可落地的SpringBoot SSE推送示例4.1 服务端带心跳和断线检测的SseEmitter把前面的思路拼起来一个完整的服务端实现大概是这样的。先定义连接管理类Component public class SseConnectionManager { private final MapString, ClientSseEmitter emitterMap new ConcurrentHashMap(); public SseEmitter subscribe(String clientId) { ClientSseEmitter emitter new ClientSseEmitter(0L); emitterMap.put(clientId, emitter); emitter.onCompletion(() - log.info(连接完成clientId{}, clientId)); emitter.onTimeout(() - { log.warn(连接超时clientId{}, clientId); emitterMap.remove(clientId); }); emitter.onError(ex - { log.warn(连接异常clientId{}, error{}, clientId, ex.getMessage()); emitterMap.remove(clientId); }); try { emitter.send(SseEmitter.event().name(connect).data(connected)); } catch (IOException e) { emitterMap.remove(clientId); emitter.complete(); } return emitter; } public void sendMessage(String clientId, String message) { ClientSseEmitter emitter emitterMap.get(clientId); if (emitter null) { return; } try { emitter.send(SseEmitter.event().name(message).data(message)); emitter.updateLastActiveTime(); } catch (IOException e) { log.warn(推送消息失败回收连接clientId{}, clientId); emitterMap.remove(clientId); emitter.complete(); } } public int activeCount() { return emitterMap.size(); } }然后Controller暴露订阅接口RestController RequestMapping(/sse) public class SseController { private final SseConnectionManager connectionManager; public SseController(SseConnectionManager connectionManager) { this.connectionManager connectionManager; } GetMapping(/subscribe) public SseEmitter subscribe(RequestParam String clientId) { return connectionManager.subscribe(clientId); } GetMapping(/active) public int activeCount() { return connectionManager.activeCount(); } }最后是心跳定时任务Component public class HeartbeatScheduler { private final SseConnectionManager connectionManager; public HeartbeatScheduler(SseConnectionManager connectionManager) { this.connectionManager connectionManager; } Scheduled(fixedRate 15_000) public void sendHeartbeat() { connectionManager.broadcastHeartbeat(); } }我再在SseConnectionManager里补一个广播心跳的方法public void broadcastHeartbeat() { for (Map.EntryString, ClientSseEmitter entry : emitterMap.entrySet()) { String clientId entry.getKey(); ClientSseEmitter emitter entry.getValue(); try { emitter.send(SseEmitter.event().name(heartbeat).data(ping)); emitter.updateLastActiveTime(); } catch (IOException e) { log.warn(心跳发送失败回收连接clientId{}, clientId); emitterMap.remove(clientId); emitter.complete(); } } }这样一个带订阅、带推送、带心跳、带断线回收的最小闭环就成型了。需要注意Scheduled默认是单线程执行的如果连接数特别多心跳发送可能成为瓶颈建议单独给定时任务配一个线程池避免所有心跳排队等待。4.2 前端EventSource的正确使用姿势前端接入其实非常轻量原生浏览器API即可const es new EventSource(/sse/subscribe?clientId${clientId}); es.addEventListener(connect, function (event) { console.log(连接建立成功, event.data); }); es.addEventListener(message, function (event) { const data JSON.parse(event.data); renderData(data); }); es.onerror function () { // EventSource会在连接断开后自动重连这里只做UI提示 showConnectionLost(); };这里有几个细节值得说。EventSource默认会自动重连而且断线后重连时会带上Last-Event-ID头服务端可以根据这个参数做断点续推。如果业务有严格顺序要求可以考虑支持一下。onerror并不等于连接彻底挂了网络抖动也可能触发。不要在这里直接做销毁否则会干扰自动重连。如果后端返回的是SSE协议标准格式EventSource就能解析。注意自定义事件要用addEventListener监听而不是onmessage。我在代码里用了自定义事件名connect、message、heartbeat所以前端都通过addEventListener来挂处理函数。4.3 curl验证断开后服务端是否真的回收了没有前端页面的时候怎么验证服务端的回收逻辑curl就是最好的工具。先启动服务然后执行curl -N --max-time 15 http://localhost:8080/sse/subscribe?clientIdtest1-N参数会取消curl的缓冲让返回的数据即时打印到终端。--max-time 15表示最多请求15秒后强制断开正好模拟客户端主动断开。执行过程中你会在终端看到connect事件数据输出。再开一个终端查看服务端活跃连接数curl http://localhost:8080/sse/activecurl的--max-time到期后客户端会主动关闭连接。此时服务端相关日志如果打印了“心跳发送失败”或者业务推送时的“推送消息失败”就说明断开被成功识别并回收了。活跃连接数也会降下来。我实测下来最常见的现象是curl挂在那里不动15秒没有数据服务端连接数却一直不减。原因多半是SseEmitter超时时间设成了0心跳也没有扫描兜底也没有那客户端断开后服务端当然无法感知。这也是为什么我强调三条回收手段要组合着用而不是只靠某一个。5. 典型问题速查与版本坑5.1 日志里的“stream disconnected before completion”怎么解读这句报错我见过很多人问第一次见时我也懵stream disconnected before completion: idle timeout waiting for sse它并非SpringBoot应用抛出的业务异常而是Tomcat容器在异步请求空闲超时后打的日志。翻译成人话就是这个SSE连接建立了但服务端一定时间内没写任何数据进去容器判定它空闲于是主动结束了连接。现象可能原因快速处理日志出现idle timeout waiting for sse容器异步超时长时间无数据写入加心跳或调大async timeout前端不定时自动重连Nginx层proxy_read_timeout过短调大代理超时并关闭缓冲连接数持续上涨SseEmitter未被回收结合IOException清理与定时扫描send时偶发IllegalStateException连接已关闭仍调用send捕获异常并做幂等清理出现这句话优先排查三件事SseEmitter创建后是不是超过容器async timeout没有发送任何数据如果是要么设置一个合理的心跳要么把超时时间用SseEmitter构造方法拉长。使用的SpringBoot版本是否已经将默认异步超时时间收紧SpringBoot 2.2之前和之后的默认表现确实有差别。反向代理层的proxy_read_timeout是否设置得太短代理层也可能主动掐断长连接。处理办法也不复杂但不要一上来就改全局async timeout。最安全的还是加心跳第4章的示例可以直接拿来用。5.2 连接数只增不减从哪开始排查如果监控发现服务端连接数一直上涨我一般按这个顺序排查。第一步先确认是不是业务正常增长。连一下/active接口看看当前连接数是否符合预期。第二步检查SseEmitter是否都被存进了Map且从未移除。重点看代码里有没有像emitterMap.remove这类清理逻辑以及是否只有complete回调里才清理。第三步检查客户端是否有默认重连导致连接重复叠加。如果服务端生成clientId的逻辑有缺陷同一个用户可能每次重连生成不同的IdMap里就攒了一堆旧连接。第四步用jstack抓一下线程栈看看是否有线程阻塞在SseEmitter的send调用或者等待队列上。这个可以快速定位是不是因为某个连接不消费导致写缓冲被填满从而拖住整个线程池。我遇到过最隐蔽的情况是客户端虽然断开了但它之前的TCP连接处于半开状态服务端每次send都能成功因为数据只进了内核写缓冲还没真正触发错误。这种情况只有等到写缓冲被填满或者通过TCP keepalive才能发现。所以最后还是要靠“最后活跃时间定时扫描”来兜底。5.3 SpringBoot 2.x与3.x在SSE上的差异如果你的项目还在用SpringBoot 2.x或者刚升级到SpringBoot 3.xSSE这块有几个差异值得留意。SpringBoot 2.x默认走javax.servletSseEmitter在spring-web模块里整体用起来没太大问题。SpringBoot 3.x换成了jakarta.servlet包名变化对业务代码无感但如果你直接操作了Servlet API比如自己拿AsyncContext加Listener就需要注意包名的改动。另外SpringBoot 3.x升级后一些老版本的Tomcat配置在application.yml里写法可能变了。比如原来设置异步请求超时的配置项新的版本建议直接看当前版本的配置文档。我踩过一次3.x下配置项写错导致所有SseEmitter连接都在30秒左右被断开。还有一个实践感受SpringBoot 3.x对连接关闭的检测比2.x更及时同样的“客户端主动刷新页面”场景3.x的onError回调触发概率更高一些。但不要因此就以为不需要try-catch和心跳网络假死场景依然要靠心跳和扫描。6. 最后分享几个我实际项目里的习惯文章写到这儿核心方案已经全部给出了最后聊几个我自己在项目里长期坚持的习惯。第一所有SseEmitter相关操作集中在一个Manager类里管理Controller只负责订阅业务模块只负责调sendMessage。这样断线清理、计数监控、日志打印都集中在一个出口排查问题的时候不用到处找代码。第二对外暴露一个活跃连接数接口不用上监控系统也能快速确认状态。技术上其实就是Map.size()但用起来真香。给运维同学排查问题时一个curl就能看到结果。第三心跳内容不要塞业务数据。我见过有人用心跳字段顺便传业务状态结果心跳和业务消息的时序混在一起前端解析出各种奇怪问题。心跳就是心跳保持简单也方便前端主动忽略。第四Nginx代理场景下记得把proxy_buffering关掉否则SSE数据会被Nginx缓冲成一坨一坨地吐出来体验非常差。location /sse/ { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_read_timeout 3600s; proxy_set_header Connection ; proxy_http_version 1.1; }这条配置能解决很大一部分“明明后端推送了前端却迟迟收不到”的问题。这个经验是我在用户反馈“大屏数据刷新有延迟”之后才发现的。第五如果你在代码里用了线程池来发送消息一定要小心线程池队列堆积的问题。SSE连接多且推送频率高时一个队列里可能积压大量发送任务而这些任务里又夹杂着SseEmitter引用一旦线程池出问题内存泄漏会比连接泄漏更严重。后来我把这套回收逻辑抽成了一个公共模块新项目接入SSE都是现成的连接数再也没出过问题。说白了SSE本身是个轻量方案但轻量不代表没有资源管理问题。把回收当成连接的核心逻辑来设计而不是事后打补丁你的服务才能真正扛住真实流量。
返回列表