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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 20:12:42