一、前言

Java 8 引入 CompletableFuture,弥补了原生 Future 的致命缺陷:无法回调、只能阻塞轮询、多任务编排复杂。 它同时实现 FutureCompletionStage,支持异步链式调用、多任务组合、异常捕获、手动完成等能力。 很多开发者只会使用 thenApplysupplyAsync,但对底层存储结构、无锁并发模型、回调链表、线程调度原理一知半解。 本文基于 OpenJDK8 源码,从底层字段存储、链表结构、CAS 原子操作、任务执行、回调触发、阻塞唤醒、多任务组合完整拆解,内容可直接用于面试复习与技术博客。

二、CompletableFuture 顶层核心字段与状态存储模型

2.1 完整类成员变量

public class CompletableFuture<T> implements Future<T>, CompletionStage<T> {
    // 实例变量:每个CompletableFuture独立持有
    volatile Object result;
    volatile Completion stack;

    // 全局静态常量
    private static final ForkJoinPool ASYNC_POOL;
    private static final Object NIL = new Object();
    private static final AltResult EXCEPTIONAL = new AltResult();

    // 异常包装内部类
    static final class AltResult {
        final Throwable ex;
        AltResult() { this.ex = null; }
        AltResult(Throwable x) { this.ex = x; }
    }

    // Unsafe CAS 相关偏移量
    private static final sun.misc.Unsafe UNSAFE;
    private static final long RESULT;
    private static final long STACK;
    static {
        try {
            UNSAFE = sun.misc.Unsafe.getUnsafe();
            Class<?> k = CompletableFuture.class;
            RESULT = UNSAFE.objectFieldOffset(k.getDeclaredField("result"));
            STACK = UNSAFE.objectFieldOffset(k.getDeclaredField("stack"));
        } catch (Exception x) { throw new Error(x); }
    }
}

2.2 result:单字段复合状态存储(核心设计)

result 是整个类最关键的变量,一个字段同时承载任务状态、正常返回值、异常信息,依靠 volatile 保证多线程可见,全程使用 CAS 无锁修改,不使用任何 synchronized 锁。

表格

result 赋值 任务状态 业务含义
null 未完成 任务未执行完毕,无结果无异常
NIL 正常完成 执行成功,返回值为 null
普通 Object 对象 正常完成 执行成功,携带业务返回值
AltResult 实例 异常完成 执行抛出异常,Throwable 存在 AltResult.ex
设计亮点

单独使用 AltResult 包装异常,解决边界问题:如果业务正常返回 Throwable 对象,不会被误判为任务执行失败,实现正常结果与异常结果隔离。

修改规则

禁止直接赋值 this.result = xxx,所有状态变更必须通过 Unsafe CAS 自旋更新:

java

运行

UNSAFE.compareAndSwapObject(this, RESULT, expectValue, newValue)

CAS 失败则循环重试,实现无锁并发。

2.3 stack:单向头插链表,存储回调与阻塞线程

stack 是链表头指针,底层是单链表,采用头插法新增节点,所有阻塞等待线程、链式回调任务全部挂载在该链表。 链表顶层抽象类 Completion

java

运行

abstract static class Completion extends ForkJoinTask<Void>
        implements Runnable, CompletableFuture.AsynchronousCompletionTask {
    volatile Completion next; // 后继指针,串联整条链表
    abstract boolean tryFire(int mode); // 尝试执行当前回调
    abstract CompletableFuture<?> getDep(); // 获取依赖的前置Future
    abstract void clean(); // 执行后资源清理
}

链表内只有两类节点:

  1. WaitNode:调用 join() / get() 阻塞时创建,保存阻塞线程;
  2. 各类 Completion 回调子类ThenApply / ThenRun / WhenComplete / ThenCombine 等,对应链式调用方法。
WaitNode 阻塞节点源码

java

运行

static final class WaitNode extends Completion {
    Thread thread;
    WaitNode(Thread t) { thread = t; }
    boolean tryFire(int mode) { return false; }
    CompletableFuture<?> getDep() { return null; }
    void clean() {}
}

作用:线程阻塞时存入链表,任务完成后通过 LockSupport.unpark() 唤醒。

三、无锁并发底层:Unsafe + CAS 完整机制

CompletableFuture 全程无重量级锁,并发安全完全依赖 CAS + volatile:

  1. resultstack 标记 volatile,写操作插入内存屏障,状态变更对其他线程立即可见;
  2. 修改字段统一使用 compareAndSwapObject,基于 CPU 指令级原子操作;
  3. CAS 竞争失败采用自旋重试,不会阻塞线程;
  4. postComplete 处理回调时会截断链表,避免遍历过程中新节点插入导致漏处理。

CAS 核心操作

  1. 修改任务状态 result

java

运行

UNSAFE.compareAndSwapObject(this, RESULT, expect, update)
  1. 修改链表头 stack(新增回调 / 等待节点)

java

运行

UNSAFE.compareAndSwapObject(this, STACK, expectHead, newHead)

四、异步任务执行全链路(supplyAsync 为例)

4.1 默认线程池:ForkJoinPool.commonPool ()

不带自定义 Executor 的 supplyAsync / runAsync 统一使用静态常量 ASYNC_POOL,即 JDK 公共 ForkJoin 池:

  • 内部采用工作窃取队列,多任务调度效率高于普通 ThreadPoolExecutor;
  • 池内线程为守护线程,JVM 退出自动回收;
  • 全局单例,复用线程,减少创建销毁开销。

4.2 任务提交流程

java

运行

public static <U> CompletableFuture<U> supplyAsync(Supplier<U> supplier) {
    return asyncSupplyStage(ASYNC_POOL, supplier);
}
  1. 创建空 CompletableFuture 对象,result = nullstack = null
  2. 封装 AsyncSupply(继承 ForkJoinTask、Runnable),持有 supplier 与当前 CF;
  3. 将 AsyncSupply 提交到 ForkJoinPool;
  4. 池线程执行 AsyncSupply.run ():
    • 执行 supplier.get() 获取业务返回值;
    • 正常执行:CAS 将 result 更新为返回值(null 则存 NIL);
    • 抛出异常:CAS 将 result 更新为 AltResult(throwable)
    • 执行 postComplete(),统一处理 stack 链表回调与阻塞线程。

五、核心方法 postComplete:回调链表统一处理入口

任务正常完成、异常完成、手动 complete 都会调用 postComplete(),是整个 CompletableFuture 的调度核心,负责遍历链表、唤醒阻塞线程、执行所有回调。

5.1 完整执行流程

  1. 开启循环,判断当前 stack 是否存在节点;
  2. CAS 将当前 CF 的 stack 设置为 null,截断整条链表,防止遍历期间并发新增节点干扰;
  3. 遍历截断后的链表(从栈顶到尾部):
    • 如果节点是 WaitNode:调用 LockSupport.unpark(node.thread) 唤醒阻塞线程;
    • 如果节点是 Completion 回调节点:调用 tryFire(NESTED) 执行回调;
  4. 回调执行时会生成新的 CompletableFuture,新回调会挂载到新 CF 的 stack;
  5. 截断链表后如果又有新节点插入 stack,循环重复执行,直到 stack 为空。

5.2 tryFire 回调执行逻辑

所有回调子类重写 tryFire,统一执行规则:

  1. 获取依赖的前置 CompletableFuture,判断是否完成;
  2. 未完成直接返回 false,节点放回链表等待下次触发;
  3. 已完成执行用户传入的 Function/Consumer/Runnable;
  4. 将执行结果 / 异常 CAS 设置到新生成的 CompletableFuture.result;
  5. 新 CF 自动调用自身 postComplete,形成链式回调传递。

六、同步回调 vs 异步回调底层差异

6.1 同步回调(thenApply /thenRun/thenAccept)

不带 Async 后缀的链式方法:

  • 不会新建线程;
  • 完成前置任务的线程内串行执行回调;
  • 性能更高,无线程切换开销;
  • 回调耗时过长会阻塞当前线程,影响后续任务。

6.2 异步回调(thenApplyAsync /thenRunAsync)

带 Async 后缀方法:

  • 将 Completion 回调任务提交到 ForkJoinPool / 自定义线程池;
  • 由池内独立线程执行回调;
  • 不会阻塞前置任务完成线程;
  • 存在线程切换、队列调度开销。

七、阻塞获取结果 join () /get () 底层原理

7.1 join () 无受检异常阻塞

  1. 循环读取 result:
    • result != null:任务已完成,判断是否为 AltResult,正常返回值,异常抛出 CompletionException;
    • result == null:任务未完成;
  2. 创建 WaitNode 节点,CAS 插入 stack 链表头部;
  3. 调用 LockSupport.park(this) 挂起当前线程;
  4. 任务完成后 postComplete 唤醒线程,再次读取 result 返回结果。

7.2 get () 带受检异常

底层阻塞逻辑和 join 完全一致,仅异常包装不同:异常封装为 ExecutionException,属于受检异常,必须 try-catch。

7.3 超时 get (long, TimeUnit)

使用 LockSupport.parkNanos() 限时阻塞,超时未完成直接抛出 TimeoutException。

八、多任务组合 allOf /anyOf 底层实现

8.1 allOf () 等待全部任务完成

  1. 创建空 CompletableFuture 作为最终返回对象;
  2. 为每一个入参 CF 注册独立的 ThenCombine 回调;
  3. 内部维护完成计数器,每个子任务完成计数器 +1;
  4. 计数器等于任务总数时,CAS 设置最终 CF 的 result,触发回调链。

8.2 anyOf () 任意一个完成即结束

  1. 创建空 CompletableFuture;
  2. 所有子任务注册竞争回调;
  3. 任意任务先完成,通过 CAS 抢占设置最终 CF 的 result;
  4. 后续完成的任务检测到最终 CF 已完成,直接跳过处理。

九、手动完成 complete /completeExceptionally

手动完成逻辑和异步任务执行完成逻辑完全统一:

  1. CAS 修改 result 为正常值 / AltResult 异常对象;
  2. CAS 修改成功后立即调用 postComplete ();
  3. 自动遍历 stack 链表,唤醒阻塞线程、执行全部注册回调;
  4. CAS 修改失败代表任务已提前完成,直接返回 false。

十、并发安全与内存模型总结

  1. 状态可见性:result、stack 使用 volatile,volatile 写屏障保证状态变更实时对其他线程可见;
  2. 无锁修改:全部字段更新依赖 Unsafe CAS,无 synchronized、AQS 锁,高并发性能优异;
  3. 链表安全:postComplete 阶段截断链表,避免并发新增节点造成漏执行;
  4. 线程阻塞:基于 LockSupport park/unpark,底层操作系统原语,无对象监视器开销;
  5. 链式传递:每个回调生成独立 CompletableFuture,各自维护独立 stack 链表,回调自动传递。

十一、CompletableFuture 与原生 Future 底层对比

  1. Future
    • 底层依赖 ThreadPoolExecutor + 阻塞队列;
    • 无内置回调存储结构,只能阻塞 / 轮询获取结果;
    • 多任务编排需要手动循环判断,代码繁琐。
  2. CompletableFuture
    • 内置 result 复合状态 + stack 回调链表双存储结构;
    • 任务完成主动推送回调,无需轮询;
    • 默认复用 ForkJoinPool 工作窃取线程池;
    • 无锁并发模型,支持链式、组合、异常全套能力。

十二、潜在内存泄漏隐患

  1. 大量注册永不执行的回调,Completion 链表无限膨胀;
  2. WaitNode 阻塞线程长期挂起,线程无法释放;
  3. 循环依赖的 CompletableFuture 互相持有引用,GC 无法回收;
  4. 自定义线程池未关闭,异步任务线程持续存活。

十三、总结

  1. 存储层:result 统一存储状态与结果,stack 单链表存储回调、阻塞线程;
  2. 并发层:Unsafe CAS + volatile 实现全程无锁;
  3. 执行层:同步回调复用完成线程,异步回调依托 ForkJoinPool;
  4. 调度层:postComplete 统一处理链表,唤醒线程、串行执行回调;
  5. 扩展层:allOf/anyOf、手动 complete、异常捕获均复用同一套链表调度逻辑。

CompletableFuture 的设计核心是主动回调替代被动轮询,依靠精巧的单字段状态模型、无锁链表结构,实现一套简洁、高性能的异步编程模型,也是面试高频底层考点。

#Java #并发编程 #CompletableFuture #JDK源码

Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐