Multicast onBackpressureBuffer Sink取消订阅后无法接收新订阅问题咨询
Reactor Sink取消最后订阅后无法重新订阅的原因解析
问题现象
当取消Sink的最后一个订阅者后,再次订阅时既无法正常接收事件,也未抛出异常,与预期不符。测试代码及输出如下:
测试代码
void playground() { Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer(); var disposer = sink.asFlux().subscribe(s -> System.out.println("1: "+s)); System.out.println("CurrentSubs:" + sink.currentSubscriberCount()); sink.tryEmitNext("Test1"); var disposer2 = sink.asFlux().subscribe(s -> System.out.println("2: "+s)); System.out.println("CurrentSubs:" + sink.currentSubscriberCount()); sink.tryEmitNext("Test2"); disposer.dispose(); System.out.println("CurrentSubs:" + sink.currentSubscriberCount()); sink.tryEmitNext("Test3"); disposer2.dispose(); System.out.println("CurrentSubs:" + sink.currentSubscriberCount()); var disposer3 = sink.asFlux().subscribe(s -> System.out.println("3: "+s)); System.out.println("CurrentSubs:" + sink.currentSubscriberCount()); sink.tryEmitNext("Test4"); disposer3.dispose(); System.out.println("CurrentSubs:" + sink.currentSubscriberCount()); }
测试输出
1: Test1 CurrentSubs:2 1: Test2 2: Test2 CurrentSubs:1 2: Test3 CurrentSubs:0 CurrentSubs:0 CurrentSubs:0
核心原因
你使用的Sinks.many().multicast().onBackpressureBuffer()底层依赖MulticastProcessor,它默认开启自动终止逻辑:当最后一个订阅者取消订阅时,MulticastProcessor会自动触发onComplete信号,进入终止状态。
一旦处理器进入终止状态:
- 后续新订阅者会立即收到
onComplete信号,不会触发onNext回调,因此看不到Test4的输出。 currentSubscriberCount()返回0,因为订阅者在收到完成信号后会自动取消订阅。- 订阅操作本身合法,Reactor不会抛出异常——订阅已终止的流属于正常场景,只会触发完成回调而非错误回调。
解决方案
如果需要在最后一个订阅者取消后仍能重新订阅并接收新事件,可以通过autoCancel(false)关闭自动终止逻辑:
Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer().autoCancel(false);
修改后,即使所有订阅者取消,Sink也不会终止,后续新订阅者可以正常接收之后发送的事件。
内容的提问来源于stack exchange,提问作者user17780952
相关产品推荐
相关产品推荐

