如何实现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
相关产品推荐
相关产品推荐

