ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

SpringBoot SseEmitter连接管理:断线感知与资源回收实战

SpringBoot SseEmitter连接管理:断线感知与资源回收实战 做后端这些年我有个很深的体会越是看着简单的技术坑起来越要命。SSEServer-Sent Events在 SpringBoot 里就是一个典型。表面上看它不过就是一个SseEmitter对象往里面 send 数据前端用 EventSource 接收完事。可真要放到生产环境里跑几天客户端一断、网络一抖、代理一超时服务端那堆没人回收的 SseEmitter 能把线程池拖垮日志里刷满stream disconnected before completion: idle timeout waiting for sse你才会反应过来这玩意儿远没有 Demo 里那么乖巧。这篇文章不打算从头讲 SSE 协议是什么我默认你已经知道它能做服务端单向推送。我重点想聊的是当客户端异常断开时服务端 SseEmitter 到底该怎么感知、怎么回收、怎么避免资源泄漏。里面会涉及 SpringBoot 的异步请求机制、Tomcat 的异步超时、线程池使用、心跳保活这些实操细节也会把我在真实项目里踩过的坑一条条拉出来说。无论你是刚接触 SSE 消息推送还是已经写了不少 SSE 接口但总觉得心里没底这篇都值得你花十分钟从头到尾看完。1. 先说清楚SSE 到底解决了什么问题1.1 从轮询到单向长连接在 SSE 出现之前服务端要给前端主动推数据最常见的方式就是轮询。前端每隔几秒发一次请求问服务端有消息了吗。这种方式实现简单但很浪费大部分请求其实都没有新数据HTTP 请求头和响应头的开销白白消耗着服务端还要承受高并发查询压力。SSE 的思路完全不一样。客户端发一次请求服务端把响应挂起保持这条 HTTP 连接不关闭然后源源不断地往响应流里写数据直到服务端主动结束或连接断开。因为连接一直是复用状态不存在反复握手的开销数据到达前端几乎是实时的。这本质上是一个单向长连接。很多初学者第一次看到 EventSource 会把它和 WebSocket 搞混面试里也经常被问。这里我用一句话区分WebSocket 是服务端和客户端双向平等对话SSE 是客户端听服务端单方面广播。如果你的需求里服务端只是单向推送SSE 是比 WebSocket 更轻量的选择。1.2 为什么选 SSE 而不是 WebSocket我在技术选型时有个习惯优先选标准 HTTP 协议能解决的方案而不是一上来就上重量级的东西。WebSocket 很强但它是独立的协议需要额外的握手、心跳、断线重连机制后端还要考虑协议兼容和网关配置。SSE 的底层还是 HTTP浏览器原生支持 EventSource 接口自带自动重连和 Last-Event-ID 恢复能力这些对消息推送场景来说太友好了。更实际的一点是SSE 在微服务和网关架构下更容易混过去。它不需要特殊的协议解析Nginx、SpringCloud Gateway 等组件对 HTTP 的支持天然就是完整的最多调一下缓冲参数。而 WebSocket 在网关层经常要做额外的升级和转发规则。我在做 AI 大模型流式输出对接时体会尤其深。大模型生成 text 是按 token 逐个吐出来的服务端给前端逐字推送SSE 简直是量身定做。前端拿到text/event-stream一个 onmessage 就能逐段渲染。如果换成 WebSocket虽然也能做但重连、消息序号管理、服务端主动关闭这些逻辑全得自己来开发成本完全不在一个量级。1.3 SSE 能用到哪些落地场景以我接触过的项目来看SSE 最适合的几类场景非常明确AI 对话流式输出大模型生成结果实时推给前端逐字展示。服务端消息广播比如订单状态变化、工单提醒、系统公告推送客户端订阅后实时接收。日志实时展示运维平台、任务调度平台把执行日志实时刷到前端页面。进度条推送批量任务、大数据导出服务端把进度百分比不断推给页面。行情数据推送股票、币价、赛事比分等高频变化但又是单向前推的数据。这些场景共同点是数据产生在服务端客户端只需要被动接收而且对实时性有要求。如果你遇到的是这种需求SSE 是一个性价比很高的技术方案。2. SseEmitter 是怎么工作的2.1 关键类不是只有 SseEmitter 一个很多人以为 SseEmitter 就是全部实际上它只是 Spring MVC 异步处理体系里的一个工具类。这套体系的核心是 Servlet 3.1 的异步请求和几个配套类SseEmitter继承自ResponseBodyEmitter专门负责按 SSE 协议格式写数据。ResponseBodyEmitter管理异步响应输出流负责把数据写到客户端连接。DeferredResult异步结果容器比 ResponseBodyEmitter 更底层一些。AsyncContextServlet 容器提供的异步上下文连接的生命周期由它维护。具体流程是这样的Controller 方法被调用后Spring MVC 发现返回值是 SseEmitter就会把当前 HttpServletRequest 标记为异步请求拿到容器创建的 AsyncContext然后立即把 Controller 线程释放回线程池。随后连接并没有关闭真正的数据推送给一个或多个业务线程去完成。搞懂这一层你就会明白一个关键点Controller 方法 return 之后连接还活着。如果你在 return 之后没有保存好 SseEmitter 引用后面想往这条连接推数据就推不了了。所以 SseEmitter 通常需要放进一个全局的会话管理器来维护。2.2 一个最小可用的 SseEmitter 接口先写一个最简版本让你对整体结构有个感觉。这是最基础的 SSE 订阅接口RestController RequestMapping(/sse) public class SseController { private final MapString, SseEmitter emitterMap new ConcurrentHashMap(); GetMapping(value /subscribe, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe(RequestParam String clientId) { SseEmitter emitter new SseEmitter(60_000L); emitterMap.put(clientId, emitter); return emitter; } PostMapping(/publish) public void publish(RequestParam String clientId, RequestBody String message) throws Exception { SseEmitter emitter emitterMap.get(clientId); if (emitter ! null) { emitter.send(SseEmitter.event().name(message).data(message)); } } }第一眼看上去挺清爽吧一个接口订阅一个接口推送好像已经能用了。但我负责任地说这段代码要是直接上生产用不了多久你就会收到报警。因为它只处理了正常工作这条路径完全没有处理客户端断开这种情况。前端一刷新、一关页面、一断网emitterMap 里的 SseEmitter 就成了孤儿连接永远躺在 Map 里占着连接和资源不放。2.3 请求的完整生命周期捋一下一条 SSE 连接从生到死会经过哪些节点创建客户端发 GET 请求服务端 new 一个 SseEmitter 并返回。注册回调框架把 SseEmitter 和当前异步上下文绑定Controller 线程释放。数据推送业务线程通过emitter.send()往输出流写 SSE 数据。超时如果设置了超时时间到达时间后容器触发 async timeout。异常连接断开、写入失败、容器报错时触发 error 回调。完成连接关闭后触发 completion 回调。这六个节点里第 4、5、6 是资源回收的关键也是最容易出错的地方。很多人对onCompletion()和onError()到底谁先触发、要不要在里面调用complete()一头雾水。我们逐个拆开讲。3. 问题核心客户端断开了服务端却蒙在鼓里3.1 报错stream disconnected before completion从哪来你在日志或前端控制台里可能会看到类似stream disconnected before completion: idle timeout waiting for sse的报错第一次遇到大概率一脸懵。这句话其实分成几段看stream disconnected before completion连接在正常完成之前被断开了。idle timeout waiting for sse连接空闲超时一直在等 SSE 数据但什么都没等到。很多情况下这其实是服务端或中间件主动断开连接的结果。客户端并没有主动关闭页面但客户端和服务端之间长时间没有数据流动触发了链路某一层的 idle timeout。这个链路可能很长客户端浏览器、Nginx/网关、Tomcat、业务代码任何一层都可能有空闲超时设置。在 SpringBoot 内嵌 Tomcat 环境下常见的超时配置有这几个配置项作用默认值参考server.tomcat.connection-timeout连接建立后等待请求到达的超时时间通常 20s~60sserver.tomcat.keep-alive-timeoutkeep-alive 连接空闲多久后关闭默认同 connection-timeoutserver.servlet.async.timeout异步请求处理最长持续时间取决于容器Tomcat 默认 30sSseEmitter 构造器传入的 timeout单个 emitter 的异步超时默认继承容器配置问题来了SSE 连接建立之后如果业务上一直没数据推或者心跳做得很随意这条连接就处于空闲状态。一旦空闲时间超过上述某一个超时值容器就把连接关了。前端 EventSource 收到异常后会自动重连如果重连后又进入空闲、又超时断开就会出现每隔一段时间断一次的诡异现象。3.2 给 SseEmitter 设置合理超时既然空闲超时会断开连接那我把超时设得无限大问题不就解决了吗理论上确实如此比如new SseEmitter(0L)或一个极大值在部分容器上可以表示不超时。但我不建议这么做原因有二一来无限超时意味着服务端必须依赖客户端断开事件来回收资源而客户端断开 TCP 连接后服务端其实很难第一时间感知。如果所有连接都无限超时断线的连接会像慢性病一样慢慢堆积最终耗尽线程池和内存。二来生产环境的上层负载均衡设备、云厂商网关不会听你程序里设置它们有自己的空闲超时策略。你程序里设置不超时网关层还是会在几分钟后掐断连接。我的建议是SseEmitter 的超时要设成一个业务可接受的值比如 5 分钟或 10 分钟同时用心跳消息来保活。心跳消息是 SSE 协议里的一种特殊写法服务端定期发送一行以冒号开头的注释数据比如: heartbeat客户端收到后既不会触发 message 事件又能让连接保持活跃。这一招既能骗过网关和容器的空闲超时又能让超时机制兜底回收那些真正失联的连接。3.3 send() 抛 IOException 是唯一可靠的断线感知你要特别注意onError()和onCompletion()都是事后通知它们确实能被触发但触发时机不一定准。我在实际项目里发现客户端直接关掉浏览器页面或者断了网服务端并不会立刻收到通知。很多时候要等到下一次往这条连接写数据、底层流写入失败抛 IOException 的时候服务端才后知后觉。所以可靠的断线回收逻辑永远要围绕send() 的异常处理来设计。每次推送数据时如果send()抛出 IOException基本可以断定这条连接已经完蛋了这时候立刻从会话管理器移除 SseEmitter并调用它的complete()释放底层资源。下面这段代码展示了这个思路public void sendEvent(String clientId, String data) { SseEmitter emitter sessions.get(clientId); if (emitter null) { return; } try { emitter.send(SseEmitter.event().name(message).data(data)); } catch (IOException e) { // 连接已断开立刻回收资源 sessions.remove(clientId); emitter.complete(); } }这里我补充一个要点SseEmitter.event()可以设置事件名、id、重连时间、数据这种结构化写法比直接用emitter.send(data)更规范。如果前端用 EventSource 监听的是默认message事件你也可以不显式指定 event 名称保持简单。4. 正确的回收姿势从回调钩子到兜底策略4.1 回调钩子onCompletion / onTimeout / onError 的准确用法SseEmitter 继承自 ResponseBodyEmitter它提供了三个注册回调的方法我把它们放在一起对比回调方法触发时机正确用法禁止事项onCompletion(Runnable)请求完成时无论正常结束、超时或异常都会触发清理业务资源、移除会话禁止在回调里调用complete()会抛 IllegalStateExceptiononTimeout(Runnable)达到超时时间时触发调用complete()关闭连接、移除会话不能什么都不做否则连接一直挂着onError(ConsumerThrowable)发生异常时触发记录异常、移除会话、必要时complete()不要吞掉异常不处理onCompletion()有个很容易踩的坑很多人以为连接完成后反正都会触发它于是在里面既移除 Map 又调complete()结果发现抛了IllegalStateException。原因是连接已经处于完成态你再调 complete 属于重复关闭。我一般把它当成资源清理通知所有核心清理逻辑放在自己封装的 remove 方法里回调只做触发不做关闭动作。onTimeout()是坑最多的一个。容器触发异步超时后并不会自动帮你关闭连接你必须自己在回调里调用complete()。如果忘了写这条连接就会一直挂着直到更外层的东西来掐断它。写完complete()之后onCompletion()还会再触发一次所以在onCompletion()里做清理时要保证幂等性比如用ConcurrentHashMap.remove天然就是幂等的。4.2 一个可落地的 SseSessionManager把上面这些经验整合起来我写了一个生产环境可用的会话管理类你可以直接参考Component public class SseSessionManager { private final MapString, SseEmitter sessions new ConcurrentHashMap(); public SseEmitter subscribe(String clientId) { SseEmitter emitter new SseEmitter(60_000L); emitter.onCompletion(() - removeQuietly(clientId)); emitter.onTimeout(() - { log.warn(sse timeout, clientId{}, clientId); emitter.complete(); }); emitter.onError(e - { log.warn(sse error, clientId{}, msg{}, clientId, e.getMessage()); emitter.complete(); }); sessions.put(clientId, emitter); return emitter; } public void send(String clientId, String data) { SseEmitter emitter sessions.get(clientId); if (emitter null) { return; } try { emitter.send(SseEmitter.event().data(data)); } catch (IOException e) { log.warn(sse send failed, client disconnected: {}, clientId); removeAndComplete(clientId); } } public void removeAndComplete(String clientId) { SseEmitter emitter sessions.remove(clientId); if (emitter ! null) { emitter.complete(); } } private void removeQuietly(String clientId) { SseEmitter emitter sessions.remove(clientId); // 连接已处于完成状态这里不做 complete防止重复关闭 } }在这个类里核心逻辑就几点用一个 ConcurrentHashMap 保存所有客户端连接key 是客户端 ID。subscribe 时注册好三个回调超时和异常路径都主动complete()。send 时捕获 IOException一旦写入失败立刻移除和关闭连接。回调里只移除不重复 complete规避 IllegalStateException。为什么用 ConcurrentHashMap因为 SSE 连接是异步访问的多个业务线程可能同时往不同客户端推数据还会有回调线程同时修改 Map非线程安全的 HashMap 在这里很快就会出问题。ConcurrentHashMap 的remove是线程安全的能保证不会有两个线程同时把自己当成成功移除者避免重复 complete。4.3 心跳线程让超时机制成为兜底而不是常态前面提到SseEmitter 设了 60 秒超时如果业务 60 秒内不推数据连接就会被超时兜底回收。但这个设计会带来一个新问题正常用户挂机超过 60 秒连接也被回收了。所以我们需要心跳定时来维持活跃。最朴素的心跳方案是启动一个后台定时线程定期遍历所有 SseEmitter向每个连接发送一行注释数据Scheduled(fixedDelay 20_000) public void heartbeat() { long now System.currentTimeMillis(); sessions.forEach((clientId, emitter) - { try { emitter.send(SseEmitter.event().comment(heartbeat)); } catch (Exception e) { removeAndComplete(clientId); } }); }SSE 协议规定以冒号开头的行是注释前端 EventSource 收到后不会触发事件但网络层确实有数据流动这就能有效避开各种 idle timeout。我把心跳间隔设为 20 秒远小于 60 秒超时这样正常情况下连接永远不会被容器超时断开只有真的出问题、心跳也推不进去的时候超时机制才会作为兜底来清理。这里有一个容易被忽略的细节如果用 Spring 的Scheduled需要确保项目开启了EnableScheduling如果不想引 Spring 的定时任务也可以用一个ScheduledExecutorService单线程池来做。我实际更推荐后者因为它不依赖 Spring 的调度环境而且对线程池的控制更直接。4.4 异步线程池与连接器配置SseEmitter 的推送最好放在专用的异步线程池里把耗时任务和 Tomcat 的工作线程隔离开。你一定不要干这种事情在某个业务线程里同步循环调用emitter.send()给几百个客户端推数据那样会把业务线程卡在 IO 上。我一般是这样设计的服务端收到上游推送事件后交给一个带线程池的 Dispatcher 处理每个客户端发送任务提交到ExecutorService线程池大小按连接数估算private final ExecutorService sseExecutor new ThreadPoolExecutor(8, 16, 60L, TimeUnit.SECONDS, new SynchronousQueue(), new ThreadFactoryBuilder().setNameFormat(sse-push-%d).build(), new ThreadPoolExecutor.CallerRunsPolicy());这个线程池有几个细节SynchronousQueue不允许任务堆积如果线程池满了新任务会走拒绝策略。拒绝策略选CallerRunsPolicy意思是让提交任务的线程自己执行好处是任务不会被静默丢弃坏处是提交线程可能被拖慢。对于 SSE 推送这种实时性要求高的场景我宁可让业务线程阻塞一下也不愿意丢失推送事件。另外SpringBoot 的server.servlet.async.timeout这个配置项也要留意。虽然 SseEmitter 构造器传入的超时时间能覆盖它但如果你在有些代码里 new SseEmitter 时忘了传超时就会退回到这个全局配置。建议在application.yml里显式配置一个合理值比如 60 秒server: servlet: async: timeout: 60000这样即使某个 SseEmitter 没显式设超时也不会无限期挂起。5. 生产环境被 SSE 坑过的真实例子5.1 案例一定时推送任务把所有线程耗尽我之前维护过一个日志实时推送系统客户端订阅后服务端会把任务执行日志实时推过去。一开始实现很简单每个任务一个线程日志来了就 send。运行一段时间后线程池报警活跃线程数居高不下。排查之后发现两个问题。第一前端 EventSource 连上来之后长时间不操作连接被 Nginx 空闲超时掐断但服务端因为没写数据感知不到Map 里全是死连接。第二定时推送线程遍历所有连接时明明连接已经断了send 方法在个别版本下并不会立刻抛异常而是卡在那里等底层 socket 超时一次遍历就拖死一个线程。解决办法就是我上面说的方案push 时严格控制超时send()包一层 try-catch定时任务每轮都尝试发送心跳失败就立刻清理Map 中的死连接量降下来之后线程池压力自然恢复了。5.2 案例二Nginx 缓冲导致前端时好时坏还有一次很有意思前端反馈说推送的数据经常一次性来一大坨而不是逐条出现。我一开始以为是服务端 push 逻辑写错了后来排查发现是 Nginx 在做缓冲。Nginx 默认会对代理响应做缓冲把数据攒到一定量再转发给客户端。SSE 这种流式响应最怕被缓冲因为你想让前端实时看到数据结果全被 Nginx 扣住了。解决方式很标准location /sse { proxy_pass http://backend; proxy_buffering off; proxy_cache off; proxy_set_header Connection ; proxy_http_version 1.1; chunked_transfer_encoding off; proxy_read_timeout 3600s; proxy_send_timeout 3600s; }核心是proxy_buffering off和proxy_read_timeout调大。如果你的域名前面还有 CDN、云负载均衡也要注意它们有没有类似的缓冲机制。SSE 链路越长要排查的中间节点就越多。5.3 问题速查表我把日常运维中比较常见的问题整理成了一张表遇到问题可以直接对照查现象根本原因处理方式日志提示 idle timeout waiting for sse连接长时间无数据容器或网关空闲超时加心跳保活调整超时配置onTimeout 回调后连接仍不释放回调里没有调用complete()在 onTimeout 中一定要complete()onCompletion 回调里调 complete 抛异常连接已完成重复关闭回调里只清理不调 complete客户端断开后服务端无感知TCP 断开通知延迟或无数据写入依赖 send 时 IOException 做兜底推送数据被前端一次性接收Nginx 或网关缓冲了响应流关闭代理缓冲调长读超时连接数持续增长不下降死连接没有被及时移除用 Map回调心跳组合清理大流量推送导致线程阻塞send 调用阻塞业务线程使用专用线程池隔离业务前端频繁自动重连连接被中间层断开后 EventSource 重连检查代理超时与服务端心跳5.4 Spring Boot 2.x 与 3.x 的差异还有一个现在绕不开的话题Spring Boot 版本。早期我用 Spring Boot 2.x 时SseEmitter 的用法和现在基本一致但有一个点要注意Servlet API 的包名从javax.servlet换成了jakarta.servlet。如果你从 Spring Boot 2.x 升级到 3.x项目里所有直接引用 Servlet API 的地方都要做包名迁移。另外Spring Boot 3.x 对异步请求的支持也做了调整整体走的是 Spring Framework 6.x 的异步处理体系。测试下来SseEmitter 的 API 层面没有破坏性变化但如果你在 WebFlux 和 WebMvc 之间切换会发现底层的抽象完全不同。WebFlux 里对应的是FluxServerSentEvent不属于本文讨论范围但你要知道它们不是同一个东西。如果是新项目且确定要长期用 Servlet 栈做 SSE用 Spring Boot 3.x 没问题如果是存量项目升级最好先在测试环境把 SSE 长连接场景回归一遍因为异步超时、连接回收这些行为在不同版本间的默认值可能有变化。5.5 用 curl 验证 SSE 连接最后分享一个调试技巧。写完之后怎么验证 SSE 接口是否正常用浏览器打开当然可以试但更适合脚本化验证的是 curl。注意一定要加-N参数告诉 curl 不要缓存输出实时打出来curl -N http://localhost:8080/sse/subscribe?clientIdtest如果接口正常你会看到控制台不断输出 SSE 格式的数据块行与行之间有空行分隔。如果等服务端超时你会看到idle timeout相关的日志或 curl 等了很久后连接被关闭。用这种方式我一般几秒钟就能判断是服务端的问题还是前端的问题。我个人在实际项目里还有个习惯把所有SseEmitter的创建、超时、移除、心跳都打上日志每次连接创建和销毁都记录 clientId配合链路追踪就能很清楚地看到每个连接是从哪个环节被回收的。线上出了诡异问题这些日志就是最直接的证据。SseEmitter 的回收从来不是靠某一个回调单打独斗而要靠超时兜底、IO异常感知、心跳保活、幂等清理这四个机制配合起来才靠谱。把这些都做到位SSE 推送才能真正用得安心。
返回列表