SseEmitter需设超时、注册回调清理、禁用阻塞操作;必须用produces指定text/event-stream;防重复连接需ConcurrentHashMap或Redis去重;纯SSE推送优先选SseEmitter而非ResponseBodyEmitter。
怎么用
建立一次有效连接
Spring Boot 里最直接的 SSE 入口就是
,它不是即发即走的响应体,而是一个“可多次写入、需手动结束”的长生命周期对象。你调用
,客户端立刻收到一条消息;但如果你忘了调用
或
,连接就一直挂着,内存泄漏风险立马出现。
常见错误现象:接口压测时内存飙升、GC 频繁、
越来越大却没人清理。
必须为每个连接设置超时时间,例如
(5 分钟),避免客户端断开后服务端还傻等
务必注册
、
、
回调,并在回调里从缓存(如
)中移除该
不要在 Controller 方法里直接 new Thread 并 sleep —— 这会阻塞 Tomcat 线程池;改用异步线程或
(注意配置线程池)
为什么不能只靠
就完事
是更底层的抽象,而
是它的子类,专为 text/event-stream 格式封装好了事件头(
、
、
)、换行符和编码处理。如果你硬用
,就得自己拼
,稍有不慎就触发浏览器解析失败——表现为 EventSource 反复重连但收不到任何
事件。
使用场景:只有当你需要混推非标准格式(比如自定义二进制流)时才考虑
;纯 SSE 推送,请无条件选
。
自动加
前缀和双换行,安全
若手动 send 字符串,必须确保末尾是
,否则 Chrome/Firefox 会卡住不触发
别忽略
,缺了这个,Nginx 或网关可能截断流
怎么防止同一用户重复建立连接
用户刷新页面、网络抖动、移动端切后台再切回,都可能导致多个
实例同时连向后端。如果服务端不做去重,一个
就对应 N 个
,推送一次消息就发 N 遍,CPU 和带宽白耗。
关键不是“禁止重复创建”,而是“识别已有连接并复用或拒绝”。实际部署中还要考虑集群问题——单机 Map 拦不住其他节点上的连接。
本地去重:用
缓存
,每次请求先
,命中则直接返回已存在的
集群去重:必须引入外部协调,比如 Redis + SETNX 或分布式锁;或者干脆把连接状态下沉到 MQ,由统一调度中心分发消息
前端配合:在
实例销毁前调用
,并在重连时带上唯一
,便于服务端做幂等判断
WebFlux 的
和传统
怎么选
如果你的业务本身就是响应式栈(比如数据来自
流、R2DBC、或要对接 Kafka Flux),那直接上
更自然,没有线程切换开销,也天然支持背压。但如果你只是定时查 DB、发几条日志、推个告警,硬套 WebFlux 反而绕远路——得学
、
、
一堆操作符,还容易写出阻塞代码(比如在
里调 JDBC)。
性能影响很实在:WebFlux 版本在万级并发下资源占用更低;但开发成本高,调试链路长,出错堆栈难读。
新项目 + 数据源已是响应式(如 MongoDB Reactive、PostgreSQL R2DBC)→ 优先
老系统改造、简单轮询、DB 查询为主 → 用
更稳,别为了“响应式”而响应式
别在
的
或
里写
或同步 IO,这会让整个 event loop 卡住
真正麻烦的从来不是怎么发第一条消息,而是怎么保证第 1000 条不丢、第 10000 个连接不垮、凌晨三点报警时你能快速定位是前端没监听
,还是 Redis 锁失效导致重复推送。
SseEmitterSseEmitteremitter.send()emitter.complete()emitter.completeWithError()emitterMapnew SseEmitter(300000L)onTimeoutonCompletiononErrorConcurrentHashMapemitter@AsyncResponseBodyEmitterResponseBodyEmitterSseEmitterdata:id:event:ResponseBodyEmitter"data: hello\n\n"messageResponseBodyEmitterSseEmitterSseEmitter.send(SseEventBuilder)data:\n\nonmessageproduces = MediaType.TEXT_EVENT_STREAM_VALUEEventSourceuserIdSseEmitterConcurrentHashMapuserId → emitterget()emitterEventSource.close()lastEventIdFluxSseEmitterReactorFlux> Flux.generatedelayElementsdoOnNextmapFluxSseEmitterFluxmapflatMapThread.sleep()error