Skip to content

CompletableFuture 全过程原理与商业异步编排

CompletableFuture 是 JDK 8 引入的异步编排工具。很多人会用 supplyAsyncthenApplyallOf,但真正上项目时更重要的是说清:

  1. 它和 Future 的本质区别是什么。
  2. 异步任务到底由哪个线程池执行。
  3. thenApplythenApplyAsync 的线程有什么区别。
  4. 多个任务怎么组合,依赖任务怎么串联。
  5. 异常怎么传播,为什么 allOf 后还要逐个取结果。
  6. 超时是不是会真正取消底层任务。
  7. 为什么不能把阻塞 IO 都丢到公共 ForkJoinPool
  8. 接口聚合、批量采集、异步入库这些商业场景怎么用才安全。

一句话先建立主线:

CompletableFuture 不是新线程,也不是性能魔法。它是“异步任务结果 + 依赖回调链”的编排对象,真正执行任务的仍然是线程池。

学习目标

学完这一页,你要能回答:

  1. Future 为什么不适合复杂异步编排。
  2. CompletableFuture 内部保存哪些状态。
  3. supplyAsyncrunAsyncthenApplythenComposethenCombineallOf 分别解决什么问题。
  4. 非 Async 回调和 Async 回调由哪个线程执行。
  5. join()get() 有什么区别。
  6. exceptionallyhandlewhenComplete 有什么区别。
  7. orTimeoutcompleteOnTimeout 的真实含义。
  8. 如何写一个可落地的订单详情接口聚合 Demo。
  9. 线上异步任务堆积、线程池打满、下游慢、异常被吞怎么排查。

Future 的痛点

JDK 5 就有 Future,但它更像“任务回执单”。

java
Future<String> userFuture = executor.submit(() -> queryUser());
String user = userFuture.get();

它能拿结果,但有几个问题:

痛点说明后果
get() 阻塞调用方线程要等结果容易把线程池卡住
不擅长组合多个 Future 要手动逐个 get代码复杂,异常难处理
不擅长依赖任务 B 依赖任务 A 的结果要手写回调容易写成阻塞等待
不擅长异常链异常处理分散在多个 get 周围兜底逻辑混乱
不能主动完成只能由任务执行完成不方便超时兜底和外部补偿

CompletableFuture 的核心增强是:它既表示未来结果,又能描述后续依赖动作。

mermaid
flowchart TD
    A["Future<br/>只能等结果"] --> B["get 阻塞"]
    C["CompletableFuture"] --> D["保存结果或异常"]
    C --> E["注册后续动作"]
    C --> F["组合多个任务"]
    C --> G["主动 complete"]

CompletableFuture 内部可以怎么理解

从学习角度,可以把一个 CompletableFuture 理解成四部分:

mermaid
flowchart TD
    A["CompletableFuture"] --> B["任务状态<br/>未完成/成功/失败/取消"]
    A --> C["结果或异常<br/>value / Throwable"]
    A --> D["依赖动作栈<br/>thenApply / thenAccept 等"]
    A --> E["执行器<br/>默认公共池或自定义线程池"]

任务完成时,它会:

  1. 把状态改成成功或失败。
  2. 保存结果或异常。
  3. 触发依赖它的后续动作。
  4. 把结果或异常继续向后传播。

流程:

mermaid
flowchart TD
    A["创建 CompletableFuture"] --> B["注册依赖动作"]
    B --> C["异步任务在线程池执行"]
    C --> D{"任务成功还是失败"}
    D -- "成功" --> E["保存正常结果"]
    D -- "失败" --> F["保存异常结果"]
    E --> G["触发后续阶段"]
    F --> H["触发异常处理阶段"]
    G --> I["继续向后传播"]
    H --> I

supplyAsync 和 runAsync

方法是否有返回值适合场景
supplyAsync查询用户、订单、库存并返回结果
runAsync异步写日志、发通知、刷新缓存

Demo:

java
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

public class SupplyRunDemo {
    private static final ExecutorService POOL = Executors.newFixedThreadPool(4);

    public static void main(String[] args) {
        CompletableFuture<String> userFuture =
                CompletableFuture.supplyAsync(() -> "user-1001", POOL);

        CompletableFuture<Void> logFuture =
                CompletableFuture.runAsync(() -> System.out.println("write log"), POOL);

        System.out.println(userFuture.join());
        logFuture.join();
        POOL.shutdown();
    }
}

生产里不要直接使用 Executors.newFixedThreadPool 作为最终方案,这里只是最小 Demo。真实项目要使用有界队列、自定义线程名、拒绝策略和监控。

非 Async 和 Async 的线程区别

这是 CompletableFuture 面试和线上排查最重要的点之一。

写法后续阶段由谁执行
thenApply通常由完成上一步的线程继续执行
thenApplyAsync 不传线程池默认提交到公共 ForkJoinPool
thenApplyAsync 传线程池提交到指定线程池

流程:

mermaid
flowchart TD
    A["任务 A 在线程 pool-1 执行"] --> B["任务 A 完成"]
    B --> C{"后续是 thenApply 还是 thenApplyAsync"}
    C -- "thenApply" --> D["可能继续由 pool-1 执行"]
    C -- "thenApplyAsync 无 executor" --> E["提交到公共 ForkJoinPool"]
    C -- "thenApplyAsync 有 executor" --> F["提交到指定线程池"]

为什么这很重要?

  1. 如果 thenApply 里做慢 IO,可能拖慢完成上一步的线程。
  2. 如果 thenApplyAsync 不指定线程池,阻塞任务会占用公共池。
  3. 公共池被占满会影响其他并行 Stream、其他 CompletableFuture 任务。

建议:

场景建议
轻量内存转换thenApply
下游 IO 调用thenComposeAsyncthenApplyAsync + 自定义 IO 线程池
CPU 计算独立 CPU 线程池
核心接口聚合必须自定义线程池和超时

thenApply、thenCompose、thenCombine 怎么选

thenApply:同步转换结果

上一步返回用户 ID,下一步只是拼接字符串:

java
CompletableFuture<String> future = CompletableFuture
        .supplyAsync(() -> 1001L)
        .thenApply(userId -> "user-" + userId);

thenApply 的返回值不是新的异步任务,而是普通转换结果。

thenCompose:串联依赖的异步任务

用户详情依赖用户 ID,下一步本身也是异步查询:

java
CompletableFuture<String> userDetail = queryUserId()
        .thenCompose(userId -> queryUserDetail(userId));

如果用 thenApply,会得到嵌套:

java
CompletableFuture<CompletableFuture<String>> wrong =
        queryUserId().thenApply(userId -> queryUserDetail(userId));

所以:

场景应该用
A 结果转成普通 BthenApply
A 结果决定下一个异步 FuturethenCompose

thenCombine:合并两个独立任务

订单和优惠券互不依赖,可以并行后合并:

java
CompletableFuture<String> merged = queryOrder()
        .thenCombine(queryCoupon(), (order, coupon) -> order + " + " + coupon);

流程图:

mermaid
flowchart TD
    A["查询订单"] --> C["thenCombine 合并"]
    B["查询优惠券"] --> C
    C --> D["生成订单详情片段"]

allOf 的真实含义

allOf 表示“等这些 Future 都完成”,但它不帮你把每个结果收集成列表。

java
CompletableFuture<Void> all =
        CompletableFuture.allOf(userFuture, orderFuture, couponFuture);
all.join();

为什么 allOf 不直接返回 List<T>

因为每个 Future 的类型可能不同:

java
CompletableFuture<UserDTO> userFuture;
CompletableFuture<OrderDTO> orderFuture;
CompletableFuture<List<CouponDTO>> couponFuture;

它只能统一表示“都完成了”,具体结果还要从每个 Future 中取。

allOf 遇到异常

只要其中一个 Future 异常,allOf().join() 也会异常。但其他任务可能已经完成,也可能还在继续执行。

mermaid
flowchart TD
    A["userFuture 成功"] --> D["allOf"]
    B["orderFuture 异常"] --> D
    C["couponFuture 仍执行"] --> D
    D --> E["allOf 以异常完成"]

所以生产代码要考虑:

  1. 哪些任务失败可以兜底。
  2. 哪些任务失败必须整体失败。
  3. 是否要记录每个子任务的异常。
  4. 是否要对慢任务设置超时。

异常处理方法区别

方法成功时调用失败时调用能否改变结果
exceptionally能,返回兜底值
handle能,返回新结果
whenComplete通常只观察,不改变上游结果

Demo:

java
CompletableFuture<String> future = CompletableFuture
        .supplyAsync(() -> {
            throw new RuntimeException("query failed");
        })
        .exceptionally(ex -> "default-user");

System.out.println(future.join());

whenComplete 常用于日志:

java
CompletableFuture<String> future = queryUser()
        .whenComplete((result, error) -> {
            if (error != null) {
                System.out.println("query user failed: " + error.getMessage());
            }
        });

注意:whenComplete 里如果自己又抛异常,可能覆盖或影响后续异常表现。日志逻辑要稳,不能在异常处理里再制造异常。

join 和 get 区别

对比join()get()
异常类型CompletionExceptionExecutionExceptionInterruptedException
是否受检异常不是
中断处理不强制你处理必须处理 InterruptedException
常见使用函数式链路里更简洁传统 Future 风格

生产代码里,如果用 get() 捕获 InterruptedException,要恢复中断标记:

java
try {
    return future.get();
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
    throw new RuntimeException("interrupted", e);
}

超时控制:不是所有超时都会杀掉底层任务

JDK 9 之后有:

方法行为
orTimeout超时后让 Future 以异常完成
completeOnTimeout超时后用默认值完成

示例:

java
CompletableFuture<String> future = CompletableFuture
        .supplyAsync(() -> slowQuery(), IO_POOL)
        .completeOnTimeout("default", 800, TimeUnit.MILLISECONDS);

重点:

Future 超时完成,不一定代表底层正在执行的任务被强制停止。

如果 slowQuery() 已经在线程中执行,并且底层 HTTP/DB 调用没有超时,它可能还会继续占用线程。真正的生产超时要两层都做:

  1. CompletableFuture 编排层超时,避免接口一直等。
  2. HTTP、DB、Redis、RPC 客户端自身超时,避免工作线程长期卡住。

商业 Demo:订单详情接口聚合

这个 Demo 模拟订单详情页并行查询用户、订单、优惠券,并处理超时和异常。

java
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;

public class OrderDetailCompletableFutureDemo {
    private static final ThreadPoolExecutor IO_POOL = new ThreadPoolExecutor(
            8,
            16,
            60,
            TimeUnit.SECONDS,
            new ArrayBlockingQueue<Runnable>(200),
            r -> new Thread(r, "order-aggregate-" + System.nanoTime()),
            new ThreadPoolExecutor.CallerRunsPolicy()
    );

    public static void main(String[] args) {
        OrderDetail detail = queryOrderDetail(1001L);
        System.out.println(detail);
        IO_POOL.shutdown();
    }

    private static OrderDetail queryOrderDetail(Long orderId) {
        CompletableFuture<String> userFuture = CompletableFuture
                .supplyAsync(() -> queryUser(orderId), IO_POOL)
                .completeOnTimeout("unknown-user", 800, TimeUnit.MILLISECONDS)
                .exceptionally(ex -> "unknown-user");

        CompletableFuture<String> orderFuture = CompletableFuture
                .supplyAsync(() -> queryOrder(orderId), IO_POOL)
                .orTimeout(1, TimeUnit.SECONDS);

        CompletableFuture<String> couponFuture = CompletableFuture
                .supplyAsync(() -> queryCoupon(orderId), IO_POOL)
                .completeOnTimeout("no-coupon", 500, TimeUnit.MILLISECONDS)
                .exceptionally(ex -> "no-coupon");

        CompletableFuture<Void> all = CompletableFuture.allOf(userFuture, orderFuture, couponFuture);

        try {
            all.join();
            return new OrderDetail(
                    userFuture.join(),
                    orderFuture.join(),
                    couponFuture.join()
            );
        } catch (Exception e) {
            throw new RuntimeException("query order detail failed, orderId=" + orderId, e);
        }
    }

    private static String queryUser(Long orderId) {
        sleep(100);
        return "user-of-" + orderId;
    }

    private static String queryOrder(Long orderId) {
        sleep(120);
        return "order-" + orderId;
    }

    private static String queryCoupon(Long orderId) {
        sleep(80);
        return "coupon";
    }

    private static void sleep(long millis) {
        try {
            TimeUnit.MILLISECONDS.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException(e);
        }
    }

    static class OrderDetail {
        private final String user;
        private final String order;
        private final String coupon;

        OrderDetail(String user, String order, String coupon) {
            this.user = user;
            this.order = order;
            this.coupon = coupon;
        }

        public String toString() {
            return "OrderDetail{user='" + user + "', order='" + order + "', coupon='" + coupon + "'}";
        }
    }
}

这个 Demo 体现的生产原则:

原则说明
自定义线程池避免阻塞公共池
有界队列防止无限堆积
线程命名方便 jstack 排查
局部兜底非核心数据可返回默认值
核心失败抛异常订单主信息失败不能假装成功
编排层超时防止接口无限等待

真实项目还要加 traceId、日志、指标、下游客户端超时、熔断限流和降级策略。

商业场景怎么选

场景是否适合 CompletableFuture注意点
订单详情接口聚合适合下游互不依赖,并且有超时兜底
批量采集多个医院接口适合但要谨慎控制并发,不能打爆目标系统
支付扣款主流程谨慎不能随便异步化,要保证状态一致
异步写日志/埋点适合可降级,不能影响主流程
大量 CPU 计算可用使用 CPU 线程池,避免和 IO 混用
MQ 消费内部并发谨慎offset 提交、幂等、失败重试要设计好

线上排查

现象一:接口用了 CompletableFuture 仍然很慢

排查路径:

mermaid
flowchart TD
    A["接口仍然慢"] --> B["看每个子任务耗时"]
    B --> C{"是否某个下游 P99 很高"}
    C -- "是" --> D["查下游超时、连接池、慢 SQL"]
    C -- "否" --> E["看聚合线程池 active 和 queue"]
    E --> F{"线程池是否满"}
    F -- "是" --> G["降并发、扩容或隔离线程池"]
    F -- "否" --> H["查是否 join/get 放错位置导致串行等待"]

常见错误:

java
String user = CompletableFuture.supplyAsync(() -> queryUser(), POOL).join();
String order = CompletableFuture.supplyAsync(() -> queryOrder(), POOL).join();

这其实是提交一个等一个,仍然接近串行。正确做法是先创建所有 Future,再统一等待。

现象二:公共 ForkJoinPool 被打满

原因:

  1. supplyAsync 没传线程池。
  2. parallelStream 和 CompletableFuture 共用公共池。
  3. 任务里做阻塞 IO。

处理:

  1. 业务异步必须传自定义线程池。
  2. IO 任务和 CPU 任务分池。
  3. 给线程命名,接入监控。

现象三:allOf 卡住或异常

常见原因:

原因处理
子任务没有超时加编排层和客户端层超时
子任务异常未处理根据业务选择兜底或整体失败
子任务线程池队列堆积看 active、queue、拒绝次数
子任务里又等待同池任务避免线程饥饿死锁

现象四:超时后线程仍然被占用

原因是 Future 超时完成不等于底层 IO 被强制中断。要检查:

  1. HTTP read timeout。
  2. DB query timeout。
  3. RPC timeout。
  4. 任务是否响应中断。
  5. 是否调用了 cancel(true),以及任务是否能响应。

常见坑

后果正确做法
不传线程池阻塞公共池自定义业务线程池
在链路中提前 join并行退化成串行先创建 Future,再统一等待
allOf 后不处理子任务异常不知道哪个任务失败每个子任务加日志和兜底
异步任务无超时线程长期占用编排层和客户端层都超时
IO 和 CPU 混用一个池互相拖慢按任务类型隔离
兜底返回默认成功掩盖核心失败区分核心和非核心依赖
忽略 traceId 传递日志链路断裂包装任务或使用上下文传递组件

面试标准回答

CompletableFuture 是什么,和 Future 有什么区别

标准回答:

text
Future 表示一个未来结果,主要通过 get 阻塞等待结果,不擅长任务组合和依赖编排。CompletableFuture 是 JDK 8 提供的异步编排工具,既表示未来结果,也可以注册 thenApply、thenCompose、thenCombine、allOf 等后续阶段,把多个异步任务组织成依赖链或并行聚合。它本身不是新线程,真正执行任务的仍然是线程池。生产使用时必须指定业务线程池,处理异常和超时,避免阻塞公共 ForkJoinPool。

thenApply、thenCompose、thenCombine 区别

标准回答:

text
thenApply 用于把上一步结果同步转换成另一个普通结果;thenCompose 用于上一步结果决定下一个异步任务,可以把嵌套的 CompletableFuture 拉平;thenCombine 用于两个互不依赖的异步任务都完成后合并结果。简单说,普通转换用 thenApply,依赖异步串联用 thenCompose,两个并行结果合并用 thenCombine。

CompletableFuture 线上使用注意什么

标准回答:

text
线上使用 CompletableFuture 要注意五点:第一,必须指定自定义线程池,避免阻塞公共 ForkJoinPool;第二,先创建多个 Future 再统一等待,避免提前 join 退化成串行;第三,每个下游任务要有超时和异常兜底;第四,区分核心依赖和非核心依赖,不能把订单主信息失败兜底成成功;第五,监控线程池 active、queue、拒绝次数和子任务 P95/P99,排查时结合 jstack 看线程是否卡在下游 IO、数据库连接或 Future.get。

关联知识点

知识点说明
CompletableFuture 基础基本 API 用法
Callable、Future 与 FutureTaskJDK 7/8 异步结果模型对比
线程池总览线程池是异步任务真正执行者
线程池生命周期与线上排查线程堆积、Future.get、jstack 排查
JDK 版本差异JDK 7 Future 到 JDK 8 CompletableFuture

本章小结

CompletableFuture 的价值是表达异步依赖关系,不是让系统无限变快。它能把串行调用改成并行聚合,也能把依赖任务串成清晰链路;但如果线程池、超时、异常、下游容量和监控没设计好,异步只会把问题藏到后台线程里。真正会用 CompletableFuture,要同时懂任务编排、线程池容量、异常传播、超时取消和商业降级策略。