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

Spring Boot SSE多订阅实现问题咨询(基于Reactor Sinks)

Spring Boot SSE 实现相关问题

我开发了一个Spring Boot后端应用,需要向前端实现SSE推送功能,要求支持多订阅者、向多用户发送多事件,且需程序化触发事件,因此采用Project Reactor的Sinks类实现,代码如下:

private final Map<Long, Many<MyNotification>> sinkMap = new ConcurrentHashMap<>();

public Flux<ServerSentEvent<MyNotification>> createSubscription(Long subscriptionId) {
    Sinks.Many<MyNotification> sink;

    if (sinkMap.get(subscriptionId) != null) {
        sink = sinkMap.get(subscriptionId);
    } else {
        sink = Sinks.many().multicast().onBackpressureBuffer();
    }

    sinkMap.put(subscriptionId, sink);

    return sink.asFlux().map(event -> ServerSentEvent.builder(event).build());
}

现咨询以下技术问题:

  1. 若应用需水平扩展,将Sinks存储在Map中是否可行?多实例并行时会有何表现?
  2. 是否可通过Redis等工具集中存储Sinks以实现订阅中心化?
  3. 如何判断订阅者是否仍处于订阅状态,以便必要时清理Map?
  4. 除Reactor Sinks外,Spring Boot中还有哪些SSE实现方案?

问题解答

1. 水平扩展时Map存储Sinks的可行性与多实例表现

完全不可行。Sinks是内存内的Reactor组件,和当前JVM进程强绑定,每个应用实例的Map都是独立隔离的。多实例部署后会出现以下问题:

  • 订阅者连接到实例A,对应的Sinks仅存在于A的本地Map;如果触发推送的请求落到实例B,B的Map里找不到该Sinks,订阅者就收不到消息。
  • 不同实例的sinkMap完全割裂,推送消息只能覆盖连接到当前实例的订阅者,无法实现全局统一的消息推送。

2. 用Redis集中存储Sinks是否可行?

不行。Sinks是Reactor的运行时内存对象,包含订阅者回调、背压处理等和当前JVM绑定的状态,根本无法序列化存储到Redis这类外部存储中。

如果要实现跨实例的订阅中心化,正确思路是用Redis做消息中转:

  • 每个应用实例启动后,订阅Redis的对应频道(按用户ID或订阅ID分组);
  • 触发推送时,将消息发送到目标Redis频道;
  • 所有订阅该频道的应用实例收到消息后,再通过本地内存中的Sinks推送给连接到自己的订阅者。

3. 判断订阅者状态并清理Map

可以通过Flux的生命周期监听操作符,在订阅终止时自动清理对应的Sinks,同时配合定时任务做兜底检查:

修改createSubscription方法,添加订阅终止监听:

public Flux<ServerSentEvent<MyNotification>> createSubscription(Long subscriptionId) {
    Sinks.Many<MyNotification> sink = sinkMap.computeIfAbsent(subscriptionId, 
        id -> Sinks.many().multicast().onBackpressureBuffer());

    return sink.asFlux()
        .map(event -> ServerSentEvent.builder(event).build())
        .doOnCancel(() -> sinkMap.remove(subscriptionId))
        .doOnError((e) -> sinkMap.remove(subscriptionId))
        .doOnComplete(() -> sinkMap.remove(subscriptionId));
}

另外可以加个定时任务,定期遍历sinkMap,通过sink.currentSubscriberCount()判断如果某个Sinks的活跃订阅数为0,就从Map中移除,避免异常终止导致的内存泄漏。

4. Spring Boot中其他SSE实现方案

  • Spring MVC的SseEmitter:Spring原生基于Servlet API的实现,适合传统MVC场景。可以通过SseEmitter对象主动向客户端发送消息,支持超时、错误处理,也能手动管理多个Emitter实例。
  • Spring WebFlux直接返回Flux:不需要Sinks的话,可直接返回Flux<ServerSentEvent<T>>,通过Flux.create()或Flux.generate()生成事件流,适合简单的推送场景。
  • 第三方消息中间件结合SSE:比如RabbitMQ、Kafka等,后端服务作为消费者接收消息,再推送给对应的SSE订阅者,适合大规模、高并发的推送场景。

内容的提问来源于stack exchange,提问作者Vlada

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 01:53:12