Quarkus中Mutiny Uni多线程并行执行不生效问题咨询
问题场景
需要将两个独立的IO/计算操作分别运行在独立线程中执行,初始组合调用代码如下:
Uni.combine().all() .unis(getItem(), getItemDetails()) .asTuple().subscribe().with(tuple -> { context.setItem(tuple.getItem1()); context.setItemDetails(tuple.getItem2()); });
初始的两个方法实现如下:
public Uni<ItemResponse> callGetItem(){ Supplier<ItemResponse> supplier = () -> itemService.getItem("item_id_1"); return Uni.createFrom().item(supplier); } public Uni<ItemDetailsResponse> callGetItemDetail(){ Supplier<ItemDetailsResponse> supplier = () -> itemService.getItemDetail("dummy_item_id"); return Uni.createFrom().item(supplier) ; }
运行后观测到callGetItem()和callGetItemDetail()两个方法全部运行在同一个线程executor-thread-0,未达到并行在不同线程执行的预期。
后续尝试自定义核心线程数为2的固定线程池,为两个Uni添加emitOn指定执行器,修改后的方法实现如下:
// 自定义固定线程池 ExecutorService executor = Executors.newFixedThreadPool(2); public Uni<ItemResponse> callGetItem(){ Supplier<ItemResponse> supplier = () -> itemService.getItem("item_id_1"); return Uni.createFrom().item(supplier).emitOn(executor); } public Uni<ItemDetailsResponse> callGetItemDetail(){ Supplier<ItemDetailsResponse> supplier = () -> itemService.getItemDetail("dummy_item_id"); return Uni.createFrom().item(supplier).emitOn(executor) ; }
修改后重新运行,两个方法仍然运行在同一个线程,未解决问题。
问题根因
Uni.createFrom().item(supplier)本身不会触发异步调度,默认直接在触发订阅的调用线程上同步执行supplier逻辑,未做线程切换配置时,两个组合的Uni会在同一个订阅链路上顺序执行,自然会落在同一个线程。- 误用了Mutiny的线程切换操作符:
emitOn的作用范围是下游消费逻辑,仅会改变该操作符之后的回调、操作符执行线程,不会影响它上游的supplier生成逻辑的执行线程,因此加了emitOn也没有把两个supplier的执行派发到线程池。 - 即使
emitOn生效,两个短任务提交到同一个线程池时,如果第一个任务执行速度极快,会在提交第二个任务前就释放线程,线程池会直接复用刚空闲的同一个线程执行第二个任务,也会出现同线程的现象,但这属于调度概率问题,不是核心根因。
正确实现方案
将emitOn替换为runSubscriptionOn即可,该操作符作用于上游生产逻辑,会在订阅触发时把当前Uni的supplier执行任务独立提交到指定线程池,两个任务会被线程池调度到空闲线程执行,只要线程池有可用线程,就能实现并行在不同线程运行的效果。
修正后的代码如下:
// 核心线程数为2的固定线程池,注意应用停止时需主动调用shutdown()回收资源 ExecutorService executor = Executors.newFixedThreadPool(2); public Uni<ItemResponse> callGetItem(){ Supplier<ItemResponse> supplier = () -> itemService.getItem("item_id_1"); // runSubscriptionOn将上游supplier执行提交到指定线程池 return Uni.createFrom().item(supplier).runSubscriptionOn(executor); } public Uni<ItemDetailsResponse> callGetItemDetail(){ Supplier<ItemDetailsResponse> supplier = () -> itemService.getItemDetail("dummy_item_id"); return Uni.createFrom().item(supplier).runSubscriptionOn(executor); }
组合调用逻辑无需修改,保持原有写法即可。
补充说明
- Mutiny两个线程切换操作符的核心差异:
runSubscriptionOn:作用于上游生产阶段,订阅触发时就将上游任务提交到指定执行器,适合包裹阻塞式的同步调用、需要异步执行的supplier逻辑emitOn:作用于下游消费阶段,仅改变后续回调、算子的执行线程,适合切换线程处理下游计算逻辑
- 如果是Quarkus环境运行,无需自定义线程池,可以直接使用内置的worker池:
.runSubscriptionOn(Infrastructure.getDefaultWorkerPool())
内容的提问来源于stack exchange,提问作者Erdem Akkuzu
相关产品推荐
相关产品推荐

