SseEmitter:从入门到源码原理深度解析

SseEmitter:从入门到源码原理深度解析
SseEmitter从入门到源码原理深度解析前言一、SSE与SseEmitter核心认知1.1 什么是SSE1.2 SseEmitter是什么1.3 SseEmitter vs WebSocket二、快速上手5分钟实现SSE推送2.1 引入依赖2.2 基础示例最简单的SSE推送2.3 客户端HTML示例三、进阶实践3.1 使用SseEventBuilder发送标准化事件3.2 连接管理存储与广播3.3 心跳机制3.4 超时时间配置四、源码原理深度解析4.1 整体流程概览4.2 建立连接Servlet异步处理机制4.3 发送数据流式写入4.4 结束连接四种终止方式4.5 超时监控原理4.6 SseEmitter vs ResponseBodyEmitter五、常见问题与解决方案5.1 数据无法实时到达客户端5.2 中文乱码5.3 内存泄漏5.4 连接数限制六、最佳实践总结6.1 核心要点回顾6.2 代码规范清单6.3 原理一句话总结前言在构建Web应用时我们经常会遇到需要服务器主动向客户端推送数据的场景——实时通知、股票行情、系统监控、大模型流式响应……传统的解决方案是轮询Polling客户端每隔几秒就发一次请求查询最新数据。这种方式不仅浪费带宽和服务器资源还存在明显的延迟。Server-Sent EventsSSE作为一种基于HTTP的轻量级服务器推送技术允许服务器通过一个长连接持续向客户端发送数据流。与WebSocket不同SSE是单向通信仅服务器→客户端但正因如此它更轻量、更易实现、与现有HTTP基础设施无缝集成。在Spring Boot中SseEmitter是Spring框架为SSE提供的高级抽象。本文将从代码实践到源码原理全面讲解Spring Boot中如何使用SseEmitter实现实时推送。一、SSE与SseEmitter核心认知1.1 什么是SSESSEServer-Sent Events是一种基于HTTP的服务器推送技术。其工作流程如下客户端通过普通HTTP请求发起连接服务器保持连接打开不关闭响应服务器在任意时刻通过该连接向客户端发送数据客户端通过EventSourceAPI接收事件SSE使用标准的text/event-stream媒体类型数据格式规范为data: 这是消息内容\n\n每个事件以data:开头以两个换行符\n\n结尾。1.2 SseEmitter是什么SseEmitter是Spring Framework 4.2引入的一个类它专门用于发送Server-Sent Events。它是ResponseBodyEmitter的子类在通用异步响应流的基础上增加了SSE协议的专用支持。核心定位SseEmitter充当了控制器与底层HTTP连接之间的桥梁将开发者从繁琐的SSE协议处理中解放出来。1.3 SseEmitter vs WebSocket对比维度SseEmitterSSEWebSocket通信方向单向服务器→客户端双向客户端↔服务器协议HTTP标准协议独立的WS协议自动重连浏览器内置支持需手动实现实现复杂度简单较复杂适用场景通知推送、监控大屏、AI流式响应聊天、游戏、双向交互应用一句话选型如果只需要服务器向客户端推送数据选SseEmitter如果需要双向实时通信选WebSocket。二、快速上手5分钟实现SSE推送2.1 引入依赖SseEmitter是spring-web模块的一部分只需引入Spring Boot Web Starter即可dependencygroupIdorg.springframework.boot/groupIdartifactIdspring-boot-starter-web/artifactId/dependency2.2 基础示例最简单的SSE推送packagecom.example.demo.controller;importorg.springframework.http.MediaType;importorg.springframework.web.bind.annotation.GetMapping;importorg.springframework.web.bind.annotation.RestController;importorg.springframework.web.servlet.mvc.method.annotation.SseEmitter;importjava.io.IOException;importjava.util.concurrent.ExecutorService;importjava.util.concurrent.Executors;RestControllerpublicclassSseController{privatefinalExecutorServiceexecutorExecutors.newSingleThreadExecutor();GetMapping(value/sse,producesMediaType.TEXT_EVENT_STREAM_VALUE)publicSseEmitterhandleSse(){// 创建SseEmitter设置超时时间为60秒SseEmitteremitternewSseEmitter(60_000L);// 在独立线程中异步发送数据释放请求线程executor.execute(()-{try{for(inti0;i10;i){// 发送数据emitter.send(Message i);Thread.sleep(1000);}// 发送完成关闭连接emitter.complete();}catch(IOException|InterruptedExceptione){// 发生错误时关闭连接emitter.completeWithError(e);}});returnemitter;}}代码要点produces MediaType.TEXT_EVENT_STREAM_VALUE必须指定响应类型为SSEnew SseEmitter(60_000L)设置超时时间为60秒异步发送使用独立线程发送数据避免阻塞Tomcat请求线程emitter.send()向客户端发送一条消息emitter.complete()正常完成关闭连接emitter.completeWithError(e)异常时关闭连接2.3 客户端HTML示例!DOCTYPEhtmlhtmlheadmetacharsetUTF-8titleSSE Demo/title/headbodyh1SSE 实时消息/h1dividmessages/divscript// 建立SSE连接consteventSourcenewEventSource(/sse);// 接收消息eventSource.onmessagefunction(event){constdivdocument.getElementById(messages);constpdocument.createElement(p);p.textContentevent.data;div.appendChild(p);};// 错误处理eventSource.onerrorfunction(){console.error(SSE连接出错);};/script/body/html浏览器原生支持EventSourceAPI连接建立后会自动接收服务器推送的消息。当连接断开时浏览器会自动重连。三、进阶实践3.1 使用SseEventBuilder发送标准化事件SseEmitter提供了事件构建器SseEventBuilder可以发送带有id、event类型、retry重连时间等标准化字段的SSE事件GetMapping(value/sse/event,producesMediaType.TEXT_EVENT_STREAM_VALUE)publicSseEmittersseWithEvent(){SseEmitteremitternewSseEmitter(60_000L);executor.execute(()-{try{// 发送标准化SSE事件emitter.send(SseEmitter.event().id(123)// 事件ID.name(notification)// 事件类型.data({\msg\:\你好\})// 事件数据.reconnectTime(3000)// 客户端重连间隔(毫秒));Thread.sleep(1000);emitter.complete();}catch(Exceptione){emitter.completeWithError(e);}});returnemitter;}字段说明id事件ID客户端断线重连时可从此ID之后继续接收name事件类型客户端可针对不同类型注册不同监听器data事件数据reconnectTime告诉客户端断开后等待多少毫秒再重连3.2 连接管理存储与广播在实际项目中需要管理多个SSE连接实现定向推送或广播packagecom.example.demo.service;importlombok.extern.slf4j.Slf4j;importorg.springframework.stereotype.Service;importorg.springframework.web.servlet.mvc.method.annotation.SseEmitter;importjava.io.IOException;importjava.util.concurrent.ConcurrentHashMap;importjava.util.concurrent.CopyOnWriteArrayList;Slf4jServicepublicclassSseService{// 存储所有连接用于广播privatefinalCopyOnWriteArrayListSseEmitteremittersnewCopyOnWriteArrayList();// 存储用户ID - 连接用于定向推送privatefinalConcurrentHashMapString,SseEmitteruserEmittersnewConcurrentHashMap();/** * 建立连接 */publicSseEmitterconnect(StringuserId){SseEmitteremitternewSseEmitter(300_000L);// 5分钟超时// 注册完成回调连接关闭时自动清理emitter.onCompletion(()-{emitters.remove(emitter);userEmitters.remove(userId,emitter);log.info(SSE连接关闭: userId{},userId);});// 注册超时回调emitter.onTimeout(()-{emitters.remove(emitter);userEmitters.remove(userId,emitter);log.warn(SSE连接超时: userId{},userId);});// 注册错误回调emitter.onError(e-{emitters.remove(emitter);userEmitters.remove(userId,emitter);log.error(SSE连接错误: userId{},userId,e);});emitters.add(emitter);userEmitters.put(userId,emitter);log.info(SSE连接建立: userId{},userId);returnemitter;}/** * 广播消息给所有连接 */publicvoidbroadcast(Objectmessage){emitters.forEach(emitter-{try{emitter.send(message);}catch(IOExceptione){// 发送失败说明连接已断开从集合中移除emitters.remove(emitter);log.warn(广播消息失败移除连接);}});}/** * 定向推送消息给指定用户 */publicvoidsendToUser(StringuserId,Objectmessage){SseEmitteremitteruserEmitters.get(userId);if(emitter!null){try{emitter.send(message);}catch(IOExceptione){userEmitters.remove(userId,emitter);log.warn(定向推送失败移除连接: userId{},userId);}}}}关键设计使用CopyOnWriteArrayList和ConcurrentHashMap保证线程安全在onCompletion、onTimeout、onError回调中自动清理资源防止内存泄漏send()方法可能抛出IOException需捕获并清理连接3.3 心跳机制SSE长连接可能因为网络中间件如Nginx的超时配置而被意外断开。心跳机制通过定期发送空数据或注释保持连接活跃ServiceSlf4jpublicclassHeartbeatService{privatefinalScheduledExecutorServiceschedulerExecutors.newSingleThreadScheduledExecutor();privatefinalMapString,SseEmitteremittersnewConcurrentHashMap();PostConstructpublicvoidstartHeartbeat(){// 每15秒发送一次心跳scheduler.scheduleAtFixedRate(()-{emitters.forEach((userId,emitter)-{try{// SSE规范中以:开头的行表示注释不会被客户端触发事件// 但可以保持连接活跃emitter.send(SseEmitter.event().comment(heartbeat));}catch(IOExceptione){// 心跳发送失败说明连接已断开emitters.remove(userId);log.debug(心跳失败移除连接: userId{},userId);}});},5,15,TimeUnit.SECONDS);}}⚠️ 重要提醒不要使用while(true) Thread.sleep()方式实现心跳——当客户端异常断开时服务端无法立即感知会导致线程持续运行、内存泄漏。应使用线程池统一管理心跳任务。3.4 超时时间配置SseEmitter的超时时间直接影响连接的存活时长// 默认超时由Spring MVC配置或容器默认值决定SseEmitteremitter1newSseEmitter();// 自定义超时30秒SseEmitteremitter2newSseEmitter(30_000L);// 永不超时需谨慎使用SseEmitteremitter3newSseEmitter(0L);配置建议场景服务器超时建议客户端retry建议高频消息推送如监控大屏300秒5分钟1-3秒低频但需保持连接60秒1分钟3-5秒大模型流式响应根据模型响应时长动态调整1-3秒关键原则服务器超时应显著大于客户端重试时间避免客户端重连时服务器已超时关闭。四、源码原理深度解析4.1 整体流程概览SseEmitter的工作流程分为三个阶段客户端发起GET请求 → Controller返回SseEmitter → Spring绑定HTTP响应不关闭连接→ 请求线程释放 → 后台线程调用send()推送数据 → 调用complete()关闭连接4.2 建立连接Servlet异步处理机制传统同步模型的问题每个请求占用一个Tomcat线程直到响应完成长连接场景下线程资源迅速耗尽。SseEmitter基于Servlet 3.0的异步处理机制AsyncContext解决了这个问题当Controller返回SseEmitter时Spring MVC的处理流程如下拦截返回类型Spring检测到返回类型是SseEmitter立即返回HTTP响应头HTTP/1.1 200 OK Content-Type: text/event-stream Cache-Control: no-cache Connection: keep-alive绑定输出流将HttpServletResponse的输出流与SseEmitter绑定释放请求线程Controller方法执行完毕后请求线程立即释放回线程池不阻塞等待保持连接打开底层TCP连接保持打开状态由SseEmitter在后台维护本质SseEmitter不是真正的非阻塞I/O而是“释放请求线程 异步写响应”的模式。4.3 发送数据流式写入当调用emitter.send()时数据被按SSE格式写入底层ServletResponse.getWriter()格式示例data: {msg:hello}\n\n数据通过已保持打开的HTTP连接流式发送给客户端⚠️ 关键限制SseEmitter不是线程安全的多个线程同时调用send()可能导致数据交错或异常。解决方案使用队列、synchronized或确保串行发送。4.4 结束连接四种终止方式连接可以通过以下方式终止方式触发时机回调顺序服务端主动完成调用emitter.complete()onCompletion()服务端发生异常调用emitter.completeWithError(e)onError()→onCompletion()超时超过设定的超时时间无活动onTimeout()→onCompletion()客户端断开客户端关闭浏览器/网络断开可能触发onError()或无回调客户端断开检测的局限性服务器无法实时感知客户端断开。只有在下次调用send()尝试写入socket时才会抛出IOException。因此必须在send()外层加try-catch并清理emitter。4.5 超时监控原理Spring内部使用**ScheduledExecutorService** 监控每个SseEmitter的超时从最后一次成功send()或连接建立开始计时超时后自动调用onTimeout()然后complete()释放Tomcat/Undertow等容器中的连接资源4.6 SseEmitter vs ResponseBodyEmitter两者都基于Servlet异步机制但用途不同特性SseEmitterResponseBodyEmitter用途专用于SSE通用异步响应流Content-Type固定为text/event-stream可自定义数据格式自动封装为SSE格式由HttpMessageConverter序列化事件支持支持id、event、retry等无事件概念五、常见问题与解决方案5.1 数据无法实时到达客户端问题调用了emitter.send()但客户端迟迟收不到数据。原因可能是响应缓冲区未刷新或代理服务器如Nginx缓冲了响应。解决方案确保Controller方法标注了produces MediaType.TEXT_EVENT_STREAM_VALUE在Nginx配置中关闭代理缓冲proxy_buffering off;5.2 中文乱码问题推送的中文内容在客户端显示为乱码。解决方案确保响应编码为UTF-8GetMapping(value/sse,producestext/event-stream;charsetUTF-8)5.3 内存泄漏问题随着连接数增加内存持续增长。原因未在onCompletion/onTimeout/onError中移除已关闭的连接使用while(true)方式实现心跳线程无法释放解决方案在所有生命周期回调中主动清理连接引用使用线程池统一管理心跳任务5.4 连接数限制问题浏览器对每个域名的SSE并发连接数有限制通常为6个。解决方案使用HTTP/2多路复用使用WebWorker或SharedWorker共享连接评估是否真的需要如此多的并发SSE连接六、最佳实践总结6.1 核心要点回顾组件/概念作用SseEmitterSpring提供的SSE核心类管理HTTP长连接produces TEXT_EVENT_STREAM_VALUE必须指定声明响应类型为SSE流异步线程池释放请求线程由后台线程推送数据生命周期回调onCompletion/onTimeout/onError自动清理资源心跳机制定期发送数据保持连接活跃防止中间件超时断开6.2 代码规范清单Controller方法必须标注produces MediaType.TEXT_EVENT_STREAM_VALUE必须设置超时时间避免连接永久占用资源必须使用异步线程发送数据不要阻塞请求线程必须注册生命周期回调在连接关闭时清理资源必须捕获send()的IOException处理客户端断开场景使用线程池实现心跳不要用while(true) sleep()6.3 原理一句话总结SseEmitter本质上是对Servlet 3.0异步机制AsyncContext的高级封装——Controller返回SseEmitter后Spring立即释放请求线程并保持HTTP连接打开后续通过后台线程调用send()将数据按SSE协议格式data: ...\n\n流式写入响应输出流连接的生命周期通过onCompletion、onTimeout、onError回调进行管理超时由Spring内置的ScheduledExecutorService监控客户端断开则通过send()时的IOException被动感知。