Spring Boot中如何将SseEmitter对象持久化到数据库?
解决Spring Boot SSE生产环境中无法内存存储SseEmitter的问题
核心问题分析
SseEmitter本质是与当前HTTP请求的响应流绑定的对象,它持有底层连接、响应状态等不可序列化的资源。序列化再反序列化只会生成全新的对象,完全丢失原有的客户端连接,同时在集群生产环境中,内存存储的SseEmitter无法跨实例共享,导致事件无法正确投递到目标客户端。
生产级优雅解决方案
以下两种方案是业界常用的分布式SSE实现思路:
方案一:基于Redis Pub/Sub的分布式事件广播
思路:
- 每个服务实例仅在本地内存维护自身持有的
SseEmitter连接 - 第三方Webhook触发事件时,将事件发布到Redis的指定频道(按
tokenId维度区分) - 所有服务实例订阅对应频道,收到事件后检查本地是否持有目标
tokenId的SseEmitter,存在则直接发送消息
修改后的代码示例:
@RestController public class EventController { private final Map<String, SseEmitter> sseEmitters = new ConcurrentHashMap<>(); private final StringRedisTemplate redisTemplate; private final RedisMessageListenerContainer listenerContainer; @Autowired public EventController(StringRedisTemplate redisTemplate, RedisMessageListenerContainer listenerContainer) { this.redisTemplate = redisTemplate; this.listenerContainer = listenerContainer; // 订阅所有用户事件频道(按tokenId前缀匹配) listenerContainer.addMessageListener((message, pattern) -> { String channel = new String(message.getChannel()); String tokenId = channel.replace("sse:event:", ""); String eventData = new String(message.getBody()); SseEmitter emitter = sseEmitters.get(tokenId); if (emitter != null) { try { emitter.send(SseEmitter.event().name("latestEvent").data(eventData)); } catch (IOException e) { // 发送失败,清理无效连接 sseEmitters.remove(tokenId); } } }, new PatternTopic("sse:event:*")); } @CrossOrigin @GetMapping(value = "/subscribe", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter subscribe(@RequestParam("tokenId") String tokenId) throws IOException { // 设置合理超时时间(示例为30分钟),客户端需实现到期重连逻辑 SseEmitter sseEmitter = new SseEmitter(30 * 60 * 1000L); sseEmitter.send(SseEmitter.event().name("latestEvent").data("INITIALIZED")); // 连接生命周期回调,及时清理本地缓存 sseEmitter.onCompletion(() -> sseEmitters.remove(tokenId)); sseEmitter.onError(e -> sseEmitters.remove(tokenId)); sseEmitter.onTimeout(() -> sseEmitters.remove(tokenId)); sseEmitters.put(tokenId, sseEmitter); return sseEmitter; } @CrossOrigin @PostMapping("/dispatchEvent") public void dispatchToClients(@RequestParam("freshEvent") String freshEvent, @RequestParam("tokenId") String tokenId) { // 将事件发布到Redis对应频道 redisTemplate.convertAndSend("sse:event:" + tokenId, freshEvent); } }
方案二:基于Spring Cloud Stream的事件驱动架构
思路:
- 使用Kafka/RabbitMQ等成熟消息中间件做事件中转
- 第三方Webhook触发事件时,将事件发送到指定主题
- 每个服务实例作为消费者订阅主题,收到事件后匹配本地
SseEmitter连接发送消息
这种方案适合需要消息持久化、重试机制、流量削峰的复杂场景,依赖Spring Cloud Stream的封装可以快速实现分布式事件投递。
生产环境补充注意事项
- 连接超时与重连:务必设置
SseEmitter超时时间,同时要求客户端实现自动重连逻辑,避免长期无效连接占用资源 - 连接清理:通过
onCompletion/onError/onTimeout回调及时清理本地SseEmitter,防止内存泄漏 - 消息幂等:分布式环境下可能出现重复事件,需在客户端或服务端实现幂等校验,避免重复处理
- 监控告警:监控本地
SseEmitter连接数、消息发送成功率,异常时及时告警
内容的提问来源于stack exchange,提问作者Akshay Thakur
相关产品推荐
相关产品推荐

