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

Spring WebFlux中Sink.asFlux()与SSE结合的重连问题及方案选型

问题分析与方案选型

首先明确问题根源:你用的share()本质是publish().refCount(1),当最后一个浏览器订阅者断开时,它会取消对上游Sink的订阅。而默认配置的Multicast Sink(onBackpressureBuffer())会在所有订阅者取消时自动终止(autoCancel=true),导致后续新订阅者无法获取任何新事件——因为Sink已经处于终止状态,不再接收新的TaskEvent。

下面逐个分析三个方案的逻辑、优缺点,以及最优选择:


方案1:关闭Sink的autoCancel

原理

修改Sink的onBackpressureBuffer参数,将autoCancel设为false。这样即使所有订阅者断开,Sink也不会自动终止,后续依然能接收新的TaskEvent,新订阅者接入后可获取后续事件。

优缺点

  • ✅ 直接从Sink层面解决问题,无需修改Flux发布逻辑,代码改动最小
  • ❌ 无订阅者时,Sink的缓冲区会持续积累TaskEvent,若长时间无订阅,可能导致内存占用过高甚至OOM
  • ❌ 若事件产生速度超过缓冲区大小,发送TaskEvent的线程会被阻塞(默认溢出策略是BUFFER),可能影响业务逻辑

方案2:用replay(0).autoConnect()替代share()

原理

replay(0)创建一个不缓存任何历史事件的ConnectableFlux,autoConnect()会在第一个订阅者接入时建立与上游的连接,且即使所有订阅者断开,连接也不会终止。新订阅者只能获取订阅之后产生的事件(因为replay(0)不缓存历史)。

优缺点

  • ✅ 解决了上游断开的问题,新订阅者可正常接入
  • ✅ 无订阅者时,buffer操作仍会每秒运行一次,消费Sink中的事件,避免Sink缓冲区溢出
  • ❌ replay(0)本质是冗余配置——因为我们不需要重播任何历史事件,用它不如直接用publish()直观

方案3:用publish().autoConnect()替代share()

原理

publish()创建一个无缓存的广播Flux,autoConnect()同样在第一个订阅者接入时建立永久连接,所有订阅者断开后上游连接依然保持活跃。新订阅者只能获取订阅后的新事件,完全符合SSE场景下“重新连接只收新事件”的需求。

优缺点

  • ✅ 完美解决问题:上游不会因订阅者断开而终止,新订阅者可正常接收后续事件
  • ✅ 无订阅者时,buffer持续消费Sink事件,避免Sink缓冲区溢出,同时不会缓存无用的历史数据
  • ✅ 代码逻辑直观,贴合业务需求(不需要重播历史)

最优选择:方案3

方案3是最贴合SSE场景的选型:

  1. 它从根本上避免了share()导致的上游断开问题,确保新订阅者始终能接入
  2. 无订阅者时,既不会让Sink缓冲区无限膨胀,也不会产生冗余的缓存逻辑
  3. 代码改动小,逻辑清晰,符合Reactor的最佳实践

如果你的业务有“重新连接时需要获取断开期间的历史事件”的需求,可考虑方案1,但必须配合缓冲区大小限制和过期清理机制,避免内存问题。否则方案3是绝对最优解。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 10:51:12