Skip to content
Go back

CompletableFuture——异步编程的链式编排

CompletableFuture:异步编程的链式编排

一句话结论(30s)

CompletableFuture 的本质是”可显式完成的 Future + 链式回调编排”,因为异步任务需要把回调从嵌套地狱变成声明式流水线。关键设计是 thenApply/thenCompose/thenCombine 等组合算子,能把串行 165ms 的多次调用压成 max(各下游) 的并发耗时。权衡是它默认跑在共享的 ForkJoinPool.commonPool 上,IO 密集任务若不显式传自定义线程池,会耗尽共享线程、拖垮全 JVM 的异步任务。

核心原理(2min)

supplyAsync/runAsync 提交任务后返回 CompletableFuture,通过 thenApply(同步转换)、thenCompose(异步串联)、thenCombine/allOf/anyOf(并行汇聚)、exceptionally/handle(异常兜底)做链式编排。先带着问题看:CompletableFuture 为什么要回调链?——因为异步任务若靠 Future.get() 阻塞等待再手写嵌套回调,会迅速陷进「回调地狱」;链式算子把「先 A 再 B 再 C」的依赖关系拍平成一行声明式代码,可读可组合。核心原则:阻塞 IO 任务必须显式传入自定义 Executor,把 commonPool 留给计算型任务,避免阻塞任务污染全 JVM 共享线程池。读到这句不妨反问:为什么 IO 任务会「污染」线程池?因为阻塞的线程占着坑不干活,池里可用的线程被耗光后,后面提交的异步任务都得排队干等。

底层深入(5-10min)

thenApply 的”链式编排”本质上是把每个算子包装成一个 Completion 节点,挂到源 future 的无锁栈上;源一完成就沿着栈逐个触发下游。下面按调用链拆开看真实 JDK 源码(基于 java.base 的 CompletableFuture.java)。

1. thenApply 入口——只是薄薄一层转发

public <U> CompletableFuture<U> thenApply(
    Function<? super T,? extends U> fn) {
    return uniApplyStage(null, fn);
}

thenApply 自身不做任何计算,只把函数 fn 转交给 uniApplyStage,且第一个参数传 null,表示”不指定 executor”。这个 null 是同步语义的关键:后续 tryFire 会因此直接在当前调用线程里执行 fn,而不是丢进线程池。对应地,thenApplyAsync 只是把这里的 null 换成 defaultExecutor() 或用户传入的 executor。这里可以想想:为什么一个 null 就能决定「同步还是异步」?因为 executor 的语义就是「fn 在哪个线程跑」——传了 executor 就把回调丢进那个池,传 null 就就地用当前线程跑,一步之差,语义完全不同。

2. uniApplyStage——已完成就立刻算,未完成就注册回调

private <V> CompletableFuture<V> uniApplyStage(
    Executor e, Function<? super T,? extends V> f) {
    Object r;
    Objects.requireNonNull(f);
    if ((r = result) != null)
        return uniApplyNow(r, e, f);
    CompletableFuture<V> d = newIncompleteFuture();
    unipush(new UniApply<T,V>(e, d, this, f));
    return d;
}

这里体现了 CompletableFuture 最核心的两条路径:如果源 future 的 result 已就绪,就直接走 uniApplyNow 同步算完,绝不额外建链;否则 newIncompleteFuture() 造出一个下游 future d,再用 unipush 把一个 UniApply 回调节点挂到 this 的栈上。无论哪条路,方法都立即返回一个新的 future,这正是链式调用能不断往下 .thenXxx 的原因。想想看,为什么非要「立即返回新 future」?因为链式调用的前提是每一步都拿到一个可继续挂回调的句柄——如果等结果算完才返回,.thenXxx 就串不起来了,异步编排也就无从谈起。

3. Completion——回调链的基类,靠 next 串成 Treiber 栈

abstract static class Completion extends ForkJoinTask<Void>
    implements Runnable, AsynchronousCompletionTask {
    volatile Completion next;      // Treiber stack link

    abstract CompletableFuture<?> tryFire(int mode);

    abstract boolean isLive();

    public final void run()                { tryFire(ASYNC); }
    public final boolean exec()            { tryFire(ASYNC); return false; }
    public final Void getRawResult()       { return null; }
    public final void setRawResult(Void v) {}
}

Completion 是所有回调节点的公共基类,它同时是一个 ForkJoinTaskRunnable,所以节点既可以被线程池执行、也能被直接 run。字段 nextvolatile 维护,把同一 future 上的多个回调连成一个无锁的 Treiber 栈;tryFire(int mode) 是触发回调的抽象入口,run/exec 只是把它包装成 ASYNC 模式交给线程池。

4. UniCompletion——用 src/dep/executor 三元组串起串行依赖

abstract static class UniCompletion<T,V> extends Completion {
    Executor executor;                 // executor to use (null if none)
    CompletableFuture<V> dep;          // the dependent to complete
    CompletableFuture<T> src;          // source for action

    UniCompletion(Executor executor, CompletableFuture<V> dep,
                  CompletableFuture<T> src) {
        this.executor = executor; this.dep = dep; this.src = src;
    }
}

thenApply 链上的每一环在内存里就是这样一个节点:src 指向上游 future,dep 指向新创建的下游 future,executor 决定回调在哪跑。所谓”串行依赖如何串联”,本质就是上一环的 dep 恰好是下一环的 src,节点之间再靠 next 挂在同一个源栈上,形成一条可逐层触发的链。不妨把这条链想象成一排多米诺骨牌:上游 future 完成是「推倒第一张」,随后沿着 dep→src 逐层点燃,每张牌倒下都会触发下一张。

5. UniApply.tryFire——真正执行 fn 并把结果写进下游

final CompletableFuture<V> tryFire(int mode) {
    CompletableFuture<V> d; CompletableFuture<T> a;
    Object r; Throwable x; Function<? super T,? extends V> f;
    if ((a = src) == null || (r = a.result) == null
        || (d = dep) == null || (f = fn) == null)
        return null;
    tryComplete: if (d.result == null) {
        if (r instanceof AltResult) {
            if ((x = ((AltResult)r).ex) != null) {
                d.completeThrowable(x, r);
                break tryComplete;
            }
            r = null;
        }
        try {
            if (mode <= 0 && !claim())
                return null;
            else {
                @SuppressWarnings("unchecked") T t = (T) r;
                d.completeValue(f.apply(t));
            }
        } catch (Throwable ex) {
            d.completeThrowable(ex);
        }
    }
    src = null; dep = null; fn = null;
    return d.postFire(a, mode);
}

这是回调真正”点火”的地方:先校验源结果已就绪、下游和函数都还在,再处理 AltResult 包装的异常传播;通过后执行 f.apply(t) 并用 completeValue 把结果写进下游 d,异常则走 completeThrowableclaim() 用 CAS 保证同一个节点只有一个线程能执行,避免重复计算。执行完把 src/dep/fn 全部置空利于 GC,最后调用 postFire 决定是否继续向后触发。

6. postComplete——源完成后逐个触发下游,避免递归爆栈

final void postComplete() {
    CompletableFuture<?> f = this; Completion h;
    while ((h = f.stack) != null ||
           (f != this && (h = (f = this).stack) != null)) {
        CompletableFuture<?> d; Completion t;
        if (STACK.compareAndSet(f, h, t = h.next)) {
            if (t != null) {
                if (f != this) {
                    pushStack(h);
                    continue;
                }
                NEXT.compareAndSet(h, t, null); // try to detach
            }
            f = (d = h.tryFire(NESTED)) == null ? this : d;
        }
    }
}

postComplete 是驱动整条链的引擎:它在一个循环里用 CAS 从 stack 弹出节点并逐个 tryFire(NESTED)tryFire 返回的下游 future 会成为新的 f 继续深挖。这样遍历只沿单条依赖链递推,其余分支压回栈里,把本可能很深的递归摊平成迭代,避免栈溢出;多个线程可并发调用 postComplete,靠 CAS 竞争出栈权保证每个节点只被处理一次。

基础用法

// 异步执行 + 链式处理
CompletableFuture<String> future = CompletableFuture
    .supplyAsync(() -> db.getUser(1L))        // 异步查数据库
    .thenApply(user -> user.getName())         // 提取名字
    .thenApply(name -> "Hello, " + name);      // 格式化

String result = future.join();  // 阻塞等待结果

核心 API

异步执行

supplyAsync(Supplier)    // 有返回值
runAsync(Runnable)       // 无返回值

默认使用 ForkJoinPool.commonPool()(共享线程池)。如果任务中包含阻塞 IO(数据库/HTTP 调用),应该传自定义线程池作为第二个参数——阻塞任务不应占用 commonPool 的共享线程。再想一层:为什么官方不干脆默认给阻塞任务配独立池?因为 JDK 无法判断你的任务里有没有 IO——这个判断权必须留给调用者,所以第二个 Executor 参数就是「性能开关」,用对是提速,用错是全 JVM 卡死。

链式转换

future.thenApply(fn)        // 同步转换:Future<T> → Future<U>
future.thenCompose(fn)      // 异步串联:Future<T> → Future<U>(fn 返回 Future)
future.thenCombine(other, fn)  // 并行汇聚:Future<T> + Future<U> → Future<V>

汇聚

// 两个任务都完成 → 合并结果
shopFuture.thenCombine(scoreFuture, (shop, score) -> new ShopVO(shop, score));

// 所有任务都完成 → 返回
CompletableFuture.allOf(f1, f2, f3).join();

// 任一任务完成 → 返回
CompletableFuture.anyOf(f1, f2, f3).join();

异常处理

future.exceptionally(ex -> fallback)        // 异常恢复(返回默认值)
      .handle((result, ex) -> { ... });     // 正常/异常都处理

实战:并发调用 4 个下游

public ShopDetailVO getShopDetail(Long shopId) {
    var shopF  = CompletableFuture.supplyAsync(() -> shopService.getShop(shopId), executor);
    var scoreF = CompletableFuture.supplyAsync(() -> scoreService.getScore(shopId), executor);
    var blogF  = CompletableFuture.supplyAsync(() -> blogService.getHotBlogs(shopId), executor);
    var nearF  = CompletableFuture.supplyAsync(() -> esService.searchNearby(shopId), executor);

    return CompletableFuture.allOf(shopF, scoreF, blogF, nearF)
        .thenApply(v -> new ShopDetailVO(shopF.join(), scoreF.join(), blogF.join(), nearF.join()))
        .join();
}

串行耗时 50+5+80+30=165ms → 并发耗时 max(50,5,80,30)=80ms。

线程池陷阱

CompletableFuture.supplyAsync() 默认用 ForkJoinPool.commonPool()——这是一个全 JVM 共享的线程池。如果多个服务无差别地往 commonPool 里提交阻塞任务,线程耗尽 → 所有 CompletableFuture 都卡住。这就是「共享资源被个别租户拖垮」的典型案例——想一想,你项目里某个慢 SQL 或超时 HTTP 调用,是不是正在悄悄耗尽 commonPool,让别的异步任务排队等死?

原则:IO 密集任务用自己的 Executors.newFixedThreadPool(N),不能污染 commonPool。

章末提问

  1. CompletableFuture 为什么要回调链? 结论先行:把嵌套回调变声明式流水线。因为异步依赖若靠 Future.get() 阻塞等待 + 手写嵌套,会成回调地狱,链式算子让「先 A 再 B 再 C」线性可读。

  2. thenApply 和 thenCompose 有什么区别? 结论先行:同步转换 vs 异步串联。因为 thenApply 的 fn 返回普通值,thenCompose 的 fn 返回 Future,用 thenApply 接异步会得到嵌套的 Future<Future>。

  3. supplyAsync 默认跑在哪个线程池? 结论先行:ForkJoinPool.commonPool,全 JVM 共享。因为不显式传 executor 时它就用这个共享池。

  4. 为什么 IO 密集任务必须自定义线程池? 结论先行:避免阻塞耗尽共享池。因为阻塞线程占着坑不干活,耗尽后全 JVM 的异步任务都会卡住。

  5. thenApply 不传 executor 为什么是同步执行? 结论先行:内部把 executor 置为 null。因为 tryFire 发现 executor 为 null 就直接在当前线程执行 fn,而不是丢进线程池。


Share this post on:

Previous Post
CountDownLatch、CyclicBarrier、Semaphore——三大并发工具辨析
Next Post
AQS——JUC所有锁的共同底座