跳到正文
hello world

41. Future、ThreadPoolExecutor、CompletableFuture 到底怎么选

发布于阅读量 0

41. Future、ThreadPoolExecutor、CompletableFuture 到底怎么选

前面分别学过:

ThreadPoolExecutor;
Future;
CompletableFuture。

这三个东西经常一起出现,所以很容易把它们混在一起。

刚开始我也会觉得:

它们是不是都在做异步任务?

既然 CompletableFuture 更强,是不是以后只用 CompletableFuture 就行?

用了 CompletableFuture,还需要 ThreadPoolExecutor 吗?

后来把几个 PDF 版本都写过一遍以后,我觉得这三个工具其实不在同一个层面。

可以先这样理解:

ThreadPoolExecutor 负责安排线程执行任务;

Future 负责保存某个异步任务将来的结果;

CompletableFuture 负责异步任务的结果和后续流程编排。

它们不是简单的替代关系。

很多时候,实际代码会同时使用:

ThreadPoolExecutor + CompletableFuture。

先看三者分别解决什么问题

ThreadPoolExecutor 解决的是线程管理问题

如果直接手动创建线程:

new Thread(() -> {
    processPdf(file);
}).start();

每来一个任务就可能创建一个线程。

文件一多,线程数量就容易失控。

ThreadPoolExecutor 解决的是:

最多创建多少线程;

任务如何排队;

线程如何复用;

队列满了怎么办;

线程池怎么关闭。

例如:

ThreadPoolExecutor pdfExecutor =
        new ThreadPoolExecutor(
                3,
                3,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(100),
                new NamedThreadFactory(
                        "pdf-worker"
                ),
                new ThreadPoolExecutor.AbortPolicy()
        );

这段代码确定了 PDF 任务的执行边界:

最多 3 个线程处理;

最多 100 个任务排队;

超过承载能力后明确拒绝。

所以 ThreadPoolExecutor 主要管理的是:

任务由哪些线程执行。

Future 解决的是任务结果问题

通过线程池提交有返回值的任务:

Future<String> future =
        pdfExecutor.submit(() -> {
            return processPdf(file);
        });

这里的 Future<String> 代表:

这个任务以后会返回一个 String。

需要结果时调用:

String result = future.get();

Future 主要解决的是:

异步任务执行完以后,怎么拿到结果。

它还能提供:

future.get(
        30,
        TimeUnit.SECONDS
);

future.cancel(true);

future.isDone();

future.isCancelled();

所以 Future 是一个异步任务的结果凭证。

它不负责创建线程,也不负责管理任务队列。

真正执行任务的仍然是线程池。


CompletableFuture 解决的是异步流程编排问题

CompletableFuture 也能表示异步任务结果。

例如:

CompletableFuture<String> future =
        CompletableFuture.supplyAsync(
                () -> processPdf(file),
                pdfExecutor
        );

但它不只是让外部主动获取结果。

还可以把后续步骤接起来:

CompletableFuture<PdfTaskResult> future =
        CompletableFuture
                .supplyAsync(
                        () -> processPdf(file),
                        pdfExecutor
                )
                .thenApply(path -> {
                    return buildDownloadUrl(path);
                })
                .thenApply(url -> {
                    return PdfTaskResult.success(
                            file.getName(),
                            url
                    );
                })
                .exceptionally(ex -> {
                    return PdfTaskResult.fail(
                            file.getName(),
                            getErrorMessage(ex)
                    );
                });

这里表达的是一条完整流程:

处理 PDF;

生成下载地址;

组装成功结果;

异常时返回失败结果。

所以 CompletableFuture 更擅长:

结果转换;

任务依赖;

多个任务合并;

批量等待;

异常兜底;

超时处理;

异步流程编排。

三者不在同一个层面

我现在会把它们分成两层。

第一层是执行资源:

ThreadPoolExecutor

它负责:

用多少个线程执行;

任务如何排队;

压力太大时怎么处理。

第二层是任务结果和流程:

Future;

CompletableFuture。

它们负责:

任务什么时候完成;

任务返回什么;

后续逻辑怎么继续。

可以画成这样:

业务任务
   ↓
Future / CompletableFuture
   ↓
ThreadPoolExecutor
   ↓
工作线程
   ↓
真正执行 processPdf()

CompletableFuture 不会凭空生成线程。

最终仍然需要线程来执行任务。

如果不指定线程池,它就使用默认公共线程池。

如果指定了 pdfExecutor,它就使用自己配置的线程池。


只用 ThreadPoolExecutor 可以吗

可以。

比如前面写过:

pdfExecutor.execute(() -> {
    processPdf(file);
});

如果我只关心任务被执行,不需要返回结果,这样已经够了。

例如:

异步写日志;

发送一个不关心结果的通知;

后台清理临时文件;

执行简单的定时任务。

这时候没必要为了使用 CompletableFuture,把代码写得更复杂。

如果需要结果,也可以使用:

Future<PdfTaskResult> future =
        pdfExecutor.submit(() -> {
            return processPdf(file);
        });

所以 ThreadPoolExecutor 本身已经能支持:

无返回值任务;

有返回值任务;

等待结果;

取消任务。

只是复杂流程写起来没那么顺手。


只用 CompletableFuture 可以吗

代码上可以这样写:

CompletableFuture.supplyAsync(() -> {
    return processPdf(file);
});

但因为没有传线程池,它会使用默认公共线程池。

如果只是一个简单示例,问题不大。

如果是 PDF 水印这种比较重的业务任务,我不会这样做。

因为我无法通过这段代码明确控制:

PDF 同时最多处理几个;

能排队多少任务;

队列满了怎么办;

线程叫什么名字;

这个业务会不会影响其他异步任务。

所以真实项目里,我更喜欢:

CompletableFuture.supplyAsync(
        () -> processPdf(file),
        pdfExecutor
);

也就是:

CompletableFuture 负责任务流程;

ThreadPoolExecutor 负责执行边界。

这两者组合使用更合适。


Future 和 CompletableFuture 是什么关系

CompletableFuture 实现了 Future 接口。

所以它也具备一些 Future 的能力:

future.get();

future.cancel(true);

future.isDone();

future.isCancelled();

同时,它又提供了更多方法:

thenApply(...);

thenCompose(...);

thenCombine(...);

allOf(...);

anyOf(...);

exceptionally(...);

handle(...);

orTimeout(...);

可以简单理解成:

CompletableFuture 包含了 Future 的一部分能力;

又增加了异步任务编排能力。

但这不代表 Future 完全没有价值。

如果需求只是:

提交任务;

稍后获取结果。

Future 反而更简单。


一个简单任务怎么选

假设我要异步处理一个 PDF,最后拿到输出路径。

没有后续异步步骤,也不需要合并其他任务。

可以直接使用:

Future<String> future =
        pdfExecutor.submit(() -> {
            return processPdf(file);
        });

String targetPath = future.get();

这段代码足够直接。

没有必要一定改成:

CompletableFuture<String> future =
        CompletableFuture.supplyAsync(
                () -> processPdf(file),
                pdfExecutor
        );

String targetPath = future.join();

两种都能完成任务。

如果只是单个任务、单个结果,我觉得 Future 没什么问题。


有后续流程时怎么选

假设 PDF 处理完成以后,还要:

上传文件;

生成下载地址;

保存任务记录;

通知用户。

这时候使用 Future,可能会写成:

Future<String> pdfFuture =
        pdfExecutor.submit(() -> {
            return processPdf(file);
        });

String targetPath =
        pdfFuture.get();

Future<String> uploadFuture =
        uploadExecutor.submit(() -> {
            return uploadFile(targetPath);
        });

String downloadUrl =
        uploadFuture.get();

saveDatabase(downloadUrl);

notifyUser(downloadUrl);

这段代码能运行。

但每一步都需要手动:

提交任务;

保存 Future;

调用 get;

拿到结果;

再创建下一个任务。

使用 CompletableFuture 会更自然:

CompletableFuture<PdfTaskResult> future =
        CompletableFuture
                .supplyAsync(
                        () -> processPdf(file),
                        pdfExecutor
                )
                .thenCompose(targetPath -> {
                    return uploadAsync(
                            targetPath,
                            uploadExecutor
                    );
                })
                .thenApply(downloadUrl -> {
                    saveDatabase(downloadUrl);

                    return PdfTaskResult.success(
                            file.getName(),
                            downloadUrl
                    );
                })
                .whenComplete((result, ex) -> {
                    notifyTaskFinished(
                            file.getName(),
                            result,
                            ex
                    );
                })
                .exceptionally(ex -> {
                    return PdfTaskResult.fail(
                            file.getName(),
                            getErrorMessage(ex)
                    );
                });

这时 CompletableFuture 的优势就明显了。


批量任务怎么选

假设有一批 PDF,需要:

同时处理;

等待全部结束;

统计成功和失败。

如果使用 Future,可以这样写:

List<Future<PdfTaskResult>> futureList =
        new ArrayList<>();

for (File file : files) {

    Future<PdfTaskResult> future =
            pdfExecutor.submit(() -> {
                return processPdf(file);
            });

    futureList.add(future);
}

List<PdfTaskResult> resultList =
        new ArrayList<>();

for (Future<PdfTaskResult> future
        : futureList) {

    try {
        resultList.add(
                future.get()
        );
    } catch (Exception e) {
        // 处理异常
    }
}

这套写法没问题。

结构也不算复杂。

如果还要增加:

单任务超时;

统一异常转换;

allOf 等待;

后续汇总流程。

可以使用 CompletableFuture

List<CompletableFuture<PdfTaskResult>>
        futureList = new ArrayList<>();

for (File file : files) {

    CompletableFuture<PdfTaskResult> future =
            CompletableFuture
                    .supplyAsync(
                            () -> processPdf(file),
                            pdfExecutor
                    )
                    .orTimeout(
                            30,
                            TimeUnit.SECONDS
                    )
                    .exceptionally(ex -> {
                        return PdfTaskResult.fail(
                                file.getName(),
                                getErrorMessage(ex)
                        );
                    });

    futureList.add(future);
}

CompletableFuture.allOf(
        futureList.toArray(
                new CompletableFuture[0]
        )
).join();

List<PdfTaskResult> resultList =
        futureList.stream()
                .map(
                        CompletableFuture::join
                )
                .collect(
                        Collectors.toList()
                );

两种方案都可以。

区别是 CompletableFuture 更方便继续扩展异步流程。


ThreadPoolExecutor 和 CompletableFuture 谁控制并发

控制并发数量的是线程池。

例如:

ThreadPoolExecutor pdfExecutor =
        new ThreadPoolExecutor(
                3,
                3,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(100),
                new NamedThreadFactory(
                        "pdf-worker"
                ),
                new ThreadPoolExecutor.AbortPolicy()
        );

这里决定:

同一时间最多有 3 个线程处理 PDF。

而:

CompletableFuture.allOf(...)

只是等待任务完成。

它不会把并发数量限制成 3。

即使创建了 1000 个 CompletableFuture,最终能同时执行多少个,还是由执行它们的线程池决定。

所以不能把:

CompletableFuture

理解成线程池。

也不能把:

allOf

理解成并发控制器。


Future.get 和 CompletableFuture.join 怎么选

Future 一般使用:

future.get();

CompletableFuture 可以使用:

future.get();

也可以使用:

future.join();

区别主要在异常形式。

get() 需要处理受检异常:

try {
    String result = future.get();
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
} catch (ExecutionException e) {
    // 处理任务异常
}

join() 不强制处理受检异常:

String result = future.join();

任务失败时通常抛出:

CompletionException

CompletableFuture 链式代码和 Stream 中,join() 写起来更方便。

例如:

List<PdfTaskResult> resultList =
        futureList.stream()
                .map(
                        CompletableFuture::join
                )
                .collect(
                        Collectors.toList()
                );

如果代码本身需要明确处理中断,我会更认真考虑 get()

如果是在完整的 CompletableFuture 流程末尾获取结果,常用 join()


不要为了异步而异步

看到 CompletableFuture 的链式写法以后,很容易把所有业务都改成异步。

比如一个简单方法:

public String buildTargetPath(
        String fileName
) {
    return "output/" + fileName;
}

没有必要写成:

public CompletableFuture<String> buildTargetPathAsync(
        String fileName
) {
    return CompletableFuture.supplyAsync(() -> {
        return "output/" + fileName;
    });
}

这里只是字符串拼接。

把它异步化不会更快,反而增加:

任务创建;

线程调度;

代码复杂度;

异常处理成本。

异步更适合:

耗时任务;

阻塞任务;

互相独立、可以并行的任务;

确实需要后台执行的任务。

简单计算和对象组装,直接同步执行通常更合适。


什么时候只用 ThreadPoolExecutor

我会在这些场景中直接使用线程池:

只需要提交 Runnable,不关心返回结果;

任务逻辑简单,没有后续异步编排;

希望通过 execute 明确看到任务执行;

需要直接使用 submit + Future 获取简单结果;

正在学习线程池本身的工作原理。

例如:

pdfExecutor.execute(() -> {
    clearTemporaryFile(file);
});

或者:

Future<String> future =
        pdfExecutor.submit(() -> {
            return processPdf(file);
        });

已经够用。


什么时候使用 Future

我会在这些场景中使用 Future

只有一个或少量异步任务;

只需要在后面获取一次结果;

没有复杂任务依赖;

不需要链式异常处理和结果转换;

希望代码简单直接。

例如:

Future<PdfTaskResult> future =
        pdfExecutor.submit(() -> {
            return processPdf(file);
        });

PdfTaskResult result =
        future.get(
                30,
                TimeUnit.SECONDS
        );

这段代码非常清楚。

没必要因为 CompletableFuture 方法更多,就认为 Future 不能用。


什么时候使用 CompletableFuture

我会在这些场景中使用 CompletableFuture

任务完成后还要继续执行下一步;

多个异步任务存在依赖关系;

两个独立任务需要合并结果;

需要等待一批任务全部完成;

需要第一个完成的任务结果;

需要统一异常兜底;

需要方便地增加超时控制;

需要把整个异步流程写成一条链。

例如:

CompletableFuture<PdfTaskResult> future =
        CompletableFuture
                .supplyAsync(
                        () -> addWatermark(file),
                        pdfExecutor
                )
                .thenCompose(path -> {
                    return uploadAsync(
                            path,
                            uploadExecutor
                    );
                })
                .thenApply(url -> {
                    return PdfTaskResult.success(
                            file.getName(),
                            url
                    );
                })
                .orTimeout(
                        60,
                        TimeUnit.SECONDS
                )
                .exceptionally(ex -> {
                    return PdfTaskResult.fail(
                            file.getName(),
                            getErrorMessage(ex)
                    );
                });

这就是 CompletableFuture 擅长的地方。


实际项目中的常见组合

我觉得真实项目里比较常见的不是三选一,而是:

ThreadPoolExecutor + CompletableFuture

例如:

ThreadPoolExecutor pdfExecutor =
        createPdfExecutor();

CompletableFuture<PdfTaskResult> future =
        CompletableFuture
                .supplyAsync(
                        () -> processPdf(file),
                        pdfExecutor
                )
                .orTimeout(
                        30,
                        TimeUnit.SECONDS
                )
                .exceptionally(ex -> {
                    return PdfTaskResult.fail(
                            file.getName(),
                            getErrorMessage(ex)
                    );
                });

这里:

ThreadPoolExecutor 控制执行资源;

CompletableFuture 组织任务流程和结果。

它们各做各的事情。

这比只用默认公共线程池更可控。


一张表整理三者区别

工具 主要职责 是否管理线程 是否表示任务结果 是否方便编排流程
ThreadPoolExecutor 管理线程、队列和拒绝策略 间接支持
Future 表示将来的任务结果 不方便
CompletableFuture 表示结果并编排异步流程

其中“CompletableFuture 不管理线程”是指:

它不是线程池;

它仍然需要执行器或公共线程池来运行任务。

我的选择顺序

现在让我选择时,我会先问几个问题。

任务是否需要后台线程执行

如果不需要,就直接同步调用。

不要先考虑 FutureCompletableFuture


是否需要控制线程数量和任务队列

如果需要,就先准备:

ThreadPoolExecutor

或者使用由 Spring 统一管理的线程池。


是否需要返回结果

不需要结果:

executor.execute(...)

需要一个简单结果:

~~~java
executor.submit(…)