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

Project Reactor带背压缓冲区的Sink异常行为问询

Reactor Sink背压与调度器问题解析

场景描述

使用默认大小(256)的带背压缓冲区的多播Sink,通过专用线程每25ms异步发送事件。三个订阅者配置如下:

  • B:通过线程休眠模拟慢业务处理(每次休眠500ms),使用publishOn绑定专用调度器
  • A、C:使用subscribeOn绑定专用调度器,无阻塞消费

核心问题

  1. 为什么A和C的subscribeOn未生效,事件仍在发送线程上发布?
  2. A和C在事件257时停止接收,此时B已处理约15个事件,按理应在约256+15=271时停止,原因是什么?
  3. 发送i=512时开始出现FAIL_OVERFLOW失败,是否因为B的publishOn缓冲区(256)和Sink的缓冲区(256)都已满?但此时B已处理约29个事件,这并不成立,原因是什么?
  4. B处理到192个事件后,A和C恢复接收事件,为何是192?且事件在B的调度器线程上接收,原因是什么?
  5. 将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 18:32:29