CompletableFuture组合多个任务

wen java案例 2

CompletableFuture组合多个任务的深度实践

目录导读

  1. 为何需要任务组合?
  2. CompletableFuture核心机制速览
  3. 组合任务的四大经典模式
    • 1 顺序组合(thenCompose)
    • 2 并行合并(thenCombine)
    • 3 任意完成(applyToEither)
    • 4 多任务聚合(allOf / anyOf)
  4. 常见问题与高性能技巧
  5. 实战案例:异步订单处理流水线
  6. QA环节:开发者最关心的5个问题

为何需要任务组合?

在微服务架构与高并发场景下,单一异步任务往往无法满足复杂业务需求。

CompletableFuture组合多个任务

  • 先获取用户信息,再根据用户ID查询订单列表
  • 同时调用支付网关、库存服务、物流服务,最后聚合结果
  • 多个外部API中任意一个成功即返回结果(如熔断降级)

Java 8引入的CompletableFuture,正是为了解决这类多阶段异步编排问题,它比传统Future更强大:支持回调、组合、异常处理,且能避免回调地狱。


CompletableFuture核心机制速览

方法 用途 类比
supplyAsync() 创建异步任务 线程池提交
thenApply() 转换结果(同步) 类似map
thenCompose() 扁平化任务组合 类似flatMap
thenCombine() 合并两个并行结果 类似zip
applyToEither() 取最快结果 竞速模式
allOf() 等待所有完成 栅栏
anyOf() 任意一个完成 竞速

组合任务的四大经典模式

1 顺序依赖:thenCompose(异步衔接)

当第二个任务依赖第一个任务的结果时,使用thenCompose避免嵌套:

CompletableFuture<String> userFuture = getUserAsync(1);
CompletableFuture<String> orderFuture = userFuture.thenCompose(user -> 
    getOrdersByUserAsync(user.getId()));

注意:误用thenApply会导致CompletableFuture<CompletableFuture<T>>嵌套。

2 并行合并:thenCombine(汇聚两个结果)

同时调用两个无关接口,并合并结果:

CompletableFuture<Double> priceFuture = getPriceAsync("itemA");
CompletableFuture<Double> discountFuture = getDiscountAsync("VIP");
priceFuture.thenCombine(discountFuture, (price, discount) -> price * discount)
           .thenAccept(System.out::println);

3 竞速模式:applyToEither(谁快用谁)

适用于缓存预热、多节点探测:

CompletableFuture<String> fastFuture = fastService.call();
CompletableFuture<String> slowFuture = slowService.call();
fastFuture.applyToEither(slowFuture, result -> "最快响应: " + result);

4 多任务聚合:allOf & anyOf

  • allOf:等待所有任务完成,适合批量数据加载
  • anyOf:任意一个成功即返回,适合降级兜底
CompletableFuture.allOf(task1, task2, task3)
    .thenRun(() -> System.out.println("所有任务完成"));

常见问题与高性能技巧

1 避免默认线程池陷阱

// 不推荐:使用默认ForkJoinPool(可能被阻塞)
CompletableFuture.supplyAsync(this::heavyTask);
// 推荐:显式指定线程池
ExecutorService executor = Executors.newFixedThreadPool(10);
CompletableFuture.supplyAsync(this::heavyTask, executor);

2 异常处理链式传递

// 方式一:exceptionally 提供降级值
CompletableFuture<Integer> safe = future.exceptionally(ex -> 0);
// 方式二:handle 统一处理结果与异常
future.handle((result, ex) -> ex == null ? result : -1);

3 超时控制(Java 9+)

future.orTimeout(3, TimeUnit.SECONDS)
      .exceptionally(ex -> "timeout");

实战案例:异步订单处理流水线

需求:用户下单后,并行校验库存、计算价格、查询历史,然后顺序执行扣库存、生成订单。

ExecutorService pool = Executors.newFixedThreadPool(20);
// 阶段1: 并行获取初始数据
CompletableFuture<Stock> stockFuture = CompletableFuture.supplyAsync(() -> checkStock(itemId), pool);
CompletableFuture<Price> priceFuture = CompletableFuture.supplyAsync(() -> computePrice(order), pool);
CompletableFuture<History> historyFuture = CompletableFuture.supplyAsync(() -> getHistory(userId), pool);
// 阶段2: 合并三个结果,并顺序执行后续步骤
CompletableFuture<Void> pipeline = CompletableFuture.allOf(stockFuture, priceFuture, historyFuture)
    .thenCompose(v -> {
        Stock s = stockFuture.join();
        Price p = priceFuture.join();
        History h = historyFuture.join();
        // 接着执行依赖上述结果的扣库存与生成订单
        return deductStock(s).thenCompose(ok -> createOrder(p, h));
    })
    .exceptionally(ex -> {
        log.error("订单处理失败", ex);
        return null;
    });

QA环节:开发者最关心的5个问题

Q1:thenComposethenApply 到底有什么区别?

  • thenApply 对结果直接做同步转换,返回 CompletableFuture<U>
  • thenCompose 对结果做异步转换,返回 U(内部是 CompletableFuture<U>
    选型规则:如果函数返回CompletableFuture,用thenCompose;否则用thenApply

Q2:多个任务如何优雅地收集所有结果(含异常)?
使用allOf + 自定义CompletableFuture,遍历时处理异常:

List<CompletableFuture<String>> futures = tasks.stream()
    .map(task -> CompletableFuture.supplyAsync(task, pool)
        .exceptionally(ex -> "ERROR: " + ex.getMessage()))
    .collect(Collectors.toList());
CompletableFuture<Void> all = CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]));
all.join(); // 等待所有(异常值已降级)
List<String> results = futures.stream().map(CompletableFuture::join).collect(Collectors.toList());

Q3:join()get() 有什么区别,该如何选择?

  • get() 抛出检查型异常(InterruptedException, ExecutionException)
  • join() 抛出非检查异常(CompletionException)
    建议:在流式操作或lambda中优先使用join()(代码简洁),在需要精确异常处理时用get()

Q4:怎么实现“所有任务完成或超时”的逻辑?
利用allOf配合orTimeout(Java 9+):

CompletableFuture<Void> all = CompletableFuture.allOf(task1, task2);
all.orTimeout(2, TimeUnit.SECONDS)
   .exceptionally(ex -> { 
       System.out.println("部分任务超时"); 
       return null; 
   });

对于Java 8,可使用Future.get(timeout, unit)CompletableFuture.delayer

Q5:在Spring Boot中如何使用CompletableFuture事务?
事务(@Transactional)默认不能跨线程传播,推荐做法:

  • 在调用supplyAsync之前获取数据
  • 或将事务逻辑放在thenCompose的回调线程中(需确保线程池使用同一个事务管理器)
  • 更稳妥:使用TransactionTemplate手动管理事务边界。

延伸阅读

  • 官方文档:CompletableFuture JavaDoc
  • 性能对比:自定义线程池通常比ForkJoinPool更适合I/O密集型任务,避免公共池被阻塞。

抱歉,评论功能暂时关闭!