
介绍
CompletableFuture 是一个高度并发化的异步编程工具,其核心机制包括:
- 使用
ForkJoinPool.commonPool()或自定义线程池来执行异步任务。 - 使用 Treiber 栈来管理依赖任务,确保高效的并发操作。
- 提供了同步和异步两种执行模式,支持任务链式调用和组合。
- 提供了丰富的异常处理机制,支持任务的正常和异常完成。
核心注意点
梳理汇总
| 常见问题场景 | 核心原因 | 关键应对建议 |
|---|---|---|
| 🐢 任务执行缓慢或阻塞 | 默认线程池 (ForkJoinPool.commonPool()) 核心线程数有限,易在 IO 密集型任务时等待 |
为 IO 密集型任务使用自定义线程池 |
| 🔒 任务未按预期执行 | 未正确获取结果(如未调用 get(), join()),或任务代码自身异常未处理 |
根据需求正确处理异步任务结果,并使用 exceptionally 等方法处理异常 |
| 🧵 资源耗尽与服务崩溃 | 自定义线程池配置不当(如核心线程数设置过小或过大) | 合理配置自定义线程池参数,避免与业务量不匹配 |
| ⚠️ 线程死锁 | 多个任务循环等待,且共用同一线程池 | 不同业务类型考虑使用独立的线程池,避免池内任务相互阻塞 |
异常捕获
仅仅在 runAsync 外部使用 try-catch 是无法捕获到异步任务内部抛出的异常的 。这是因为异常发生在另一个线程中。
错误示例:
public void asyncException() {
try {
CompletableFuture.runAsync(() -> {
int i = 1 / 0; // 这里会抛出 ArithmeticException
}, customThreadPool);
// 这里的 try-catch 抓不到上面的除零异常
} catch (Exception e) {
System.out.println("捕获到异常: " + e.getMessage()); // 这行不会执行
}
}
正确的处理方式之一(在异步内部处理):
CompletableFuture.runAsync(() -> {
try {
int i = 1 / 0;
} catch (Exception e) {
System.out.println("在异步任务内部处理异常: " + e.getMessage());
// 记录日志、进行补偿操作等
}
}, customThreadPool);
或者,使用 **exceptionally** 方法 :
CompletableFuture.runAsync(() -> {
int i = 1 / 0;
}, customThreadPool).exceptionally(e -> {
System.out.println("处理异步执行中的异常: " + e.getMessage());
return null; // 因为 runAsync 没有返回值,这里返回 null
});
细节汇总说明
原生CompletableFuture.runAsync默认的线程名为:

ForkJoinPool.commonPool-worker-xxx
CompletableFuture相关API
1、任务完成回调
CompletableFuture.runAsync(() -> {
System.out.println("异步任务执行");
})
.thenRun(() -> System.out.println("任务完成回调"))
.thenRunAsync(() -> System.out.println("异步任务完成回调"));
2、异常处理
CompletableFuture.runAsync(() -> {
throw new RuntimeException("出错了!");
})
.exceptionally(ex -> {
System.out.println("捕获异常: " + ex.getMessage());
return null;
});
3、任务组合
CompletableFuture<Void> task1 = CompletableFuture.runAsync(() -> System.out.println("任务1"));
CompletableFuture<Void> task2 = CompletableFuture.runAsync(() -> System.out.println("任务2"));
// 等待所有完成
CompletableFuture<Void> all = CompletableFuture.allOf(task1, task2);
// 等待任意一个完成
CompletableFuture<Object> any = CompletableFuture.anyOf(task1, task2);
4、结果转换(虽然runAsync无返回值,但可以链式转换)
CompletableFuture.runAsync(() -> System.out.println("原始任务"))
.thenApply(v -> {
System.out.println("转换操作");
return "转换后的结果";
})
.thenAccept(result -> System.out.println("消费结果: " + result));
5、超时控制
CompletableFuture.runAsync(() -> {
try { Thread.sleep(5000); } catch (InterruptedException e) {}
})
.orTimeout(2, TimeUnit.SECONDS) // 2秒超时
.exceptionally(ex -> {
System.out.println("任务超时: " + ex.getMessage());
return null;
});
完整链式调用示例
CompletableFuture.runAsync(() -> {
System.out.println("执行核心业务");
})
.thenRun(() -> System.out.println("后续操作1"))
.thenRunAsync(() -> System.out.println("后续操作2"))
.exceptionally(ex -> {
System.out.println("错误处理: " + ex.getMessage());
return null;
})
.thenRun(() -> System.out.println("最终清理"))
.join(); // 等待整个链完成
实战案例demo
1、原生CompletableFuture.runAsync跑任务
public class CompletableFutureDemo {
public static void main(String[] args) throws InterruptedException {
System.out.println("主线程开始 - " + Thread.currentThread().getName());
List<CompletableFuture<Void>> futures = new ArrayList<>();
for (int i = 0; i < 1000; i++) {
// 使用runAsync执行异步任务
CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
System.out.println("异步任务开始 - " + Thread.currentThread().getName());
try {
TimeUnit.SECONDS.sleep(2); // 模拟耗时操作
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("异步任务完成 - " + Thread.currentThread().getName() + "\n");
});
futures.add(future);
}
for (CompletableFuture<Void> future : futures) {
future.join();
}
System.out.println("所有任务完成");
}
}

看现象是默认线程池中线程数量为7个,且在不断的去执行中。
2、自定义线程池案例
public class CompletableFutureWithCustomPool {
public static void main(String[] args) {
// 创建自定义线程池
ExecutorService customExecutor = Executors.newFixedThreadPool(3);
System.out.println("=== 使用自定义线程池 ===");
// 使用自定义线程池执行异步任务
CompletableFuture<Void> future1 = CompletableFuture.runAsync(() -> {
System.out.println("异步任务开始 - " + Thread.currentThread().getName());
try {
TimeUnit.SECONDS.sleep(2); // 模拟耗时操作
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("异步任务完成 - " + Thread.currentThread().getName() + "\n");
}, customExecutor);
CompletableFuture<Void> future2 = CompletableFuture.runAsync(() -> {
System.out.println("异步任务开始 - " + Thread.currentThread().getName());
try {
TimeUnit.SECONDS.sleep(2); // 模拟耗时操作
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("异步任务完成 - " + Thread.currentThread().getName() + "\n");
}, customExecutor);
CompletableFuture<Void> future3 = CompletableFuture.runAsync(() -> {
System.out.println("异步任务开始 - " + Thread.currentThread().getName());
try {
TimeUnit.SECONDS.sleep(2); // 模拟耗时操作
} catch (InterruptedException e) {
e.printStackTrace();
}
System.out.println("异步任务完成 - " + Thread.currentThread().getName() + "\n");
}, customExecutor);
// 等待所有任务完成
CompletableFuture.allOf(future1, future2, future3).join();
// 关闭线程池
customExecutor.shutdown();
}
}

核心原理分析
默认线程池:使用 ForkJoinPool.commonPool()
线程数量:默认是 CPU 核心数 - 1
任务队列:使用工作窃取(work-stealing)算法
异步执行:立即返回 CompletableFuture 对象,不阻塞调用线程
线程池的使用
ForkJoinPool 在处理大量可分解的并行任务时(如递归计算、流处理)具有极高的性能,但在处理大量独立短期任务时可能不如传统的 ThreadPoolExecutor 高效。
CompletableFuture 默认使用 ForkJoinPool.commonPool() 作为线程池来执行异步任务。如果 commonPool 的并行度小于 2(即系统只有一个核心),则会退化为每个任务创建一个新线程(通过 ThreadPerTaskExecutor 实现)。
**ForkJoinPool.commonPool()**:
-
- 这是一个共享的线程池,通常用于执行轻量级任务。
- 如果并行度大于 1,则使用
commonPool;否则,使用单线程模式。
// 编译的时候确认
// 这个“并行度”通常等于 CPU 的可用核心数 - 1。例如,在一个 8 核的机器上,它通常返回 7。
// 它也可以通过 JVM 启动参数 -Djava.util.concurrent.ForkJoinPool.common.parallelism 来显式设置。
// 若是false情况:
// 单核 CPU:系统只有一个处理器核心,公共池的并行度被计算为 1;
// 显式配置:通过 -Djava.util.concurrent.ForkJoinPool.common.parallelism=0 或 1 手动设置了很低的并行度。
private static final boolean useCommonPool =
(ForkJoinPool.getCommonPoolParallelism() > 1);
// =====================
// 根据机器实际情况进行选择
private static final Executor asyncPool = useCommonPool ?
ForkJoinPool.commonPool() : new ThreadPerTaskExecutor();
// =====================
// 这个线程池每来一个任务就创建一个线程去跑
static final class ThreadPerTaskExecutor implements Executor {
public void execute(Runnable r) { new Thread(r).start(); }
}
源码如下:
private static final Executor asyncPool = useCommonPool ?
ForkJoinPool.commonPool() : new ThreadPerTaskExecutor();
public static CompletableFuture<Void> runAsync(Runnable runnable) {
return asyncRunStage(asyncPool, runnable);
}
情况1:new ThreadPerTaskExecutor()
// 编译的时候确认
// 这个“并行度”通常等于 CPU 的可用核心数 - 1。例如,在一个 8 核的机器上,它通常返回 7。
// 它也可以通过 JVM 启动参数 -Djava.util.concurrent.ForkJoinPool.common.parallelism 来显式设置。
// 若是false情况:
// 单核 CPU:系统只有一个处理器核心,公共池的并行度被计算为 1;
// 显式配置:通过 -Djava.util.concurrent.ForkJoinPool.common.parallelism=0 或 1 手动设置了很低的并行度。
private static final boolean useCommonPool =
(ForkJoinPool.getCommonPoolParallelism() > 1);
// 这个线程池每来一个任务就创建一个线程去跑
static final class ThreadPerTaskExecutor implements Executor {
public void execute(Runnable r) { new Thread(r).start(); }
}
情况2:使用ForkJoinPool.commonPool()
// ForkJoinPool
common = java.security.AccessController.doPrivileged
(new java.security.PrivilegedAction<ForkJoinPool>() {
public ForkJoinPool run() { return makeCommonPool(); }});
ForkJoinPool 类中初始化公共池(common pool)的核心代码,它涉及Java中一个比较高级和重要的概念:特权操作(Privileged Action)。
该默认线程池实现底层逻辑参考:https://www.yuque.com/changlu-azwmm/ib7lmr/va78bvw4et2pswsp
参考文章
[1]. CompletableFuture.runAsync()使用不当导致生产问题:https://blog.csdn.net/my_1996/article/details/141901820#comments_35857450
评论区请在客户端页面查看