Skip to content

CompletableFuture 异步编排

提出问题

假设你正在开发一个订单查询接口,需要同时从三个地方拿数据:用户信息(用户服务)、订单详情(订单服务)、商品快照(商品服务),等三个都拿到后,再组合成一个统一的订单视图返回。

用传统的 Future 实现的话,代码大概长这样:

java
ExecutorService executor = Executors.newFixedThreadPool(3);
Future<User> userFuture = executor.submit(() -> userService.getUser(uid));
Future<Order> orderFuture = executor.submit(() -> orderService.getOrder(orderId));
Future<Product> productFuture = executor.submit(() -> productService.getProduct(pid));

User user = userFuture.get();      // 阻塞
Order order = orderFuture.get();   // 阻塞
Product product = productFuture.get(); // 阻塞

return new OrderVO(user, order, product);

三个 get() 都是阻塞调用。虽然三个任务确实并行执行了,但主线程在等待期间干不了别的。更糟的是,如果你想在拿到用户信息后立即查优惠券(时序依赖),Future 完全无能为力——你必须手动 get() 用户信息后再提交新任务,代码变成一层又一层的回调嵌套。

数据量级:一个典型订单查询接口,用户服务 RPC 耗时 30-80ms,订单服务 20-50ms,商品快照 10-30ms,优惠券查询 15-40ms。如果串行执行,P99 总耗时约 150-200ms;如果并行 + 异步编排,P99 可降至 60-80ms,降低 60% 以上。对于日调用量 100 万次的接口,这相当于每天节省 20-30 小时的 CPU 等待时间。

JDK 8 的 CompletableFuture 就是为了解决这类问题而生的。它让异步任务可以像 Stream API 一样链式编排,用声明式的方式描述"先做 A、A 完了做 B、同时做 C 和 D、等 B 和 C 都完成再合并"这种复杂的异步流程,而不需要显式地管理线程、锁或阻塞等待。

对比:Future vs CompletableFuture

维度FutureCompletableFuture
获取结果阻塞 get(),超时 get(timeout, unit)非阻塞 join(),链式回调,不阻塞调用线程
任务编排不支持,必须手动 get() 后提交新任务thenApply / thenCompose / thenCombine 等声明式编排
异常处理通过 ExecutionException 捕获,异常链不清晰exceptionally / handle / whenComplete,异常沿链传播
多任务组合不支持,需手动 CountDownLatch + get()allOf / anyOf 原生支持
超时控制get(timeout) 抛出 TimeoutExceptionJDK 9+ orTimeout / completeOnTimeout 声明式控制
主动完成不支持,只能等任务结束complete() / completeExceptionally() / obtrudeValue() 可以手动完成
取消语义cancel(true/false),中断 vs 不中断cancel(true) 等价于 completeExceptionally(new CancellationException())
底层实现基于 AbstractQueuedSynchronizer(AQS)基于 volatile result + CAS 无锁栈

分析问题

异步任务创建:supplyAsync vs runAsync

java
// 有返回值
CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    return "Hello";
});

// 无返回值
CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
    System.out.println("Hello");
});

两个方法都有一个带 Executor 参数的重载版本。如果不传 Executor,默认使用 ForkJoinPool.commonPool()

这是一个常见的性能陷阱——commonPool 的线程数等于 Runtime.getRuntime().availableProcessors() - 1。如果异步任务涉及 IO(数据库查询、RPC 调用、HTTP 请求),IO 等待会占用 commonPool 的工作线程,导致其他并行流和 CompletableFuture 任务被阻塞,甚至引发线程饥饿死锁。

生产案例:某团队在 8 核容器上部署订单服务,所有 supplyAsync 都未指定线程池。高峰时 7 个 commonPool 线程全被 RPC 等待占满,导致 parallelStream 的并行流任务无法执行,整条链路超时从 50ms 飙升到 2s。修复方式就是引入自定义 IO 线程池:

java
// 自定义线程池,隔离 IO 和 CPU 型任务
ExecutorService ioPool = Executors.newFixedThreadPool(20);
CompletableFuture.supplyAsync(() -> userService.getUser(uid), ioPool);

线程池大小估算:IO 密集型场景,线程数 = 2 × CPU 核心数 × (1 + 等待时间 / 计算时间)。以订单查询为例,RPC 等待约 60ms,计算约 5ms,比例 12:1,8 核机器建议线程数 = 2 × 8 × (1 + 12) = 208。但实际还需考虑下游服务容量,一般控制在 20-50 之间。

链式回调:thenApply / thenAccept / thenRun

拿到一个异步结果后,最常见的需求是"对它做转换"或"消费它":

java
CompletableFuture.supplyAsync(() -> getUser(uid))
    .thenApply(user -> user.getAddress())       // 转换:User → Address
    .thenAccept(address -> saveToCache(address)) // 消费:Address → void
    .thenRun(() -> log("cache updated"))         // 执行:不依赖前序结果
    .join();  // 等待整个链完成
方法函数签名返回值依赖前序结果
thenApplyFunction<T, R>CompletableFuture<R>
thenAcceptConsumer<T>CompletableFuture<Void>
thenRunRunnableCompletableFuture<Void>

执行线程的陷阱thenApply 默认在完成当前任务的线程上执行回调——可能是 ForkJoinPool 的 worker,也可能是主动调用 complete() 的那个线程。如果你希望回调在另一个线程池中执行,使用 thenApplyAsync 会提交到 ForkJoinPool 重新调度。thenApplyAsync(executor) 则允许指定线程池。

java
// thenApply → 在 supplyAsync 的线程上执行
// thenApplyAsync → 在 ForkJoinPool.commonPool 上执行
// thenApplyAsync(executor) → 在指定线程池上执行

多任务组合

组合是 CompletableFuture 最强大的能力:

thenCompose:扁平化时序组合

java
// 不优雅:CompletableFuture<CompletableFuture<Order>>
CompletableFuture<CompletableFuture<Order>> bad = 
    getUser(uid).thenApply(user -> getOrder(user.getOrderId()));

// 优雅:扁平化
CompletableFuture<Order> good = 
    getUser(uid).thenCompose(user -> getOrder(user.getOrderId()));

thenCompose 相当于 flatMap——它把返回 CompletableFuture 的回调拍平,避免嵌套。底层通过 uniComposeStage 实现,本质是注册一个 UniCompletion 节点在当前 CF 的栈上,前序完成后自动触发。

thenCombine:两个任务并行,结果合并

java
CompletableFuture<String> userFuture = getUser(uid);
CompletableFuture<String> orderFuture = getOrder(orderId);

CompletableFuture<OrderVO> result = userFuture
    .thenCombine(orderFuture, (user, order) -> new OrderVO(user, order));

thenCombine 等待两个 CompletableFuture 都完成,然后用 BiFunction 合并结果。两个任务是并行执行的(假设它们各自使用了不同的线程)。

时序图

时间线 →
用户查询:  ├─────────────────────┤
订单查询:  ├───────────┤
优惠券查询:               ├──────┤
          (用户查完才开始)
合  并:                   └── 三路都完成 → 触发组装

allOf:等待全部完成

java
CompletableFuture<User> userFuture = getUser(uid);
CompletableFuture<Order> orderFuture = getOrder(orderId);
CompletableFuture<Product> productFuture = getProduct(pid);

CompletableFuture<Void> allDone = CompletableFuture.allOf(
    userFuture, orderFuture, productFuture
);

// allOf 返回 Void,需要手动获取结果
allDone.thenRun(() -> {
    User user = userFuture.join();       // 此时不会阻塞,因为已经完成
    Order order = orderFuture.join();
    Product product = productFuture.join();
    return new OrderVO(user, order, product);
});

allOf 的底层实现:内部创建一个 AllOf 内部类实例,构造一个 Cnt(计数器),初始值为传入的 CF 数量。每个 CF 完成后,通过 CAS 递减计数器,当计数器归零时,调用 complete() 完成结果 CF。整个过程无锁,靠 volatile 保证可见性。

join() vs get()join() 抛出非受检的 CompletionExceptionget() 抛出受检的 InterruptedExceptionExecutionException。在链式编程中 join() 更简洁,但如果异常需要传播给调用方处理,get() 的受检异常更合适。

anyOf:任意一个完成

java
// 多个缓存查询,哪个先返回就用哪个
CompletableFuture<Object> first = CompletableFuture.anyOf(
    queryCache("redis"), 
    queryCache("local"), 
    queryCache("memcached")
);

异常处理

CompletableFuture 的异常处理有三种方式,行为差异微妙:

java
// 1. exceptionally:只在异常时执行,提供降级值
CompletableFuture.supplyAsync(() -> riskyCall())
    .exceptionally(e -> {
        log.error("调用失败", e);
        return "fallback";
    });

// 2. handle:无论成功失败都执行,自行决定返回值
CompletableFuture.supplyAsync(() -> riskyCall())
    .handle((result, e) -> {
        if (e != null) {
            log.error("失败", e);
            return "fallback";
        }
        return result;
    });

// 3. whenComplete:无论成功失败都执行,但不改变结果
CompletableFuture.supplyAsync(() -> riskyCall())
    .whenComplete((result, e) -> {
        if (e != null) {
            log.error("失败", e);
        } else {
            log.info("成功: {}", result);
        }
    })
    .join(); // 异常仍然会抛出
方法异常时执行成功时执行是否改变结果
exceptionally是(返回降级值)
handle是(返回任意值)
whenComplete否(结果不变)

生产环境最容易踩的坑是异常静默丢失:如果 thenApply 里抛了异常,后续的 thenApply 不会执行,但异常也不会被打印——除非链上有一个 exceptionallyhandle。调试时可以通过 get() 捕获 ExecutionException 来查看异常,但生产代码中,建议每个异步链的末端都挂一个 exceptionally 兜底

进阶:completeExceptionallyobtrudeValue 的区别

completeExceptionally(Throwable ex) 通过 CAS 将 result 设为 AltResult(ex)只有任务尚未完成时才生效。obtrudeValue(T value)强制设置结果,即使任务已经完成也会覆盖。obtrudeValue 是破坏性的,一般只用于测试或异常恢复场景。

超时控制(JDK 9+)

java
// orTimeout:超时后抛出 TimeoutException
CompletableFuture.supplyAsync(() -> slowCall())
    .orTimeout(3, TimeUnit.SECONDS)
    .exceptionally(e -> "timeout fallback");

// completeOnTimeout:超时后使用默认值完成
CompletableFuture.supplyAsync(() -> slowCall())
    .completeOnTimeout("default", 3, TimeUnit.SECONDS);

JDK 8 没有原生超时方法,需要手动用 get(timeout)FutureTask 包装,JDK 9 的 orTimeoutcompleteOnTimeout 是更优雅的解法。orTimeout 底层使用 Delayer(一个单线程 ScheduledExecutorService)来延迟触发 completeExceptionally(new TimeoutException())

JDK 8 的超时替代方案

java
// JDK 8 兼容方式:手动用 ScheduledExecutorService 实现超时
ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);
CompletableFuture<String> timeoutFuture = future.applyToEither(
    CompletableFuture.supplyAsync(() -> {
        try { Thread.sleep(3000); } catch (Exception e) {}
        return "timeout";
    }, scheduler),
    Function.identity()
);

底层实现:Completion 栈

CompletableFuture 的底层是一个无锁的依赖链机制。每个 CompletableFuture 内部维护一个 volatile Object result(结果值或 AltResult 包装的异常)和一个 Completion 栈(等待当前任务完成后的回调链表)。

关键字段

java
volatile Object result;       // null(未完成) | T(正常结果) | AltResult(异常)
volatile Completion stack;    // 回调链表栈顶(CAS 操作)

完整的执行流程

  1. 创建阶段supplyAsync(runnable) 创建 CompletableFuture 实例,提交任务到 ForkJoinPool.commonPool()(或自定义线程池)
  2. 注册阶段future.thenApply(fn) 创建 UniApply<T,R> 继承自 Completion,通过 CAS 将自身压入前序 CF 的 stack 字段(栈顶)
  3. 完成阶段supplyAsync 的任务执行完毕,调用 complete(T value)
    • 通过 UNSAFE.compareAndSwapObjectresultnull 设为 T
    • 调用 postComplete() 遍历 stack 链表
    • 对每个 Completion 节点调用 tryFire(),返回 true 则继续传播
  4. 传播阶段tryFire() 执行回调函数,更新当前 CF 的 result,然后继续调用 postComplete() 触发后续回调

LIFO 顺序stack 是链表结构,push() 使用 CAS 插入栈顶,因此回调的顺序是后进先出——最后一个注册的 thenApply 最先执行。这是 CompletableFuture 的一个微妙特性。

CAS 无锁的好处:相比 FutureTask 基于 AQS 的阻塞队列,CompletableFuture 的 Completion 栈在单次完成场景下完全不需要锁,避免了线程挂起/唤醒的开销。在高吞吐场景下,这种无锁设计可以显著减少上下文切换。

实战:订单查询接口(完整版)

把开头的问题用 CompletableFuture 重写,加入超时控制、异常兜底、自定义线程池:

java
public OrderVO getOrderDetail(String uid, String orderId, String pid) {
    ExecutorService ioPool = Executors.newFixedThreadPool(10);
    
    CompletableFuture<User> userFuture = CompletableFuture
        .supplyAsync(() -> userService.getUser(uid), ioPool);
    CompletableFuture<Order> orderFuture = CompletableFuture
        .supplyAsync(() -> orderService.getOrder(orderId), ioPool);
    CompletableFuture<Product> productFuture = CompletableFuture
        .supplyAsync(() -> productService.getProduct(pid), ioPool);
    
    // 用户信息查完后,立即查优惠券(时序依赖)
    CompletableFuture<Coupon> couponFuture = userFuture
        .thenApplyAsync(user -> couponService.getCoupon(user.getLevel()), ioPool);
    
    // 等待三个主要数据 + 优惠券全部就绪,组装结果
    CompletableFuture<OrderVO> result = CompletableFuture
        .allOf(userFuture, orderFuture, productFuture, couponFuture)
        .thenApplyAsync(v -> {
            User user = userFuture.join();
            Order order = orderFuture.join();
            Product product = productFuture.join();
            Coupon coupon = couponFuture.join();
            return new OrderVO(user, order, product, coupon);
        }, ioPool)
        .exceptionally(e -> {
            log.error("组装订单视图失败", e);
            return OrderVO.empty();  // 降级返回空对象
        })
        .orTimeout(5, TimeUnit.SECONDS)
        .exceptionally(e -> {
            log.warn("订单查询超时,uid={}, orderId={}", uid, orderId);
            return OrderVO.empty();
        });
    
    return result.join();
}

这个设计的优势

  • 三个无依赖的查询并行执行,总耗时 = max(用户服务, 订单服务, 商品服务) ≈ 60ms
  • 用户信息一拿到,自动触发优惠券查询(不阻塞,时序依赖自动处理)
  • 所有数据就绪后自动组装(不主动轮询)
  • 异常和超时都有兜底,不会让调用方空等
  • ioPool 隔离了 IO 线程,不影响 commonPool 的并行流计算

更现实的线程池隔离策略

java
// 按下游服务隔离线程池,避免一个服务雪崩拖垮其他服务
ExecutorService userPool = Executors.newFixedThreadPool(5, r -> new Thread(r, "user-rpc-"));
ExecutorService orderPool = Executors.newFixedThreadPool(5, r -> new Thread(r, "order-rpc-"));
ExecutorService productPool = Executors.newFixedThreadPool(3, r -> new Thread(r, "product-rpc-"));

高级用法

生产级重试模式

java
public <T> CompletableFuture<T> withRetry(
        Supplier<T> action, int maxRetries, Executor executor) {
    CompletableFuture<T> cf = CompletableFuture.supplyAsync(action, executor);
    for (int i = 0; i < maxRetries; i++) {
        cf = cf.exceptionally(e -> {
            log.warn("重试 {}/{}", i + 1, maxRetries, e);
            return CompletableFuture.supplyAsync(action, executor).join();
        });
    }
    return cf;
}

熔断式超时(JDK 9+)

java
// 先快速返回缓存,后台异步刷新
CompletableFuture<String> cf = CompletableFuture.supplyAsync(() -> fetchFromDb())
    .completeOnTimeout(cache.get(key), 100, TimeUnit.MILLISECONDS)
    .thenApply(result -> {
        cache.put(key, result); // 后台刷新缓存
        return result;
    });

总结

关键点清单

  • supplyAsync/runAsync 创建异步任务,务必指定自定义线程池处理 IO 型任务,避免公共池被阻塞
  • thenApply(转换)、thenAccept(消费)、thenRun(执行)构建链式回调;Async 后缀的版本切换到新线程执行
  • thenCompose 扁平化时序组合,避免 CompletableFuture<CompletableFuture> 嵌套
  • thenCombine 并行执行两个任务后合并结果;allOf 等待全部完成;anyOf 等待任意一个完成
  • 异常处理三件套:exceptionally(异常降级)、handle(成功失败都处理)、whenComplete(感知但不改变结果)
  • 异步链末端必须挂异常处理,否则异常会静默丢失
  • JDK 9+ 的 orTimeoutcompleteOnTimeout 提供原生超时控制
  • 底层基于 volatile result + CAS 无锁传播的 Completion

面试话术示例

"CompletableFuture 的核心价值是把异步编排从回调地狱变成声明式链式调用。底层的 Completion 栈是无锁的,通过 CAS 更新 volatile result 触发后续回调传播。生产中最关键的是两件事:一是 IO 密集型任务必须用自定义线程池,别用 commonPool 把所有并行流都堵住——我见过 8 核机器上 commonPool 7 个线程全被 RPC 等待占满导致整条链路超时翻 40 倍的案例;二是每个异步链末端必须放一个 exceptionally 兜底,否则异常会丢得无声无息。JDK 9 的 orTimeout 是超时控制的最佳实践,JDK 8 只能靠 get(timeout)ScheduledExecutorService 手动模拟。allOf 底层通过 CAS 计数器实现无锁完成检测,join()get() 在链式编程中更简洁因为抛的是非受检异常。"

参考:CompletableFuture 源码 java.util.concurrent.CompletableFutureCompletionStage 接口文档、ForkJoinPool.commonPool() 源码、JDK 9 JEP 266(CompletableFuture 增强)

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。