Project Reactor带背压缓冲区的Sink异常行为问询
场景描述
使用默认大小(256)的带背压缓冲区的多播Sink,通过专用线程每25ms异步发送事件。三个订阅者配置如下:
- B:通过线程休眠模拟慢业务处理(每次休眠500ms),使用
publishOn绑定专用调度器 - A、C:使用
subscribeOn绑定专用调度器,无阻塞消费
核心问题
- 为什么A和C的
subscribeOn未生效,事件仍在发送线程上发布? - A和C在事件257时停止接收,此时B已处理约15个事件,按理应在约256+15=271时停止,原因是什么?
- 发送i=512时开始出现
FAIL_OVERFLOW失败,是否因为B的publishOn缓冲区(256)和Sink的缓冲区(256)都已满?但此时B已处理约29个事件,这并不成立,原因是什么? - B处理到192个事件后,A和C恢复接收事件,为何是192?且事件在B的调度器线程上接收,原因是什么?
- 将B的
publishOn改为subscribeOn后,发送线程会适配最慢的消费者(B),但此时subscribeOn仍未生效,原因是什么?
问题解答
1. subscribeOn未生效的原因
subscribeOn的核心作用是指定订阅逻辑(从下游到上游的订阅信号传递)的执行线程,以及冷发布者中事件生成的线程。但Sinks.Many.multicast()是热发布者,事件推送逻辑由发送线程调用tryEmitNext主动触发,直接将事件推送给所有订阅者的onNext方法,这个推送过程完全在发送线程中完成,不受subscribeOn影响——subscribeOn不会改变热发布者的事件推送线程,仅能改变订阅环节的线程。
2. A和C在257事件时停止的原因
Sinks.many().multicast().onBackpressureBuffer()的256缓冲区是多播组共享的全局缓冲区,而非每个订阅者单独拥有。背压策略取所有订阅者中最小的请求量来控制Sink的发送:
- A、C无阻塞,会立即发送
request(Long.MAX_VALUE) - B因
publishOn缓冲区(默认256)+慢处理,请求量被限流
当Sink的共享缓冲区被填满(256个事件),且B的publishOn缓冲区也已满时,Sink无法再接收新事件,因此停止向所有订阅者推送,A、C自然在257事件时停止接收。你计算的256+15是错误逻辑,共享缓冲区是全局的,不叠加单个订阅者的处理量。
3. i=512时出现FAIL_OVERFLOW的原因
Sink的共享缓冲区(256)和B的publishOn缓冲区(256)是两个独立的异步缓冲区:
- 发送线程以每25ms一个的速度推送事件,先填满Sink的共享缓冲区(256个),接着填满B的
publishOn缓冲区(256个),累计512个事件 - 此时B的处理速度远慢于发送速度,
publishOn缓冲区无法释放空间,导致B停止向Sink发送请求,Sink的共享缓冲区也无法再接收新事件,因此发送i=512时触发FAIL_OVERFLOW
你提到的B已处理约29个事件,是因为publishOn是异步推送,B的处理线程在处理事件时,publishOn缓冲区早已填满256个事件,不影响两个缓冲区总容量512的计算。
4. B处理到192个事件时A、C恢复的原因
- 为什么是192:
publishOn的默认缓冲区大小为256,当B处理完192个事件后,缓冲区空闲出192个位置,publishOn会向上游(Sink)发送request(192),请求填充缓冲区。Sink收到请求后,会将共享缓冲区中的事件推送给所有订阅者,A、C因此恢复接收。 - 事件在B的调度器线程上接收:Reactor中背压请求是从下游向上游传递的,请求的执行线程(即B的
publishOn调度器线程)会成为事件推送的线程,因此A、C的onNext会在该线程上执行。
5. 改为subscribeOn后看似未生效的原因
将B的publishOn改为subscribeOn后,subscribeOn的作用已经生效:B的onNext逻辑会在BBB调度器线程上执行(可通过日志线程名验证)。你觉得“未生效”是误解了subscribeOn的作用:
- 热发布者的事件推送线程仍由发送线程控制,
subscribeOn不会改变这个逻辑 - 发送线程适配B的速度,是因为B的慢处理导致其向Sink发送请求的速度变慢,Sink根据最小请求量限流,这是背压策略的结果,和
subscribeOn是否生效无关。
可运行代码
import org.slf4j.Logger; import org.slf4j.LoggerFactory; import reactor.core.publisher.Sinks; import java.util.concurrent.ThreadFactory; import java.util.concurrent.atomic.AtomicInteger; import static java.util.concurrent.Executors.newScheduledThreadPool; import static reactor.core.scheduler.Schedulers.fromExecutorService; public class SinkBackpressureBuffer { private static final Logger LOGGER = LoggerFactory.getLogger("SinkLogger"); private static void sleep(long ms) { try { Thread.sleep(ms); } catch (InterruptedException e) { throw new RuntimeException(e); } } public static void main(String[] args) throws InterruptedException { final Sinks.Many<Integer> sink = Sinks.many().multicast().onBackpressureBuffer(); new Thread(() -> { sleep(1000); for (int i = 1; i <= 550; i++) { // has to be bigger than 256 + 256 final Sinks.EmitResult emitResult = sink.tryEmitNext(i); if (emitResult != Sinks.EmitResult.OK) { LOGGER.error("Emit for i={}, res={}", i, emitResult); } else { LOGGER.info("Emit for i={}, res={}", i, emitResult); } sleep(25); } }).start(); sink.asFlux() .subscribeOn(fromExecutorService(newScheduledThreadPool(5, new NamedThreadFactory("AAA")))) // has no effect .subscribe(i -> { LOGGER.info("A: {}", i); }); sink.asFlux() .publishOn(fromExecutorService(newScheduledThreadPool(5, new NamedThreadFactory("BBB")))) .subscribe(i -> { sleep(500); LOGGER.info("B: {}", i); }); sink.asFlux() .subscribeOn(fromExecutorService(newScheduledThreadPool(5, new NamedThreadFactory("CCC")))) // has no effect .subscribe(i -> { LOGGER.info("C: {}", i); }); LOGGER.info("done"); sleep(500000); } public static class NamedThreadFactory implements ThreadFactory { private final AtomicInteger sequence = new AtomicInteger(1); private final String prefix; public NamedThreadFactory(String prefix) { this.prefix = prefix; } @Override public Thread newThread(Runnable r) { Thread thread = new Thread(r); int seq = sequence.getAndIncrement(); thread.setName(prefix + (seq > 1 ? "-" + seq : "")); thread.setDaemon(true); return thread; } } }
内容的提问来源于stack exchange,提问作者Tomask

