拓冰建站拓冰建站
首页 / 资讯中心 / 正文

SSE生产级实战:断线重连、心跳保活与超时降级全解析

先说个结论SSE这玩意儿看起来简单到不行不就是返回一个text/event-stream流吗但绝大多数的Java项目SSE一上生产就出问题。要么连接断了一直不复位要么消息推着推着就断了要么用户在线数越积越多最后内存被拖爆。我见过太多团队把SSE当成“高级轮询”来用却不知道SSE最核心的价值——协议级的断线重连和事件续传——完全被浪费了。这篇文章我把SSE在生产环境最要命的三件事讲透断线重连怎么做、超时降级怎么做、以及一套可以直接抄的完整方案。同时把面试中关于SSE的追问点全都覆盖到让你面试讲得清楚生产扛得住。1. 为什么你的SSE一上生产就断1.1 先搞清楚SSE到底是什么SSEServer-Sent Events是基于HTTP的服务器单向推送协议。客户端通过EventSource对象发起一个普通的GET请求服务端在响应头里声明Content-Type: text/event-stream然后这个连接就变成了一条“长连接”服务端可以随时往这条连接里写入事件数据。客户端收到数据后通过事件监听器处理。很多人对比SSE和WebSocket时会说“SSE就是单向的WebSocket是双向的”。这个说法没错但不够准确。我更愿意把SSE理解成“基于HTTP的、自带断线重连机制的消息推送协议”。它最大的杀手锏不是推送本身而是协议层内置了重连、事件ID、重连间隔控制这些能力而这些能力WebSocket全都没有都得自己造轮子。维度SSEWebSocket方向服务端到客户端单向双向实时通信底层协议HTTPtext/event-streamWSTCP帧协议断线重连协议内置自动重连需要自己实现事件ID/续传支持Last-Event-ID不支持二进制消息不支持支持实现复杂度低一个GET连接高需要处理握手、心跳、帧、关闭对于“服务端给单个用户推通知、推状态、推变更”这类场景SSE的复杂度和稳定性是明显优于WebSocket的。这也是为什么我要把SSE讲透因为它在绝大多数业务系统里就是比上WebSocket性价比高得多的方案。1.2 四个把SSE搞挂的隐藏杀手先说最经典的“最小实现”GetMapping(value /sse, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter sse() { return new SseEmitter(); }这段代码看着没毛病能跑通浏览器也能收到消息。但上了生产环境它最多只能活几十秒。杀手一默认异步超时。Spring Boot对异步请求有默认超时时间spring.mvc.async.request-timeout默认是30秒。你用new SseEmitter()不指定超时等于告诉框架走默认配置30秒内没有写入数据连接自动超时断开。也就是说你每隔3分钟推一条消息的业务第30秒就被断开了。更坑的是断开之后浏览器EventSource会自动重连但重连完30秒又断。你感觉系统好像在工作但消息实际上一直在“重连-断开-重连”的循环里打转用户那边什么都收不到。杀手二没有心跳保活。即使你设置了超时时间甚至设置为不超时只要连接上长时间没有TCP数据流动Nginx、SLB、云厂商的负载均衡等中间设备就会把这条空闲连接拆掉。常见的报错就是热词里那个idle timeout waiting for sse或者stream disconnected before completion: idle timeout waiting for sse翻译过来就是连接因为空闲超时被断了但协议流还没结束。杀手三Nginx代理默认开启缓冲。proxy_buffering默认是开启的。开启缓冲意味着Nginx会把后端传过来的数据攒在缓冲区里等攒够了或者等连接结束才一次性发给客户端。SSE是流式响应如果开着缓冲前几条事件根本不会实时到达浏览器直到缓冲区满或者连接被断开。这就会造成非常诡异的“消息延迟”“消息一次性涌出一堆”的现象。杀手四不管理连接生命周期。很多人写SSE只往一个静态Map里放了SseEmitter然后只做send不注册onCompletion、onTimeout、onError回调。连接断了、超时了、异常了map里的SseEmitter对象还留着用户下次重连又放一个新的进去。这个map越积越大最终内存溢出。这就是SSE最常见的“连接泄漏”。这四个杀手任何一个都能让你的SSE上不了生产更别说同时踩中四个了。下面我从断线重连和超时降级两个维度把完整的解法讲清楚。2. 断线重连让SSE连接“断了也能续上”2.1 先吃透协议自带的重连能力SSE协议里服务端可以向客户端发送三类信息事件数据、事件ID、重连时间。先看一段标准的SSE响应id: 1001 event: message data: {a: 1} retry: 10000这段响应拆开来看id表示这条事件的唯一ID客户端收到后会自动记住如果连接断开再重连时会在请求头里带上Last-Event-ID: 1001。event表示事件类型浏览器端可以注册对应的监听器比如addEventListener(message, ...)。data是事件内容如果有多行data浏览器会用换行拼接。retry告诉浏览器如果连接断开等多少毫秒后再发起重连。也就是说断线重连不是让你自己造轮子协议已经给你铺好了路。关键是服务端要不要配合。这里最容易被忽略的是id字段。服务端如果每发一条事件都带一个全局递增或者至少趋势递增的ID客户端就会自动记住“我收到哪一条了”重连时自动把Last-Event-ID带给服务端。服务端看到这个值就知道该从哪儿开始补发这就是“断线续传”的底层机制。2.2 服务端断线补发基于Last-Event-ID接上一段服务端要做的不光是“收到 Last-Event-ID 然后打个日志”而是要根据这个ID把客户端缺失的事件重新推送一遍。实现上需要一个“事件流水”存储。最简单的方案是内存环形队列每条事件一个自增ID保留最近N条。生产环境建议用Redis List或Redis Stream原因后面多实例部署那节会讲。伪代码先放在这里GetMapping(value /subscribe, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe(RequestParam(userId) String userId, RequestHeader(value Last-Event-ID, required false) String lastEventId) { SseEmitter emitter new SseEmitter(0L); // 注册回调清理连接... emitters.put(userId, emitter); // 有 Last-Event-ID从历史流水里补发 if (StringUtils.hasText(lastEventId)) { ListStoredEvent missed eventHistory.pollEventsAfter(lastEventId); for (StoredEvent event : missed) { emitter.send(SseEmitter.event() .id(event.getId()) .name(event.getName()) .data(event.getPayload(), MediaType.APPLICATION_JSON)); } } return emitter; }这段代码的核心思路是建立新连接时先查一下这个客户端上次收到的事件ID把之后的所有事件补给它。这样即使客户端掉线了5分钟重新连上来也能把中间漏掉的关键消息捞回来。需要注意pollEventsAfter这个方法的返回结果一定要按ID升序发送否则客户端收到的事件顺序会乱。另外如果历史事件里包含你已经不关心的过期事件可以直接跳过只补发对业务有意义的变更事件。2.3 前端重连策略指数退避 最大重试次数EventSource默认自带自动重连但如果服务端挂了、网络闪断了单纯靠它默认的重连可能在服务端还未恢复的几十秒里拼命打请求。所以前端一定要做自己的重连控制。我说的“自己控制”不是关掉EventSource的自动重连而是主动替换它的行为监听到onerror时主动close()然后自己调度重连。这样你可以控制重连间隔、最大重试次数以及重试耗尽后的降级动作。核心思路是指数退避class ReconnectableEventSource { constructor(url, options {}) { this.url url; this.options options; this.maxRetries options.maxRetries || 5; this.retryCount 0; this.eventSource null; this.connect(); } connect() { if (this.eventSource) { this.eventSource.close(); } this.eventSource new EventSource(this.url); this.eventSource.onopen () { this.retryCount 0; // 连接成功重置重试次数 }; this.eventSource.onmessage (e) { this.handleEvent(e.data); }; this.eventSource.onerror (e) { this.eventSource.close(); if (this.retryCount this.maxRetries) { this.onMaxRetriesExceeded(); return; } const delay Math.min(1000 * Math.pow(2, this.retryCount), 30000); this.retryCount; setTimeout(() this.connect(), delay); }; } handleEvent(data) { // 处理事件业务自己实现 } onMaxRetriesExceeded() { // 重试耗尽执行降级 } close() { this.eventSource this.eventSource.close(); } }几点实测心得1000 * Math.pow(2, retryCount)算出来是 1秒、2秒、4秒、8秒、16秒……封顶30秒这是比较合理的退避节奏。onopen里把retryCount清零很重要。如果连接恢复后又断开重试从头开始算否则会出现“连接明明恢复了但重试次数已经很大”的尴尬。不要只依赖服务端retry字段来自动重连因为服务端的retry只能控制重连间隔控制不了“重试多少次后降级”。3. 超时降级从连接到业务的全链路兜底3.1 三类超时你一定要认清它们我在排查SSE生产问题时遇到过最多的是下面三类超时。很多人把它们混为一谈盲调参数越调越乱。超时类型含义典型表现常见原因异步连接超时服务端单个连接的总生命周期固定时间点连接断开比如正好30秒SseEmitter未设置超时时间走了Spring Boot默认值空闲超时连接上长时间没有数据传输无固定规律中间设备随机断开没有心跳Nginx的proxy_read_timeout到了读写超时服务端发送数据失败发送时报IOException客户端断开、网络抖动、数据量过大注意第一类和第二类经常叠加出现。你设置了足够长的异步连接超时但没做心跳那么空闲超时会先把你干掉。你做了心跳但把超时时间设成了0L不超时那么某个客户端网络异常后连接永远不会被服务端主动回收泄漏风险反而更高。所以我推荐的组合策略是总超时时间设得比心跳间隔大几倍比如60秒超时、30秒心跳。心跳能保活就保活保不了就让超时机制把连接回收掉靠前端重连拉新连接。3.2 心跳保活让空闲连接持续“呼吸”心跳的本质是在连接空闲时主动发送一条客户端可以忽略的数据让TCP层产生数据传输从而重置各个中间设备的空闲计数器。最简单的实现是发送SSE注释行emitter.send(SseEmitter.event().comment(heartbeat));注释行在SSE协议里的格式是:heartbeat\n\n浏览器收到后会直接忽略不会触发任何事件监听器。但它确实会让TCP连接上产生数据包这就够了。Spring的实现里SseEmitter.event().comment(heartbeat)最终生成的字节确实包含换行实测对各个网关的保活都有效。心跳用ScheduledExecutorService来做比较可控ScheduledExecutorService heartbeatExecutor Executors.newSingleThreadScheduledExecutor(r - { Thread t new Thread(r, sse-heartbeat); t.setDaemon(true); return t; }); heartbeatExecutor.scheduleAtFixedRate(this::sendHeartbeat, 0, 30, TimeUnit.SECONDS);注意一个细节scheduleAtFixedRate的任务里一定要捕获所有异常。如果某个用户的心跳发送失败且你没捕获异常整个调度线程会退出之后所有人的心跳都会停掉这是一个线上事故级的坑。我一开始就踩过排查了半天发现是某一次send抛出的IOException把心跳线程弄死了。正确的发送逻辑private void sendHeartbeat() { emitters.forEach((userId, emitter) - { try { if (emitter ! null) { emitter.send(SseEmitter.event().comment(heartbeat)); } } catch (Exception e) { // 心跳失败说明连接已经断了清理掉 log.warn(SSE heartbeat failed, remove userId{}, userId); emitters.remove(userId); emitter.completeWithError(e); } }); }3.3 消息推送失败与慢消费者降级心跳是保活但业务消息发送时也会失败。最常见的是客户端已经断开但连接还没被服务端感知到此时调用emitter.send会抛出IOException。我的兜底策略很简单发送失败后立即移除本地连接。把这条消息写入待补发历史防止客户端重连后丢失。不再继续重试因为重试只是在加快连接清理节奏。public boolean sendToUser(String userId, String eventName, String payload) { SseEmitter emitter emitters.get(userId); if (emitter null) { return false; } try { String eventId snowflakeId(); emitter.send(SseEmitter.event() .id(eventId) .name(eventName) .data(payload, MediaType.APPLICATION_JSON)); eventHistory.add(new StoredEvent(eventId, userId, eventName, payload)); return true; } catch (Exception e) { log.warn(SSE send failed, userId{}, userId, e); emitters.remove(userId); try { emitter.completeWithError(e); } catch (Exception ignore) { } return false; } }如果业务场景是“同一个用户快速产生大量状态变更”比如订单状态、任务进度那么还有一个降级策略每个用户在内存里只保留“最新状态”而非所有历史事件推送时先判断是否有比当前更早的事件在队列里未发送有则直接丢弃只发送最新状态。这就是“合并推送”或者说“最终状态覆盖”。它能在高并发状态下有效防止慢消费者拖垮服务端。3.4 最坏情况SSE降级为轮询前端重试次数用尽仍然连不上SSE怎么办不能直接放弃业务还需要消息可达性。此时降级为HTTP轮询这是“最坏情况下的兜底”。前端判断条件很简单onMaxRetriesExceeded触发后启动一个定时器每隔几秒调用普通HTTP接口拉取增量数据。后端轮询接口可以复用SSE断线续传的历史事件流水GetMapping(/history) public ListStoredEvent history(RequestParam(userId) String userId, RequestParam(value lastId, required false) String lastId) { return eventHistory.pollEventsAfter(userId, lastId); }这样设计的好处是SSE通道恢复后前端可以立刻切回SSE并且由于最后收到的事件ID被记着切换过程不会丢消息。我实际项目里会把SSE降级轮询的门槛放大一点连续重试5次失败才降级而且要加一个“静默期”降级后每5秒轮询一次每次轮询到数据后至少再轮询3次确认SSE确实恢复不了才彻底停留在轮询模式。4. 面试/生产双杀一套可以直接抄的完整方案4.1 服务端完整实现连接管理 心跳 续传 干净回收下面是一个Spring Boot SseEmitter的完整服务端方案包含我前面所有提到的要点。这套代码我直接在生产项目里改过多次整理出的版本相对干净可以直接参考。Component public class SseConnectionManager { private static final Logger log LoggerFactory.getLogger(SseConnectionManager.class); // userId - SseEmitter private final ConcurrentMapString, SseEmitter emitters new ConcurrentHashMap(); // 事件历史流水生产环境请替换为Redis List/Stream private final DequeStoredEvent eventHistory new ConcurrentLinkedDeque(); private static final int HISTORY_SIZE 1000; private final ScheduledExecutorService heartbeatExecutor Executors.newSingleThreadScheduledExecutor(r - { Thread t new Thread(r, sse-heartbeat); t.setDaemon(true); return t; }); PostConstruct public void init() { heartbeatExecutor.scheduleAtFixedRate(this::sendHeartbeat, 0, 30, TimeUnit.SECONDS); } public SseEmitter subscribe(String userId, String lastEventId) { // 同一用户重复连接时先清理旧连接 SseEmitter old emitters.remove(userId); if (old ! null) { old.complete(); } // 0L表示不主动超时由心跳物理连接状态控制生命周期 SseEmitter emitter new SseEmitter(0L); emitter.onCompletion(() - { log.info(SSE completed, userId{}, userId); emitters.remove(userId, emitter); }); emitter.onTimeout(() - { log.warn(SSE timeout, userId{}, userId); emitter.complete(); }); emitter.onError(e - { log.error(SSE error, userId{}, userId, e); emitters.remove(userId, emitter); }); emitters.put(userId, emitter); // 从Last-Event-ID补发 if (StringUtils.hasText(lastEventId)) { replayMissedEvents(userId, emitter, lastEventId); } return emitter; } public boolean sendToUser(String userId, String eventName, Object payload) { SseEmitter emitter emitters.get(userId); if (emitter null) { return false; } try { String eventId UUID.randomUUID().toString(); emitter.send(SseEmitter.event() .id(eventId) .name(eventName) .data(payload, MediaType.APPLICATION_JSON)); recordEvent(new StoredEvent(eventId, userId, eventName, payload)); return true; } catch (Exception e) { log.warn(SSE send failed, userId{}, userId, e); emitters.remove(userId); try { emitter.completeWithError(e); } catch (Exception ignore) { } return false; } } public int onlineCount() { return emitters.size(); } private void sendHeartbeat() { emitters.forEach((userId, emitter) - { try { emitter.send(SseEmitter.event().comment(heartbeat)); } catch (Exception e) { log.warn(SSE heartbeat failed, remove userId{}, userId); emitters.remove(userId); try { emitter.completeWithError(e); } catch (Exception ignore) { } } }); } private void replayMissedEvents(String userId, SseEmitter emitter, String lastEventId) { for (StoredEvent event : eventHistory) { if (event.userId.equals(userId) event.id.compareTo(lastEventId) 0) { try { emitter.send(SseEmitter.event() .id(event.id) .name(event.name) .data(event.payload, MediaType.APPLICATION_JSON)); } catch (Exception e) { break; } } } } private void recordEvent(StoredEvent event) { eventHistory.addLast(event); while (eventHistory.size() HISTORY_SIZE) { eventHistory.removeFirst(); } } }这里有几个关键设计面试时是加分项emitters.remove(userId, emitter)用了ConcurrentMap的two-arg版本防止误删新连接。比如旧连接的回调晚触发把用户刚建立的新连接给删了。onTimeout里调用的是emitter.complete()而不是直接remove这样会让前端收到连接关闭信号从而触发EventSource的自动重连。事件ID用UUID还是雪花ID都行但必须保证在“同一个用户的消息序列里”比较单调这样Last-Event-ID的compareTo才有意义。如果用的是入库的自增ID直接存long值更好。4.2 前端完整实现EventSource 手动重连 轮询降级前端跟服务端配套要同时覆盖“鉴权、重连、降级”三个需求。这里有一个我特别想强调的点EventSource的API不支持自定义Header。如果你想用Authorization: Bearer xxx这种标准方式鉴权EventSource做不到。可选方案有两个URL参数携带token比如/subscribe?userId123tokenxxx。简单缺点是token会出现在访问日志里所以token要短时效或者只用于建立连接。用fetch ReadableStream代替EventSource可以自定义Header解析SSE格式。如果你们公司对安全要求中等URL参数完全可以接受很多大厂在SSE上也这么干的配合短期token 重连时刷新token即可。如果鉴权要求比较高用fetch方案。下面是fetch方案的完整实现class SSEStreamClient { constructor({ url, token, maxRetries 5, pollInterval 5000 }) { this.url url; this.token token; this.maxRetries maxRetries; this.pollInterval pollInterval; this.retryCount 0; this.pollTimer null; this.closed false; this.lastEventId null; this.usePolling false; } connect() { if (this.closed) return; this.connectWithFetch(); } async connectWithFetch() { try { const resp await fetch(this.url, { headers: { Authorization: Bearer ${this.token}, Accept: text/event-stream, }, }); if (!resp.ok || !resp.body) { throw new Error(HTTP ${resp.status}); } this.retryCount 0; this.usePolling false; const reader resp.body.getReader(); const decoder new TextDecoder(utf-8); let buffer ; while (true) { const { value, done } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const events buffer.split(/\r?\n\r?\n/); buffer events.pop(); for (const raw of events) { this.parseEvent(raw); } } // 正常结束后主动重连 this.scheduleReconnect(); } catch (e) { this.scheduleReconnect(); } } parseEvent(raw) { const event {}; const lines raw.split(/\r?\n/); for (const line of lines) { if (line.startsWith(:)) continue; // 注释行忽略 const idx line.indexOf(:); if (idx -1) continue; const field line.slice(0, idx); let value line.slice(idx 1); if (value.startsWith( )) value value.slice(1); event[field] value; } if (event.id) this.lastEventId event.id; if (event.data) { try { this.onEvent(event.event || message, JSON.parse(event.data)); } catch { this.onEvent(event.event || message, event.data); } } } scheduleReconnect() { if (this.retryCount this.maxRetries) { this.startPolling(); return; } const delay Math.min(1000 * Math.pow(2, this.retryCount), 30000); this.retryCount; setTimeout(() this.connect(), delay); } startPolling() { if (this.usePolling) return; this.usePolling true; this.pollTimer setInterval(() this.poll(), this.pollInterval); } async poll() { try { const params new URLSearchParams({ userId: this.userId }); if (this.lastEventId) { params.set(lastId, this.lastEventId); } const resp await fetch(/history?${params.toString()}, { headers: { Authorization: Bearer ${this.token} }, }); const list await resp.json(); list.forEach(item this.onEvent(item.name, item.payload)); this.lastEventId list.length ? list[list.length - 1].id : this.lastEventId; } catch (e) { // 轮询失败暂不处理等下一轮 } } onEvent(type, data) { // 业务处理由调用方覆盖 } close() { this.closed true; if (this.pollTimer) { clearInterval(this.pollTimer); } } }几个实现细节我补充一下buffer decoder.decode(value, { stream: true })这里必须用{ stream: true }否则多字节UTF-8字符在跨chunk时会被切割成乱码。按空行分割事件时events.pop()把不完整的一块留在buffer里等下一批数据来了再拼上去。这是处理所有流式协议的标准做法。parseEvent里对字段名和值都做了trim处理否则会遇到data: {a: 1}和data:{a:1}两种格式导致解析不一致的问题。4.3 多实例部署SSE的“分布式”难题单体应用里用一个ConcurrentHashMap存emitter就够了但生产环境肯定不是单实例。一旦你上了Nginx负载均衡用户连接和业务产生事件的请求可能落在不同的实例上。用户连接在实例A业务事件却发到了实例B此时B的本地Map里找不到这个用户的emitter消息就丢了。所以SSE的多实例核心思路是本地存连接广播发事件。本地Map负责维护本实例上的客户端连接事件发送走消息广播让所有实例都尝试往本地连接上推。最简单的广播方案就是Redis Pub/SubComponent public class SseEventPublisher { private final StringRedisTemplate redisTemplate; private final SseConnectionManager connectionManager; public SseEventPublisher(StringRedisTemplate redisTemplate, SseConnectionManager connectionManager) { this.redisTemplate redisTemplate; this.connectionManager connectionManager; } public void publish(String userId, String eventName, Object payload) { // 同时发到Redis其他实例收到广播后再推送 redisTemplate.convertAndSend(sse:notify, JSON.toJSONString(new PushEvent(userId, eventName, payload))); // 本地直接推送 connectionManager.sendToUser(userId, eventName, payload); } // 每个实例都订阅同一个频道 Bean MessageListener messageListener() { return (message, pattern) - { PushEvent event JSON.parseObject(message.getBody(), PushEvent.class); connectionManager.sendToUser(event.getUserId(), event.getEventName(), event.getPayload()); }; } }这种方案的要点本地直接推一份再通过Redis广播一份。如果用户就连接在当前实例本地推送立刻到达如果不是则本地的sendToUser返回false因为本地map里没有这个用户Redis广播后用户实际连接的实例会把消息推过去。Redis Pub/Sub的缺点是广播消息不持久化如果某个实例刚好重启它上面存的连接会丢失客户端会重连到其他实例。重连后如果有Last-Event-ID就可以从历史流水补发。所以历史事件流水一定不要只放在本地内存要么放Redis List、要么放数据库不然重启等于丢全部历史。会话粘滞Session Stickiness对SSE也是一个可选项。如果你不想引入广播机制可以配置负载均衡策略保证同一个用户ID的请求都打到同一个实例上。但HTTP连接一旦在网关层被重新分配比如实例重启、扩缩容粘滞策略就会失效。所以更稳妥的还是Redis广播 事件流水存储。4.4 Nginx与网关配置清单SSE上生产Nginx配置是绕不开的一环。下面是经过验证的最小配置location /sse/ { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Connection ; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; proxy_set_header Last-Event-ID $http_last_event_id; # 关键关闭代理缓冲保证事件实时推送 proxy_buffering off; proxy_cache off; chunked_transfer_encoding on; # 超时时间建议大于心跳间隔 proxy_read_timeout 3600s; proxy_send_timeout 3600s; # 如果有gzip建议在此location关闭 gzip off; }proxy_set_header Last-Event-ID $http_last_event_id;这行很多人不知道。Nginx默认不会把客户端的Last-Event-ID头转给后端。不配置这行后端通过RequestHeader(Last-Event-ID)永远拿不到值断线续传就失效了。另外proxy_read_timeout不是设置一次就一劳永逸。如果网关前面还有一层云负载均衡它的空闲超时通常默认是几秒到几十秒这个也要配长一些。配到3600秒只是兜底真正保证连接长期的还是心跳。5. 常见问题与排查实录5.1 “idle timeout waiting for sse”到底是怎么回事热词里反复出现before completion: idle timeout waiting for sse和stream disconnected before completion: idle timeout waiting for sse。这个报错通常来自HTTP客户端库比如WebClient、Netty、OkHttp等发生在客户端已经发起了SSE请求但在设定时间内没有收到任何数据。排查步骤先确认服务端是否真的有数据在流动。如果业务就是几分钟才一条消息那十有八九是服务端没发心跳导致的。加心跳立刻能解决。再看是不是中间网络设备的问题。在服务端上用tcpdump抓包观察连接上是否有周期性TCP数据包。有但客户端还是报idle timeout检查客户端自己的读超时配置把客户端的读超时时间设置得比心跳间隔大一些。最后看是否存在“假连接”。客户端建立的连接确实存在但服务端这个连接没有进入正常的SSE处理流程。比如本地连接管理里压根没这条记录那自然不会有任何数据发过来。最直接有效的验证方式用curl测一段视频流或直接测SSE接口观察是否能持续收到空行心跳。curl -N -H Accept: text/event-stream http://localhost:8080/sse/subscribe?userId1如果curl能不断收到:heartbeat注释行基本可以断定服务端心跳是通的接下来排查客户端和中间链路。5.2 连接只增不减泄漏排查SSE连接泄漏是生产环境最隐蔽的问题。现象是用户量只涨不跌内存也跟着涨GC频率越来越高。排查思路先看连接管理器的Map大小。我习惯给SSE做一个/sse/online监控接口实时返回当前在线连接数配合监控系统画曲线。如果用户下线后曲线不回落说明有泄漏。顺着泄漏去查回调。重点看onCompletion、onTimeout、onError里是否真的把emitter从Map里删掉了。特别是onCompletion里只打了日志没做清理这最坑。查心跳线程。如果心跳发送失败只是打了error日志但没把坏连接删掉下一次心跳会再次失败连接就永远留在Map里。我这边有个教训当时为了快速上线把emitter.onCompletion的回调写成了emitters.remove(userId);而不是emitters.remove(userId, emitter)。结果用户快速重连时新连接刚放进去旧连接的completion回调把新连接也删了。用户一直收不到消息还疯狂重连最后把服务端连接数打爆了。后来改成带value比较的remove(userId, emitter)这个诡异问题才消失。5.3 消息推不过去CORS、鉴权、格式问题有些场景下SSE连接建得起来但消息就是推不过去。CORS浏览器跨域请求时EventSource跨域也是受CORS限制的。服务端要在响应头里加上Access-Control-Allow-Origin而且跨域场景下EventSource默认不带cookie如果你用Cookie鉴权还得设置withCredentials true。fetch方案同理。鉴权EventSource不支持自定义Header要么URL参数要么Cookie。URL参数要注意日志脱敏。Cookie方案要注意CSRF防护因为SSE的GET请求是可以被第三方页面发起的。格式SSE协议要求事件与事件之间用空行分隔单条事件内部各行用\n结尾不能是\r\n吗实际上浏览器和Spring都兼容\n和\r\n但你如果自己拼字符串拼错了比如缺少空行客户端会一直等下一个空行才触发事件表现为“数据收不到”。如果你用Spring的SseEmitter.event().data(...)构建事件格式一般是没问题的。怕的是有人自己拼String.format(data:%s\n\n, json)来发送这种情况最容易漏掉空行或者漏掉id:字段导致解析错乱。最后再分享一点个人体会这套方案我在生产环境迭代过很长时间被各种超时和断线问题折磨过。如果让我总结一个最核心的教训那就是SSE的关键不在推送而在连接生命周期的管理。你推送业务消息的代码可能只占20%剩下80%都在处理“连接什么时候断、断了怎么补、补不了怎么降级”。另外提一句如果只是想给前端发个通知别一上来就上SSE。先评估消息量和实时性要求。实时性要求不高普通轮询够用实时性高但是单向推送SSE比WebSocket省事得多如果要做双向交互再去考虑WebSocket。工具选型这种事永远是越简单越稳。这套方案里的代码核心逻辑你直接抄下来再根据业务改改事件类型和历史流水存储方式基本就能扛住生产的考验。祝你的SSE上线顺利。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门