Quarkus Mutiny响应式代码线程数仅10而非20,性能不及命令式代码
问题分析与修正方案
线程数量不足&性能低下的核心原因
1. 线程配置与调度逻辑不匹配
命令式代码直接创建20个独立线程,每个线程持续执行循环任务;而你的响应式代码存在以下问题:
- 如果
THREADS变量值不是20,线程池大小自然不足;即使设为20,runSubscriptionOn(executor)仅控制订阅阶段的线程,流的实际处理逻辑(repeating循环)可能没有真正占用线程池的并发资源,导致线程利用率不足。 disjoint()操作符会把processMessages(t)返回的列表拆分为单个元素逐个处理,把批量处理拆成串行单元素操作,直接拖慢处理速度。
2. 响应式模型适配错误
你用repeating().until(List::isEmpty)的模式,本质是把持续循环的任务拆成了多次流元素发射,额外增加了响应式流的分发、订阅开销,完全没有发挥响应式的优势,反而比命令式多了一层不必要的封装。
修正后的响应式代码(对齐命令式逻辑)
方案一:贴近命令式的直接任务执行
这种方式最接近原命令式代码的逻辑,保证20个并发线程持续处理任务:
// 创建大小为20的固定线程池,和命令式线程数一致 ExecutorService executor = Executors.newFixedThreadPool(20); IntStream.range(0, 20) .forEach(t -> Multi.createFrom().runner(subscriber -> { int records = RANGE; // 循环处理直到任务完成或订阅取消 while (records > 0 && !subscriber.isCancelled()) { records = processRecords(t, Thread.currentThread().getId()); } subscriber.onComplete(); }) // 指定任务运行在自定义线程池 .runSubscriptionOn(executor) .subscribe() .with(item -> System.out.println(Thread.currentThread().getName()))); // 程序结束时务必关闭线程池 // executor.shutdown();
方案二:修正repeating模式的使用
如果坚持用repeating API,需要去掉冗余操作符,调整任务逻辑:
ExecutorService executor = Executors.newFixedThreadPool(20); IntStream.range(0, 20) .forEach(t -> Multi.createBy() .repeating() // 直接返回处理结果,替代原有的uni模式 .supplier(() -> processRecords(t, Thread.currentThread().getId())) // 直到records为0时终止循环 .until(records -> records <= 0) // 指定任务运行线程 .runSubscriptionOn(executor) .subscribe() .with(item -> System.out.println(Thread.currentThread().getName())));
关键修正点说明
- 线程池大小明确设为20:和命令式的并发线程数保持一致,确保足够的并发度。
- 移除
disjoint()操作符:避免批量处理被拆分为串行单元素操作,保留批量处理的效率。 - 选择合适的流创建方式:用
runner直接封装循环任务,或者调整repeating的供应商逻辑,减少响应式流的额外开销。 - 响应式规范兼容:在循环中加入
subscriber.isCancelled()判断,支持响应式的取消与背压机制。
内容的提问来源于stack exchange,提问作者user21684014
相关产品推荐
相关产品推荐

