ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

SSE生产级实战:高并发、断线重连与超时降级全方案

SSE生产级实战:高并发、断线重连与超时降级全方案 1. 这不是“能跑就行”的SSE是生产环境里真刀真枪的存活方案你写过SSEServer-Sent Events吗Java后端同学大概率都写过——用ResponseBody返回ResponseEntityFluxString或者手撸SseEmitter前端用EventSource一接控制台打印出几条“data: hello”心里一松“成了”。但真正把这套逻辑扔进VCP服务器集群、扛住每秒300连接、持续运行72小时不掉线、凌晨三点告警电话打进来还能稳住的——不到10%。这不是夸张。我去年在金融级实时风控平台做SSE网关重构时翻了全团队17个存量SSE服务的代码14个没做连接保活12个没设超时兜底9个连Content-Type: text/event-stream都拼错了大小写最离谱的是一个电商秒杀推送服务用SseEmitter塞了5000条未消费消息进内存队列OOM直接拖垮整个Pod。标题里说的“90%写不上生产”不是打击信心而是划清一条线能本地curl通 ≠ 能上生产能推消息 ≠ 能扛住断线、重连、超时、背压、鉴权、扩容、灰度。今天这篇不讲协议定义、不列RFC文档、不画抽象架构图。我们只干一件事把一套已在3个高并发业务日均消息量2.4亿稳定运行18个月的SSE生产级方案掰开、揉碎、贴着代码和日志讲清楚。你会看到为什么SseEmitter原生API在生产中是“定时炸弹”——它根本没暴露连接状态机断线重连不是前端EventSource.onclose里location.reload()那么简单而是要精确到毫秒级的指数退避服务端会话续传“超时降级”不是加个Timeout注解就完事而是要在连接空闲、消息阻塞、下游依赖超时三个维度分别设防VCP服务器部器部署时Nginx、Tomcat、Spring Boot三层缓冲区怎么调——调错一个参数stream disconnected before completion: idle timeout waiting for sse错误率飙升300%。如果你正被面试官问“SSE怎么保证可靠性”或者刚收到运维告警“SSE连接数突增500%”又或者正在设计AI Agent的实时反馈通道——这篇就是为你写的。它不教你怎么“学会”而是告诉你怎么“活下来”。2. 为什么原生SseEmitter在生产里注定失败——从源码和线程模型看本质缺陷2.1 SseEmitter的“黑盒”状态机你根本不知道连接死了Spring Framework 5.2 提供的SseEmitter看似简单GetMapping(value /events, produces MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter handleEvents() { SseEmitter emitter new SseEmitter(30_000L); // 30秒超时 emitter.send(SseEmitter.event().name(init).data(connected)); return emitter; }但问题藏在底层。我们看SseEmitter核心字段// org.springframework.web.servlet.mvc.method.annotation.SseEmitter private final AsyncWebRequest webRequest; // 封装了HttpServletResponse private volatile boolean completed false; private final CountDownLatch latch new CountDownLatch(1);关键点来了completed字段只在complete()或超时后置为true但连接断开如用户关浏览器、网络闪断时这个字段永远保持false。为什么因为SseEmitter依赖Servlet容器的AsyncContext回调而Tomcat/Jetty对TCP连接异常中断的检测有延迟通常15~60秒且不触发onTimeout或onError。我实测过用curl -N http://localhost:8080/events发起连接然后kill -9 curl进程Spring日志里零记录completed仍为falseemitter对象在内存里挂着直到30秒超时才释放。更糟的是如果此时你调用emitter.send()会抛出IllegalStateException: The emitter has already been completed——但你根本不知道它“已经完成”因为你没监听到任何事件。提示别信SseEmitter.onCompletion(Runnable)——它只在显式调用complete()或超时后触发对被动断连完全失能。2.2 线程模型陷阱一个连接一个线程不是一个线程池里的“幽灵线程”SseEmitter默认使用WebMvcConfigurer配置的TaskExecutor。很多人直接用Executors.newCachedThreadPool()以为“线程够用就行”。错。CachedThreadPool会无限创建线程而每个SSE连接在Tomcat里占用一个AsyncContext背后绑定一个NioEndpoint的Poller线程。当连接数激增线程数爆炸CPU 100%GC频繁最终OutOfMemoryError: unable to create new native thread。我们线上曾出现过单节点Tomcat配置maxThreads200但SSE连接峰值达1200CachedThreadPool创建了1500线程系统负载飙到40所有HTTP接口响应超时。正确做法是为SSE单独配一个有界线程池并强制关联到连接生命周期。比如这样Bean public TaskExecutor sseTaskExecutor() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(50); // 核心线程数预估并发连接数×0.3 executor.setMaxPoolSize(200); // 绝对上限防雪崩 executor.setQueueCapacity(1000); // 队列容量存待发送消息 executor.setThreadNamePrefix(sse-async-); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); executor.initialize(); return executor; }但光配线程池不够——你得确保SseEmitter的send()操作走这个池子。Spring默认不保证必须手动指定GetMapping(/events) public SseEmitter events() { SseEmitter emitter new SseEmitter(30_000L); // 关键绑定自定义执行器 ((SseEmitterImpl) emitter).setTaskExecutor(sseTaskExecutor()); return emitter; }注意SseEmitterImpl是Spring内部类强转有风险。生产环境建议用ResponseBodyEmitter替代或直接封装StreamingResponseBody——后者可控性更强。2.3 内存泄漏三连击未清理的引用、堆积的消息、失控的监听器SseEmitter的三大内存泄漏源静态Map缓存Emitter为支持重连很多同学把SseEmitter存进ConcurrentHashMapString, SseEmitterkey是用户ID。但连接断开后如果不主动remove()emitter对象及其持有的AsyncContext、HttpServletResponse、OutputStream全留在堆里。消息队列无界堆积SseEmitter.send()内部用LinkedBlockingQueue暂存消息若前端消费慢或断连队列无限增长。监听器未注销SseEmitter.onTimeout()、onError()注册的Runnable如果没在complete()时清理会持有外部类引用导致整个Controller实例无法GC。我们线上一个案例某客服系统用static MapString, SseEmitter存连接高峰期每分钟新增2000连接但断连后仅5%被remove()72小时后老年代占用92%Full GC每3分钟一次。解决方案不是“加内存”而是用弱引用定时清理背压控制// 使用WeakReference避免强引用泄漏 private final MapString, WeakReferenceSseEmitter emitterCache new ConcurrentHashMap(); // 启动定时任务每30秒扫描过期连接 Scheduled(fixedDelay 30_000) public void cleanupStaleEmitters() { emitterCache.entrySet().removeIf(entry - { SseEmitter emitter entry.getValue().get(); if (emitter null || isCompleted(emitter)) { log.warn(Cleanup stale emitter for user: {}, entry.getKey()); return true; } return false; }); } private boolean isCompleted(SseEmitter emitter) { try { // 反射读取completed字段Spring 5.3 Field completedField SseEmitter.class.getDeclaredField(completed); completedField.setAccessible(true); return completedField.getBoolean(emitter); } catch (Exception e) { return false; } }3. 断线重连不是前端retry而是服务端会话续传客户端精准锚点3.1 前端EventSource的“伪重连”陷阱标准EventSource配置const es new EventSource(/api/events?userId123); es.addEventListener(message, e console.log(e.data)); es.addEventListener(error, () console.log(reconnecting...));你以为error事件触发时浏览器会自动重连不完全是。Chrome/Firefox默认3秒后重试但重试URL带Last-Event-ID头服务端需解析Safari重试间隔随机且不带Last-Event-ID更致命的是error事件只在连接彻底断开时触发而TCP半开连接如NAT超时时前端认为“还连着”服务端却收不到心跳消息永久丢失。我们抓包发现某次运营商网络抖动前端EventSource.readyState卡在1OPENING长达8分钟期间服务端已断连但前端毫无感知。3.2 服务端会话续传用Redis实现连接状态与消息快照真正的断线重连核心是服务端记住“谁、在哪个时间点、收到了哪条消息”。我们采用三级状态管理层级存储作用TTL连接会话Redis Hash (sse:session:{userId})记录lastEventId、connectTime、ip、userAgent24h消息快照Redis Stream (sse:stream:{userId})存储最近100条消息按eventId有序永久自动滚动全局索引Redis Sorted Set (sse:index)按scoretimestamp存所有活跃会话用于广播/踢人24h关键代码// 创建会话时写入Redis public void createSession(String userId, String ip, String userAgent) { MapString, String session Map.of( lastEventId, 0, connectTime, String.valueOf(System.currentTimeMillis()), ip, ip, userAgent, userAgent ); redisTemplate.opsForHash().putAll(sse:session: userId, session); redisTemplate.opsForZSet().add(sse:index, userId, System.currentTimeMillis()); } // 发送消息时先写Stream再更新lastEventId public void sendMessage(String userId, String eventId, String data) { // 写入Stream自动按eventId排序 MapRecordString, String, String record StreamRecords.string( Map.of(eventId, eventId, data, data, timestamp, String.valueOf(System.currentTimeMillis())) ).withStreamKey(sse:stream: userId); redisTemplate.stream().add(record); // 更新会话lastEventId redisTemplate.opsForHash().put(sse:session: userId, lastEventId, eventId); } // 客户端重连时根据Last-Event-ID拉取缺失消息 public ListMapString, String getMissedMessages(String userId, String lastEventId) { // 从Stream中读取eventId lastEventId的所有消息 Long startId Long.parseLong(lastEventId); return redisTemplate.opsForStream().read( Consumer.from(group, client), StreamReadOptions.empty().count(100), StreamOffset.from(sse:stream: userId, 0-0) // 实际用XREADGROUP ).stream() .flatMap(list - list.getMessages().stream()) .filter(msg - Long.parseLong(msg.getValue().get(eventId)) startId) .map(msg - msg.getValue()) .collect(Collectors.toList()); }注意Redis Stream的XREADGROUP需提前创建消费者组否则首次读取会失败。我们用redisTemplate.execute()执行Lua脚本原子化创建。3.3 客户端精准锚点用Date头自定义Event-ID双保险前端不能只靠Last-Event-ID因为Safari不支持。我们加一层Date头校验GetMapping(/events) public ResponseEntityStreamingResponseBody events( RequestParam String userId, RequestHeader(value Last-Event-ID, required false) String lastEventId, HttpServletResponse response) { // 设置必要头 response.setContentType(text/event-stream); response.setCharacterEncoding(UTF-8); response.setHeader(Cache-Control, no-cache); response.setHeader(Connection, keep-alive); // 关键写入当前时间戳前端用作重连锚点 long now System.currentTimeMillis(); response.setHeader(X-Server-Time, String.valueOf(now)); StreamingResponseBody body outputStream - { PrintWriter writer new PrintWriter(outputStream, false); // 发送初始化事件含服务端时间戳 writer.write(event: init\n); writer.write(data: {\serverTime\: now }\n\n); writer.flush(); // 主循环拉取消息并推送 while (!Thread.currentThread().isInterrupted()) { ListMapString, String messages getMissedMessages(userId, lastEventId); for (MapString, String msg : messages) { writer.write(id: msg.get(eventId) \n); writer.write(event: message\n); writer.write(data: msg.get(data) \n\n); writer.flush(); lastEventId msg.get(eventId); // 更新lastEventId } Thread.sleep(100); // 避免空轮询 } }; return ResponseEntity.ok().body(body); }前端接收时let lastEventId localStorage.getItem(lastEventId) || 0; let serverTime 0; const es new EventSource(/api/events?userId${userId}lastEventId${lastEventId}); es.addEventListener(init, e { const init JSON.parse(e.data); serverTime init.serverTime; // 同步本地时间与服务端 const offset serverTime - Date.now(); }); es.addEventListener(message, e { lastEventId e.id; // 更新lastEventId localStorage.setItem(lastEventId, lastEventId); console.log(e.data); }); es.addEventListener(error, () { // 重连时带上服务端时间戳避免时钟漂移导致漏消息 const reconnectUrl /api/events?userId${userId}lastEventId${lastEventId}t${Date.now() offset}; es.close(); setTimeout(() { es new EventSource(reconnectUrl); }, 3000); });4. 超时降级三层防御体系——连接空闲、消息阻塞、下游依赖4.1 连接空闲超时Nginx/Tomcat/Spring三层缓冲区协同调优stream disconnected before completion: idle timeout waiting for sse错误90%源于Nginx代理超时。默认配置# nginx.conf location /api/events { proxy_pass http://backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; # 缺少关键配置 }必须加这三行proxy_read_timeout 300; # Nginx读超时单位秒 proxy_send_timeout 300; # Nginx写超时 proxy_buffering off; # 关闭缓冲SSE必须流式传输Tomcat层面server.xml需调Connector port8080 protocolHTTP/1.1 connectionTimeout20000 keepAliveTimeout60000 !-- 保持连接60秒 -- maxKeepAliveRequests100 !-- 单连接最多100请求 -- asyncTimeout300000 !-- 异步超时5分钟 -- /Spring Bootapplication.ymlserver: tomcat: connection-timeout: 20000 async-timeout: 300000 spring: web: resources: cache: period: 0关键proxy_read_timeout必须 ≥asyncTimeout否则Nginx先断连Tomcat还不知道。4.2 消息阻塞超时用Reactor背压超时熔断Flux发送消息时如果前端消费慢Flux会堆积在Queue里。我们用onBackpressureBuffer限容timeout熔断GetMapping(value /events-reactor, produces MediaType.TEXT_EVENT_STREAM_VALUE) public FluxServerSentEventString eventsReactor(RequestParam String userId) { return Flux.generate( () - new State(userId, 0), // 初始状态 (state, sink) - { ListMessage messages getMessageFromRedis(state.userId, state.lastId); if (messages.isEmpty()) { sink.next(ServerSentEvent.builder() .event(ping) .data(keepalive) .build()); state.lastId 0; } else { for (Message msg : messages) { sink.next(ServerSentEvent.builder() .id(msg.getId()) .event(message) .data(msg.getContent()) .build()); state.lastId msg.getId(); } } }, state - state // 状态传递 ) .onBackpressureBuffer(100, BufferOverflowStrategy.DROP_LATEST) // 最多缓100条 .timeout(Duration.ofSeconds(30), Flux.just( // 超时发心跳 ServerSentEvent.builder() .event(timeout) .data(connection idle) .build() )) .doOnNext(event - log.debug(Send event: {}, event.getEvent())) .doOnError(error - log.error(SSE error, error)); }4.3 下游依赖超时Hystrix/Fallback 降级消息兜底SSE消息常依赖DB、Redis、RPC。如果下游超时不能让整个连接卡死。我们用Resilience4j// 配置熔断器 CircuitBreakerConfig config CircuitBreakerConfig.custom() .failureRateThreshold(50) // 错误率50%触发熔断 .waitDurationInOpenState(Duration.ofSeconds(60)) .ringBufferSizeInHalfOpenState(10) .build(); CircuitBreaker circuitBreaker CircuitBreaker.of(sse-db, config); // 发送消息时包装 public FluxMessage getMessages(String userId, String lastId) { return Mono.fromSupplier(() - { // 从DB查消息带熔断 return circuitBreaker.executeSupplier(() - messageRepository.findByUserIdAndIdGreaterThan(userId, lastId) ); }) .onErrorResume(throwable - { // 熔断时从Redis快照读取降级 log.warn(DB fallback to Redis for user: {}, userId, throwable); return Mono.just(getMessagesFromRedis(userId, lastId)); }) .flux(); }降级消息格式统一{ type: DEGRADED, reason: db_timeout, data: 服务暂时不可用请稍后重试 }前端收到DEGRADED类型可展示友好提示而非白屏。5. 生产级部署实操VCP服务器部器TomcatSpring Boot三阶调优5.1 VCP服务器部器Nginx配置详解我们线上VCP服务器部器配置精简版upstream sse_backend { server 10.0.1.10:8080 max_fails3 fail_timeout30s; server 10.0.1.11:8080 max_fails3 fail_timeout30s; keepalive 32; # 保持32个长连接 } server { listen 443 ssl; server_name api.example.com; location /api/events { proxy_pass http://sse_backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; proxy_set_header Connection upgrade; proxy_set_header Host $host; proxy_set_header X-Real-IP $remote_addr; proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for; # SSE关键配置 proxy_cache off; proxy_buffering off; proxy_read_timeout 300; proxy_send_timeout 300; proxy_connect_timeout 30; # 防DDoS限制单IP连接数 limit_conn addr 100; limit_conn_status 503; } }注意limit_conn addr 100是硬性保护防止恶意刷连接。addr基于$remote_addrVPC内网IP直连无需担心代理IP问题。5.2 Tomcat JVM参数调优专为长连接优化catalina.sh添加JAVA_OPTS$JAVA_OPTS -Xms4g -Xmx4g JAVA_OPTS$JAVA_OPTS -XX:UseG1GC -XX:MaxGCPauseMillis200 JAVA_OPTS$JAVA_OPTS -XX:UnlockExperimentalVMOptions -XX:UseCGroupMemoryLimitForHeap JAVA_OPTS$JAVA_OPTS -Djava.net.preferIPv4Stacktrue # 关键减少GC对长连接影响 JAVA_OPTS$JAVA_OPTS -XX:G1NewSizePercent30 -XX:G1MaxNewSizePercent60 JAVA_OPTS$JAVA_OPTS -XX:G1HeapRegionSize2M线程池配置server.xmlExecutor namesse-executor namePrefixsse-executor- maxThreads200 minSpareThreads50 prestartminSpareThreadstrue maxIdleTime60000 queueCapacity1000/ Connector executorsse-executor port8080 protocolorg.apache.coyote.http11.Http11NioProtocol connectionTimeout20000 keepAliveTimeout60000 maxKeepAliveRequests100 asyncTimeout300000 maxConnections10000/5.3 Spring Boot Actuator监控埋点暴露关键指标management: endpoints: web: exposure: include: health,metrics,prometheus,threaddump endpoint: metrics: show-details: always prometheus: enabled: true自定义SSE指标Component public class SseMetrics { private final Counter activeConnections; private final Timer sendLatency; public SseMetrics(MeterRegistry registry) { this.activeConnections Counter.builder(sse.connections.active) .description(Active SSE connections) .register(registry); this.sendLatency Timer.builder(sse.send.latency) .description(SSE message send latency) .register(registry); } public void incrementActive() { activeConnections.increment(); } public void recordSend(long durationMs) { sendLatency.record(durationMs, TimeUnit.MILLISECONDS); } }Prometheus查询示例# 当前活跃连接数 sum(rate(sse_connections_active_total[5m])) # 发送延迟P95 histogram_quantile(0.95, rate(sse_send_latency_seconds_bucket[5m]))6. 常见问题与排查技巧实录从日志到火焰图的实战指南6.1 问题速查表高频报错与根因定位报错信息根因排查命令解决方案stream disconnected before completion: idle timeout waiting for sseNginxproxy_read_timeout TomcatasyncTimeoutnginx -t nginx -s reload调大Nginx超时检查Tomcat配置java.lang.IllegalStateException: The emitter has already been completed连接断开后仍调用send()grep SseEmitter.send catalina.out | tail -20加try-catch或用isCompleted()判断OutOfMemoryError: unable to create new native threadCachedThreadPool线程数爆炸ps -eLf | grep java | wc -l改用有界线程池设maxPoolSizeConnection reset by peer客户端强制关闭连接服务端未捕获tcpdump -i any port 8080 -w sse.pcap在SseEmitter.onError()里清理资源Redis Stream消息重复消费XREADGROUP未ACKXRANGE sse:stream:123 - COUNT 10检查消费者组ACK逻辑用XACK6.2 日志分析三板斧从access.log到gc.logStep 1Nginx access.log定位异常IP# 查找5xx错误最多的IP awk $9500 {print $1} /var/log/nginx/access.log \| sort \| uniq -c \| sort -nr \| head -10 # 查看某IP的SSE请求详情 grep 192.168.1.100 /var/log/nginx/access.log \| grep /api/events \| tail -20Step 2Tomcat catalina.out抓线程堆栈# 找出SSE相关线程 jstack -l pid \| grep -A 20 sse-async # 查看线程阻塞点 jstack -l pid \| grep -A 10 BLOCKED \| grep SseStep 3GC日志分析内存泄漏开启GC日志-XX:PrintGCDetails -XX:PrintGCTimeStamps -Xloggc:/opt/logs/gc.log用gceasy.io上传分析重点关注Survivor区占用率持续100% → 对象没被回收Metaspace持续增长 → 动态代理类加载泄漏GC pause时间1s → 堆太大或GC策略不当。6.3 火焰图实战定位SSE发送性能瓶颈生成火焰图# 安装async-profiler wget https://github.com/jvm-profiling-tools/async-profiler/releases/download/v2.9/async-profiler-2.9-linux-x64.tar.gz tar -xzf async-profiler-2.9-linux-x64.tar.gz # 采集30秒CPU火焰图 ./profiler.sh -e cpu -d 30 -f /tmp/flame.html pid常见瓶颈点org.springframework.web.servlet.mvc.method.annotation.SseEmitter.send占比高 → 消息序列化慢换JacksonObjectWriter复用java.io.OutputStream.write占比高 → 网络IO慢检查Nginx缓冲区redis.clients.jedis.Jedis.get占比高 → Redis查询慢加索引或改用Pipeline。6.4 实操心得那些文档里不会写的坑坑1Spring Boot 2.6 的WebMvcConfigurer失效新版本WebMvcConfigurer的addInterceptors()不生效必须用Bean WebMvcRegistrations注册拦截器否则SseEmitter的onCompletion不触发。坑2Docker容器内Nginx DNS解析超时VCP服务器部器用upstream域名容器内DNS慢导致连接失败。解决方案docker run --dns10.0.0.2指定内网DNS或用/etc/hosts硬编码。坑3iOS Safari的EventSource兼容性iOS 15.4以下版本EventSource不支持withCredentials:true。解决方案改用fetch readableStream手动解析SSE或降级为WebSocket。坑4K8s Service的sessionAffinity: ClientIP不生效VCP服务器部器层做了会话保持但K8s Service默认轮询。必须加service: sessionAffinity: ClientIP sessionAffinityConfig: clientIP: timeoutSeconds: 10800最后分享个小技巧上线前用wrk压测SSE连接稳定性# 模拟1000并发持续5分钟 wrk -t12 -c1000 -d300s --latency http://your-api.com/api/events?userId1关注Latency Distribution里99%是否1sRequests/sec是否稳定。如果Connect错误率1%立刻检查Nginx和Tomcat超时配置。这套方案我们已沉淀为内部vcptoolbox的SSE-PRO模块开箱即用。它不追求“最炫技术”只解决一个问题让SSE在生产里像自来水一样稳定流淌。
返回列表