Sinks.Many首次订阅数为0后失效的原因及解决方案咨询
为什么multicast()的Sink会自动完成?
你使用的Sinks.many().multicast().onBackpressureBuffer()创建的是MulticastProcessor,它的默认行为是:当最后一个订阅者取消或完成订阅时,处理器会自动进入完成状态。
你的T1线程订阅用了take(5),当接收完5个元素后,take操作符会主动取消订阅,此时sink的订阅者数变为0,触发了MulticastProcessor的自动完成逻辑。所以后续T2订阅时,拿到的Flux已经处于完成状态,会立即执行onComplete回调。
为什么tryEmitNext仍返回成功?
对于已经完成的MulticastProcessor,调用tryEmitNext确实会返回SUCCESS,但发射的元素会被直接丢弃。这是该处理器的设计特性:允许在终止状态下调用发射方法,但不会对元素做任何处理,也不会抛出异常。
实现支持订阅者动态进出的持久化Sink
要让Sink在订阅者全部离开后仍保持活跃,支持后续订阅者接收新元素,有两种方案:
方案1:修改multicast()的自动取消行为
在创建multicast sink时,添加autoCancel(false)配置,关闭“最后一个订阅者离开时自动完成”的逻辑:
Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer().autoCancel(false);
这样即使所有订阅者都取消,sink仍会保持活跃状态,新的订阅者可以正常接收后续发射的元素。
方案2:使用replay()类型的Sink
你尝试的Sinks.many().replay().limit(Duration.ZERO)能正常工作,是因为ReplayProcessor默认不会因订阅者全部离开而终止。它会一直保持活跃,直到你手动调用tryEmitComplete或tryEmitError,或者sink被销毁。这种类型的Sink适合需要支持订阅者动态加入、且不需要重放历史元素的场景(limit(Duration.ZERO)表示不保留任何历史元素)。
内容的提问来源于stack exchange,提问作者kaqqao

