Skip to content

CompletableFuture

CompletableFuture 是 Java 8 提供的异步编排工具。它常用于并行调用多个接口、异步执行耗时任务、组合多个任务结果。相比手动创建线程和等待,CompletableFuture 更适合表达异步任务之间的依赖关系。

本页先讲基础 API。如果要深入理解线程池选择、非 Async/Async 回调线程、异常传播、allOf 失败策略、超时是否真正取消底层任务、接口聚合线上排查,继续看:CompletableFuture 全过程原理与商业异步编排

为什么需要 CompletableFuture

假设一个订单详情页需要查询用户、订单、优惠券、物流信息。如果串行调用,每个接口耗时 200ms,总耗时可能接近 800ms。

mermaid
flowchart TD
    A["查询用户 200ms"] --> B["查询订单 200ms"]
    B --> C["查询优惠券 200ms"]
    C --> D["查询物流 200ms"]
    D --> E["总耗时约 800ms"]

如果这些查询互不依赖,可以并行执行。

mermaid
flowchart TD
    A["发起请求"] --> B["查询用户"]
    A --> C["查询订单"]
    A --> D["查询优惠券"]
    A --> E["查询物流"]
    B --> F["汇总结果"]
    C --> F
    D --> F
    E --> F

CompletableFuture 执行原理

CompletableFuture 可以理解为“一个将来会完成的结果”。它内部保存任务状态、结果、异常和后续依赖动作。

mermaid
flowchart TD
    A["提交异步任务"] --> B["线程池执行任务"]
    B --> C["任务成功或失败"]
    C --> D["完成 Future 状态"]
    D --> E["触发后续 thenApply 或 thenAccept"]
    E --> F["继续向后传播结果"]

它不是创建了一个神秘的新执行模型,真正干活的仍然是线程池。CompletableFuture 负责把多个任务的完成关系组织起来。

基本使用

java
import java.util.concurrent.CompletableFuture;

public class CompletableFutureBasicDemo {
    public static void main(String[] args) {
        CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
            sleep(500);
            return "user";
        });

        String result = future.join();
        System.out.println(result);
    }

    private static void sleep(long millis) {
        try {
            Thread.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

supplyAsync 有返回值,runAsync 没有返回值。

thenApply、thenAccept、thenRun

java
import java.util.concurrent.CompletableFuture;

public class CompletableFutureChainDemo {
    public static void main(String[] args) {
        CompletableFuture.supplyAsync(() -> "Tom")
                .thenApply(name -> "Hello " + name)
                .thenAccept(System.out::println)
                .thenRun(() -> System.out.println("处理完成"))
                .join();
    }
}
方法作用
thenApply接收上一步结果并返回新结果
thenAccept接收上一步结果但不返回
thenRun不接收结果,也不返回

thenCompose 和 thenCombine

thenCompose 用于有依赖的异步任务,thenCombine 用于合并两个独立任务结果。

java
import java.util.concurrent.CompletableFuture;

public class ComposeCombineDemo {
    public static void main(String[] args) {
        CompletableFuture<String> detailFuture = queryUserId()
                .thenCompose(ComposeCombineDemo::queryUserDetail);

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

        System.out.println(detailFuture.join());
        System.out.println(mergedFuture.join());
    }

    private static CompletableFuture<Long> queryUserId() {
        return CompletableFuture.supplyAsync(() -> 1001L);
    }

    private static CompletableFuture<String> queryUserDetail(Long userId) {
        return CompletableFuture.supplyAsync(() -> "user-" + userId);
    }

    private static CompletableFuture<String> queryOrder() {
        return CompletableFuture.supplyAsync(() -> "order");
    }

    private static CompletableFuture<String> queryCoupon() {
        return CompletableFuture.supplyAsync(() -> "coupon");
    }
}

allOf 汇总多个任务

java
import java.util.concurrent.CompletableFuture;

public class AllOfDemo {
    public static void main(String[] args) {
        CompletableFuture<String> userFuture = CompletableFuture.supplyAsync(() -> "user");
        CompletableFuture<String> orderFuture = CompletableFuture.supplyAsync(() -> "order");
        CompletableFuture<String> couponFuture = CompletableFuture.supplyAsync(() -> "coupon");

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

        String result = userFuture.join() + ", " + orderFuture.join() + ", " + couponFuture.join();
        System.out.println(result);
    }
}

allOf 本身不直接返回每个任务结果,只表示所有任务都完成。结果需要从各个 future 中取。

异常处理

java
import java.util.concurrent.CompletableFuture;

public class FutureExceptionDemo {
    public static void main(String[] args) {
        String result = CompletableFuture.supplyAsync(() -> {
            if (true) {
                throw new RuntimeException("查询失败");
            }
            return "success";
        }).exceptionally(ex -> {
            System.out.println("异常:" + ex.getMessage());
            return "default";
        }).join();

        System.out.println(result);
    }
}

常见方法:

方法作用
exceptionally异常时返回兜底值
handle成功和异常都处理,并返回结果
whenComplete成功和异常都能观察,但通常不改变结果

线程池为什么必须显式指定

不指定线程池时,异步任务默认使用公共线程池。业务系统里不建议重要任务都挤在公共线程池里。

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

public class FutureExecutorDemo {
    private static final ExecutorService IO_POOL = Executors.newFixedThreadPool(8);

    public static void main(String[] args) {
        CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
            // 模拟 IO 查询
            return "result";
        }, IO_POOL);

        System.out.println(future.join());
        IO_POOL.shutdown();
    }
}

生产项目不要直接用 Executors.newFixedThreadPool 后就不管,还需要考虑线程命名、队列长度、拒绝策略和监控。这里为了演示保持简单。

超时控制

Java 9 之后可以使用 orTimeoutcompleteOnTimeout

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

public class FutureTimeoutDemo {
    public static void main(String[] args) {
        String result = CompletableFuture.supplyAsync(() -> {
            sleep(3000);
            return "slow result";
        }).completeOnTimeout("timeout default", 1, TimeUnit.SECONDS).join();

        System.out.println(result);
    }

    private static void sleep(long millis) {
        try {
            Thread.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
}

远程调用必须有超时,否则异步任务也会堆积。

商业 Demo:订单详情并行聚合

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

public class OrderDetailAggregateDemo {
    private static final ExecutorService POOL = Executors.newFixedThreadPool(8);

    public static void main(String[] args) {
        CompletableFuture<String> user = CompletableFuture.supplyAsync(() -> "用户信息", POOL);
        CompletableFuture<String> order = CompletableFuture.supplyAsync(() -> "订单信息", POOL);
        CompletableFuture<String> logistics = CompletableFuture.supplyAsync(() -> "物流信息", POOL);

        CompletableFuture.allOf(user, order, logistics).join();

        String detail = user.join() + "," + order.join() + "," + logistics.join();
        System.out.println(detail);
        POOL.shutdown();
    }
}

真实项目里要加超时、异常兜底、线程池隔离、日志追踪。

常见风险

问题后果建议
不指定线程池公共线程池被阻塞业务线程池隔离
异步任务里做长时间阻塞线程池耗尽设置超时和容量
忘记处理异常join 时抛包装异常使用 exceptionallyhandle
任务依赖关系写错结果提前读取或死等区分 thenComposethenCombine
无限并行调用下游打爆外部接口限流和线程池隔离

本章小结

CompletableFuture 用来表达异步任务和任务依赖。它适合接口聚合、并行查询、异步处理,但不是性能魔法。生产使用时必须显式线程池、处理异常、设置超时,并控制对下游的并发压力。

深入学习:CompletableFuture 全过程原理与商业异步编排