如何用CompletableFuture分离处理Either类型集合中的异常与结果?
实现方案(Java 11)
你的需求完全可以实现,以下是具体的代码示例和关键说明:
1. 修正转换方法签名
首先注意你的convertBusinessObjectJson方法签名存在矛盾(注释说明返回List<Either>,但代码写的返回List<String>),这里假设使用Vavr的Either(Java标准库无内置Either,若为自定义类型可直接替换),修正后的方法如下:
import io.vavr.control.Either; import java.util.ArrayList; import java.util.List; public class FutureImpl { public static List<Either<RuntimeException, String>> convertBusinessObjectJson(List<BusinessObject> businessObjList) { List<Either<RuntimeException, String>> eitherList = new ArrayList<>(); // 替换为你的实际转换逻辑: // 遍历businessObjList,转换成功则添加Either.right(jsonStr), // 捕获异常则添加Either.left(exception) for (BusinessObject obj : businessObjList) { try { String json = convertToJson(obj); // 你的转换逻辑 eitherList.add(Either.right(json)); } catch (RuntimeException e) { eitherList.add(Either.left(e)); } } return eitherList; } // 示例转换方法(需替换为你的实际逻辑) private static String convertToJson(BusinessObject obj) { return "{\"id\": " + obj.getId() + "}"; } }
2. 异步处理与结果拆分
在调用代码中,我们通过CompletableFuture.supplyAsync发起异步转换,然后用thenAccept消费转换后的Either列表,拆分异常和成功结果并分别执行对应操作:
import io.vavr.control.Either; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.stream.Collectors; public class AsyncProcessor { public static void processMainList(List<List<BusinessObject>> mainList) { // 收集所有异步任务的Future,便于后续等待全部完成 List<CompletableFuture<Void>> futures = new ArrayList<>(); mainList.forEach(sublist -> { // 异步执行转换 CompletableFuture<List<Either<RuntimeException, String>>> convertFuture = CompletableFuture.supplyAsync(() -> FutureImpl.convertBusinessObjectJson(sublist)); // 处理转换结果:拆分异常和成功数据,执行对应操作 CompletableFuture<Void> processingFuture = convertFuture.thenAccept(eitherList -> { // 提取所有异常 List<RuntimeException> exceptions = eitherList.stream() .filter(Either::isLeft) .map(Either::getLeft) .collect(Collectors.toList()); // 提取所有成功转换的JSON结果 List<String> successResults = eitherList.stream() .filter(Either::isRight) .map(Either::getRight) .collect(Collectors.toList()); // 推送异常到Kafka(非空时执行) if (!exceptions.isEmpty()) { sendExceptionsToKafka(exceptions); } // 写入成功结果到数据库(非空时执行) if (!successResults.isEmpty()) { writeResultsToDb(successResults); } }); futures.add(processingFuture); }); // 可选:等待所有异步任务完成,避免主线程提前退出 CompletableFuture.allOf(futures.toArray(new CompletableFuture[0])).join(); } // 替换为你的Kafka推送实现 private static void sendExceptionsToKafka(List<RuntimeException> exceptions) { // 示例逻辑 exceptions.forEach(ex -> System.out.println("推送异常到Kafka: " + ex.getMessage())); } // 替换为你的数据库写入实现 private static void writeResultsToDb(List<String> results) { // 示例逻辑 results.forEach(json -> System.out.println("写入DB: " + json)); } }
3. 关键注意事项
- 使用
thenAccept而非thenApply:thenApply用于返回新的结果,而我们的操作是消费型的(无返回值),thenAccept更贴合场景。 - 处理Future自身的异常:如果
convertBusinessObjectJson方法本身抛出未被捕获的异常(未封装到Either中),可以通过exceptionally补充处理:CompletableFuture<Void> processingFuture = convertFuture .exceptionally(ex -> { // 处理转换方法本身抛出的未封装异常 sendExceptionsToKafka(List.of((RuntimeException) ex)); return List.of(); // 返回空列表继续后续流程 }) .thenAccept(eitherList -> { /* 原有处理逻辑 */ }); - 异步操作的组合:如果推送Kafka或写入DB的操作本身是异步的(返回
CompletableFuture),可以用thenCompose和CompletableFuture.allOf组合执行:CompletableFuture<Void> processingFuture = convertFuture.thenCompose(eitherList -> { List<RuntimeException> exceptions = /* 提取逻辑 */; List<String> successResults = /* 提取逻辑 */; CompletableFuture<Void> kafkaFuture = exceptions.isEmpty() ? CompletableFuture.completedFuture(null) : sendExceptionsToKafkaAsync(exceptions); CompletableFuture<Void> dbFuture = successResults.isEmpty() ? CompletableFuture.completedFuture(null) : writeResultsToDbAsync(successResults); // 并行执行两个异步任务,等待全部完成 return CompletableFuture.allOf(kafkaFuture, dbFuture); });
内容的提问来源于stack exchange,提问作者Nishikant Tayade
相关产品推荐
相关产品推荐

