Quarkus Mutiny如何并发调用多个Uni返回最先完成结果
问题根源
你的代码始终返回第一个Uni的结果,和组合API本身无关,核心是你构造的两个Uni从未真正并行执行:
Uni.createFrom().item(Supplier)默认不会异步调度任务,会在订阅发生的当前线程同步执行Supplier里的逻辑。- Mutiny组合多个Uni时,会按传入顺序依次触发订阅。你把阻塞1秒的
uniSlow放在参数第一位,订阅它时Thread.sleep(1000)直接把当前线程完全阻塞,排在后面的uniFast根本没有机会被订阅启动。等1秒后uniSlow执行完成返回结果,first()逻辑已经拿到了第一个可用返回值,自然永远输出one。 - 真实业务场景下如果使用响应式数据库客户端,返回的Uni本身就是异步执行的,不会出现这类阻塞串行的问题;但如果调用的是阻塞式客户端(比如传统JDBC),必须手动将阻塞逻辑调度到独立线程池,才能实现并行调用。
修复方案
给每个需要并行执行的Uni加上runSubscriptionOn,将执行逻辑提交到独立的工作线程池,保证两个任务可以同时启动,再用first()取最先完成的结果即可。
修正后的可运行代码:
import io.smallrye.mutiny.Uni; import org.junit.jupiter.api.Test; import java.time.Duration; import static io.smallrye.mutiny.infrastructure.Infrastructure.getDefaultWorkerPool; public class UniJoinTest { @Test public void testUniJoin(){ // 慢任务,提交到默认工作线程池异步执行 var uniSlow = Uni.createFrom().item(() -> { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return null; } return "one"; }).runSubscriptionOn(getDefaultWorkerPool()); // 快任务,同样提交到工作线程池,与慢任务并行启动 var uniFast = Uni.createFrom().item(() -> { try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); return null; } return "two"; }).runSubscriptionOn(getDefaultWorkerPool()); // 取第一个成功返回的结果 var resp = Uni.join().first(uniSlow, uniFast).withItem().await().atMost(Duration.ofSeconds(2)); System.out.println(resp); // 输出 two var resp2 = Uni.combine().any().of(uniSlow, uniFast).await().atMost(Duration.ofSeconds(2)); System.out.println(resp2); // 输出 two } }
注意事项
- 禁止在响应式调用链中直接使用
Thread.sleep、阻塞IO这类会卡住线程的操作,如果必须调用阻塞逻辑,一定要通过runSubscriptionOn调度到专门的阻塞工作线程池,避免堵死Vert.x事件循环线程。 Uni.join().first()在拿到第一个完成的结果后,会自动取消其他尚未完成的Uni,避免不必要的资源消耗。Uni.join().first(xxx).withItem()和Uni.combine().any().of(xxx)效果完全一致,二者都会返回多个组合Uni中第一个成功完成的结果。- 如果你在Quarkus中使用官方提供的响应式数据库客户端,返回的Uni默认已经做了异步调度,不需要额外加
runSubscriptionOn,直接组合即可拿到最先返回的数据库结果,不会等待慢库响应。
内容的提问来源于stack exchange,提问作者Konstantin Rezchikov
相关产品推荐
相关产品推荐

