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

如何实现Mutiny中Uni任务的并行执行?

解决Mutiny中Uni任务并行运行的问题

你的代码存在两个核心问题导致任务无法并行执行:

1. 阻塞操作同步执行

你的consultrecord方法中,Thread.sleep(5000)是在创建Uni之前直接执行的,这会导致调用consultrecord(r)时,当前线程被立即阻塞,四个任务只能逐个同步完成,完全没有异步并行的机会。

2. 未触发Uni的执行

Mutiny的Uni是惰性求值的对象,仅创建Uni不会触发任何任务执行,必须通过订阅(subscribe())或等待(await())操作来启动任务流。你的代码仅创建了合并后的Uni,但没有触发执行,不过因为前面的同步阻塞,你看到了逐个执行的现象。


修正后的代码

public void initialize(@Observes StartupEvent ev) {
    UniJoin.Builder<Integer> builder = Uni.join().builder();

    Stream.of(1,2,3,4).forEach(r -> {
        builder.add(consultrecord(r));
    });

    Uni<List<Integer>> all = builder.joinAll().usingConcurrencyOf(10).andFailFast();

    // 触发并行任务执行并等待结果(StartupEvent环境允许阻塞等待)
    List<Integer> result = all.await().indefinitely();
    System.out.println("FINISH - 执行结果: " + result);
}

public Uni<Integer> consultrecord(Integer input) {
    // 将阻塞逻辑封装到Uni的异步执行体中
    return Uni.createFrom().item(() -> {
        try {
            System.out.println("PRUEBA DE TEXT - 任务" + input);
            Thread.sleep(5000);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt(); // 保留中断状态
            throw new RuntimeException("任务被中断", e);
        }
        return input;
    }).runSubscriptionOn(Infrastructure.getDefaultWorkerPool()); // 指定工作线程池执行阻塞操作
}

关键修正点说明

  • 异步封装阻塞操作:把Thread.sleep放到Uni.createFrom().item()的Supplier参数中,这样阻塞逻辑只会在Uni被订阅时才执行,且由指定的工作线程池处理,避免阻塞主线程。
  • 触发Uni执行:通过await().indefinitely()触发并行任务执行并等待所有任务完成,确保在StartupEvent初始化阶段完成所有任务。如果不需要同步等待,也可以用subscribe()异步处理结果:
    all.subscribe().with(
        result -> System.out.println("FINISH - 执行结果: " + result),
        error -> System.err.println("任务执行失败: " + error.getMessage())
    );
    

内容的提问来源于stack exchange,提问作者Carlos045

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 05:26:12