为何Reactor BoundedElastic调度器线程数过少会触发任务容量超限异常?
问题描述
我是ProjectReactor新手(使用reactor-core:3.4.18),尝试并行化Flux消费者订阅。当创建最大线程数为2的Scheduler时出现任务容量超限异常,设为4则正常运行。
相关代码:
Scheduler schedulers = Schedulers.newBoundedElastic(2, 2, "PublishedThread"); Flux.range(1, 10) .parallel() .runOn(schedulers) .doOnNext(e -> printName(e)) .subscribe();
抛出的异常信息:
[ERROR] (main) Operator called default onErrorDropped - reactor.core.Exceptions$ErrorCallbackNotImplemented: reactor.core.Exceptions$ReactorRejectedExecutionException: Task capacity of bounded elastic scheduler reached while scheduling 1 tasks (3/2) reactor.core.Exceptions$ErrorCallbackNotImplemented: reactor.core.Exceptions$ReactorRejectedExecutionException: Task capacity of bounded elastic scheduler reached while scheduling 1 tasks (3/2) Caused by: reactor.core.Exceptions$ReactorRejectedExecutionException: Task capacity of bounded elastic scheduler reached while scheduling 1 tasks (3/2) at reactor.core.Exceptions.failWithRejected(Exceptions.java:277)
请问为何线程数较少时会触发该异常?
原因分析
这是因为Schedulers.newBoundedElastic()的第二个参数是任务队列的容量上限,你设置的队列容量是2,同时最大线程数也是2,调度器能承载的总任务量(正在执行+等待排队)是2+2=4。
而Flux.parallel()默认并行度等于CPU核心数(通常至少4核),这会有多个并行流同时向调度器提交任务。当提交的任务总数超过调度器的总承载量时,就会触发拒绝策略,抛出任务容量超限的异常。
当你把最大线程数设为4时,总承载量变成4+2=6,默认并行度下同时提交的任务数不会超过这个值,因此不会触发异常。如果并行度更高(比如超过6),还是会出现同样的问题。
另外,你可以通过Flux.parallel(int parallelism)手动指定并行度,比如把并行度设为2,配合线程数2、队列容量2的调度器,总承载量4,此时10个元素会被逐步处理,队列不会被填满,就不会触发异常。
核心要点:
newBoundedElastic(maxThreads, queueCapacity)的总可容纳任务数是maxThreads + queueCapacity- 当
Flux.parallel()的并行度超过调度器总承载量时,任务提交速度会超过调度器的处理+排队能力,触发任务拒绝
内容的提问来源于stack exchange,提问作者Vijay Manohar
相关产品推荐
相关产品推荐

