You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 00:52:40