
当面对一个Uni<List<T>>并希望对列表中的每个T执行异步操作时,一个常见的误区是尝试直接通过map将List<T>转换为List<Uni<R>>,然后使用Uni.join().all(unis).andCollectFailures()来合并结果,最后通过subscribe()进行消费。这种方法在Mutiny的链式操作中是可行的,但如果后续没有适当的机制来保持主线程的活跃,例如在简单的main方法中,程序可能会在所有异步任务完成之前退出,导致部分或全部异步操作未能执行或其结果未被观察到。
要正确地将Uni<List<T>>中的每个元素转换为一个独立的异步Uni<R>并进行并发处理,我们需要利用Mutiny的流式处理能力,或者采用阻塞机制来等待所有操作完成。
这种方法是Mutiny推荐的、更符合响应式编程范式的处理方式。它通过将包含列表的Uni转换为一个Multi流,然后对流中的每个元素进行异步转换和合并,实现并发处理。
Mutiny的Multi类型非常适合处理元素流。通过以下步骤,我们可以将Uni<List<T>>转换为Multi<T>,对流中的每个元素独立应用异步转换,并利用transformToUniAndMerge实现并发处理:
以下示例演示了如何使用线程池模拟异步操作,并结合Mutiny的Multi进行非阻塞流式处理。
import io.smallrye.mutiny.Multi;
import io.smallrye.mutiny.Uni;
import java.time.Duration;
import java.util.List;
import java.util.Random;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.CountDownLatch; // 用于在main方法中等待所有异步任务完成
public class AsyncListProcessor {
private final ExecutorService executor = Executors.newFixedThreadPool(3); // 示例用线程池,限制并发数
// 模拟一个返回Future的耗时操作
private Future<String> processFuture(String s) {
return executor.submit(() -> {
System.out.println("开始处理 (Future): " + s + " on thread " + Thread.currentThread().getName());
try {
Thread.sleep(new Random().nextInt(3000) + 1000); // 模拟1-4秒延迟
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("处理中断", e);
}
System.out.println("结束处理 (Future): " + s + " on thread " + Thread.currentThread().getName());
return s.toUpperCase(); // 假设处理后返回大写
});
}
// 将Future封装成Uni
private Uni<String> processItemAsUni(String item) {
return Uni.createFrom().future(processFuture(item));
}
public void processListReactive(List<String> items, CountDownLatch latch) {
System.out.println("\n--- 启动非阻塞流式处理 ---");
Uni.createFrom()
.item(items)
// 将 Uni<List<String>> 转换为 Multi<String>
.onItem().transformToMulti(Multi.createFrom()::iterable)
// 对 Multi 中的每个元素进行异步处理,并合并结果
.onItem().transformToUniAndMerge(this::processItemAsUni)
// 订阅并打印每个完成的结果
.subscribe()
.with(
s -> System.out.println("接收到结果 (Reactive): " + s),
failure -> System.err.println("处理失败: " + failure.getMessage()),
() -> {
System.out.println("所有非阻塞流式处理完成.");
latch.countDown(); // 通知主线程所有任务已完成
}
);
}
public static void main(String[] args) throws InterruptedException {
AsyncListProcessor processor = new AsyncListProcessor();
List<String> data = List.of("apple", "banana", "cherry", "date", "elderberry");
// 使用CountDownLatch等待所有异步任务完成
CountDownLatch latch = new CountDownLatch(1);
processor.processListReactive(data, latch);
// 等待所有异步任务完成
latch.await();
System.out.println("主线程继续执行,所有异步任务已完成或失败。");
processor.executor.shutdown(); // 关闭线程池
}
}这种方式是非阻塞的,非常适合构建响应式应用程序。它允许任务并发执行,且不会阻塞主调用线程。在非Web服务器等环境中(如简单的main方法),为了确保程序不会在异步操作完成前退出,需要使用CountDownLatch、await()或其他同步机制来等待所有任务完成。在基于Mutiny的框架(如Quarkus)中,这些通常由框架的调度器和生命周期管理。
如果你的需求是等待所有异步操作完成,并将它们的结果收集到一个列表中,然后才能继续执行后续逻辑,那么可以使用阻塞式的方法。
这种方法同样利用Multi进行元素的异步转换,但在最后阶段,它会阻塞当前线程,直到所有异步操作完成并将结果聚合到一个列表中。
import io.smallrye.mutiny.Multi;
import io.smallrye.mutiny.Uni;
import java.time.Duration;
import java.util.List;
import java.util.Random;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
public class AsyncListProcessorBlocking {
private final ExecutorService executor = Executors.newFixedThreadPool(3); // 示例用线程池
// 模拟一个返回Future的耗时操作 (同上)
private Future<String> processFuture(String s) {
return executor.submit(() -> {
System.out.println("开始处理 (Future): " + s + " on thread " + Thread.currentThread().getName());
try {
Thread.sleep(new Random().nextInt(3000) + 1000); // 模拟1-4秒延迟
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("处理中断", e);
}
System.out.println("结束处理 (Future): " + s + " on thread " + Thread.currentThread().getName());
return s.toUpperCase();
});
}
// 将Future封装成Uni
private Uni<String> processItemAsUni(String item) {
return Uni.createFrom().future(processFuture(item));
}
public List<String> processListBlocking(List<String> items) {
System.out.println("\n--- 启动阻塞式处理 ---");
List<String> results = Uni.createFrom()
.item(items)
.onItem().transformToMulti(Multi.createFrom()::iterable)
.onItem().transformToUniAndMerge(this::processItemAsUni)
.collect().asList() // 收集所有结果到一个 Uni<List<String>>
.await().indefinitely(); // 阻塞当前线程直到所有结果收集完毕
System.out.println("--- 阻塞式处理完成 ---");
return results;
}
public static void main(String[] args) {
AsyncListProcessorBlocking processor = new AsyncListProcessorBlocking();
List<String> data = List.of("alpha", "beta", "gamma", "delta", "epsilon");
List<String> processedResults = processor.processListBlocking(data);
System.out.println("所有处理结果 (阻塞式): " + processedResults);
processor.executor.shutdown(); // 关闭线程池
}
}await().indefinitely()会阻塞调用线程。虽然它能确保所有异步操作完成,但在响应式系统中应谨慎使用,因为它可能导致线程阻塞,降低系统的并发能力。它更适用于启动时的数据加载、测试场景或需要等待所有结果才能继续的特定批处理任务。在Web应用中,应避免在处理请求的线程中使用await(),以防阻塞请求处理。
Mutiny提供了灵活且强大的机制来处理Uni<List<T>>中的元素异步操作。选择哪种
以上就是Mutiny异步处理Uni中元素的最佳实践的详细内容,更多请关注php中文网其它相关文章!
每个人都需要一台速度更快、更稳定的 PC。随着时间的推移,垃圾文件、旧注册表数据和不必要的后台进程会占用资源并降低性能。幸运的是,许多工具可以让 Windows 保持平稳运行。
Copyright 2014-2025 https://www.php.cn/ All Rights Reserved | php.cn | 湘ICP备2023035733号