CompletableFuture.runAsync 详解

coverImg

介绍

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默认的线程名为:

img

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("所有任务完成");
    }
}

img

看现象是默认线程池中线程数量为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();
    }
}

img

核心原理分析

默认线程池:使用 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

评论区请在客户端页面查看