跳到正文
hello world

31. allOf:等待一批异步任务全部完成

发布于阅读量 0

31. allOf:等待一批异步任务全部完成

上一节用 join() 等待了一个 CompletableFuture 完成。

如果只有一个 PDF 任务,这样写没问题:

String result = future.join();

但批量 PDF 处理时,通常不是一个任务,而是一批任务。

比如:

a.pdf
b.pdf
c.pdf
d.pdf
e.pdf

每个 PDF 都对应一个 CompletableFuture

这时候我就需要一个更明确的表达:

等这一批异步任务全部完成以后,再继续往下执行。

这就是 CompletableFuture.allOf() 的作用。


allOf 是干什么的

allOf() 的基本写法是:

CompletableFuture.allOf(future1, future2, future3).join();

它的意思是:

等待 future1、future2、future3 全部完成。

注意,是全部完成。

只要还有一个任务没结束,allOf(...).join() 就会继续等待。

这个方法很适合批量任务,比如:

批量处理 PDF;
批量调用接口;
批量查询数据;
批量生成报表;
批量上传文件。

先写一个简单例子

新建类:

com.succos.completablefuture.AllOfDemo

代码如下:

package com.succos.completablefuture;

import java.util.concurrent.CompletableFuture;

public class AllOfDemo {

    public static void main(String[] args) {

        CompletableFuture<String> future1 = CompletableFuture.supplyAsync(() -> {
            System.out.println("任务1开始:" + Thread.currentThread().getName());
            sleep(3000);
            System.out.println("任务1完成:" + Thread.currentThread().getName());
            return "a.pdf 处理完成";
        });

        CompletableFuture<String> future2 = CompletableFuture.supplyAsync(() -> {
            System.out.println("任务2开始:" + Thread.currentThread().getName());
            sleep(5000);
            System.out.println("任务2完成:" + Thread.currentThread().getName());
            return "b.pdf 处理完成";
        });

        CompletableFuture<String> future3 = CompletableFuture.supplyAsync(() -> {
            System.out.println("任务3开始:" + Thread.currentThread().getName());
            sleep(2000);
            System.out.println("任务3完成:" + Thread.currentThread().getName());
            return "c.pdf 处理完成";
        });

        System.out.println("main 已经提交 3 个异步任务");

        CompletableFuture<Void> allFuture = CompletableFuture.allOf(
                future1,
                future2,
                future3
        );

        allFuture.join();

        System.out.println("3 个任务全部完成");

        System.out.println(future1.join());
        System.out.println(future2.join());
        System.out.println(future3.join());
    }

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

这个例子里,三个任务耗时不一样:

任务1:3 秒;
任务2:5 秒;
任务3:2 秒。

allOf().join() 会等最慢的任务也完成。

所以大概要等 5 秒以后,才会打印:

3 个任务全部完成

allOf 返回的是 CompletableFuture

这里要注意一个点:

CompletableFuture<Void> allFuture = CompletableFuture.allOf(
        future1,
        future2,
        future3
);

allOf() 返回的是:

CompletableFuture<Void>

也就是说,它本身不直接返回每个任务的结果。

它只是告诉我:

这一批任务是否都完成了。

如果我要拿每个任务的结果,还是要对原来的 future1future2future3 调用 join()

所以代码里才会写:

System.out.println(future1.join());
System.out.println(future2.join());
System.out.println(future3.join());

这一点很重要。

allOf() 负责等全部完成,不负责把所有结果自动合成一个集合。


为什么 allOf 后再 join 单个任务不会阻塞太久

这段代码看起来好像又对每个 future 调用了一次 join()

allFuture.join();

System.out.println(future1.join());
System.out.println(future2.join());
System.out.println(future3.join());

但这里不会再重复等待很久。

因为前面:

allFuture.join();

已经确认所有任务都完成了。

后面再调用:

future1.join()

只是直接拿已经完成的结果。

所以这个写法是可以接受的。


批量任务里怎么使用 allOf

实际批量任务里,不可能手动写 future1future2future3

通常会有一个集合:

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

然后遍历 PDF 文件,每个文件创建一个 CompletableFuture,放进集合。

最后用:

CompletableFuture<Void> allFuture = CompletableFuture.allOf(
        futureList.toArray(new CompletableFuture[0])
);

allFuture.join();

这里的:

futureList.toArray(new CompletableFuture[0])

是因为 allOf() 接收的是数组参数:

CompletableFuture<?>... cfs

所以要把 List 转成数组。


批量 PDF 示例

新建类:

com.succos.completablefuture.AllOfPdfDemo

代码如下:

package com.succos.completablefuture;

import java.io.File;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;

public class AllOfPdfDemo {

    public static void main(String[] args) {

        File inputDir = new File("input");

        File[] files = inputDir.listFiles(file ->
                file.isFile() && file.getName().toLowerCase().endsWith(".pdf")
        );

        if (files == null || files.length == 0) {
            System.out.println("input 目录下没有 PDF 文件");
            return;
        }

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

        long start = System.currentTimeMillis();

        for (File file : files) {

            CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
                System.out.println(Thread.currentThread().getName()
                        + " 开始处理:"
                        + file.getName());

                sleep(3000);

                String targetPath = "output/"
                        + file.getName().replace(".pdf", "-watermark.pdf");

                System.out.println(Thread.currentThread().getName()
                        + " 处理完成:"
                        + file.getName());

                return targetPath;
            });

            futureList.add(future);
        }

        CompletableFuture<Void> allFuture = CompletableFuture.allOf(
                futureList.toArray(new CompletableFuture[0])
        );

        allFuture.join();

        System.out.println("全部异步任务已经完成,开始收集结果");

        for (CompletableFuture<String> future : futureList) {
            String targetPath = future.join();
            System.out.println("输出路径:" + targetPath);
        }

        long end = System.currentTimeMillis();

        System.out.println("--------------------------------");
        System.out.println("全部 PDF 处理完成");
        System.out.println("总文件数:" + files.length);
        System.out.println("总耗时:" + (end - start) + " ms");
    }

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

这个版本比上一节一个个遍历 join() 更直观。

因为它明确表达了:

先等待这一批任务全部完成;
再统一收集结果。

allOf 和循环 join 的区别

其实下面这种写法也能等所有任务完成:

for (CompletableFuture<String> future : futureList) {
    future.join();
}

它不是不能用。

allOf() 表达得更清楚:

CompletableFuture.allOf(...).join();

这行代码一看就知道:

等待全部任务完成。

而循环 join() 更像是在一个个等待。

在简单场景里两种都可以。

但如果后面还要接下一步流程,allOf() 会更自然。

比如:

CompletableFuture.allOf(
        futureList.toArray(new CompletableFuture[0])
).thenRun(() -> {
    System.out.println("所有 PDF 都处理完了,开始汇总");
});

这就很像异步流程编排。


allOf 后面接 thenRun

allOf() 返回的是 CompletableFuture<Void>,所以后面可以接 thenRun()

比如:

CompletableFuture<Void> allFuture = CompletableFuture.allOf(
        futureList.toArray(new CompletableFuture[0])
);

allFuture.thenRun(() -> {
    System.out.println("所有任务都完成了,执行汇总逻辑");
}).join();

这里的 thenRun() 不接收上一步结果。

因为 allOf() 本来就没有业务返回值。

它只是表示:

等所有任务完成后,再执行这段代码。

这个写法比手动等完再写一堆逻辑更像一条流程。


allOf 不会自动处理每个结果

这个地方再强调一次。

allOf() 不会把这些结果:

a.pdf 的输出路径
b.pdf 的输出路径
c.pdf 的输出路径

自动组合成一个 List<String>

如果想要 List<String>,要自己收集。

比如:

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

如果你的 Java 版本低于 16,没有 toList(),可以写成:

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

这样才真正把所有任务结果变成一个集合。


用 allOf 收集 PDF 输出路径

可以这样写:

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

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

for (String targetPath : targetPathList) {
    System.out.println("输出路径:" + targetPath);
}

这里的顺序是:

先等所有任务完成;
再把每个 future 的结果取出来;
最后得到一个输出路径列表。

这个写法适合后面做汇总。

比如:

生成下载清单;
保存数据库;
返回给前端;
统计成功数量。

allOf 遇到异常会怎么样

如果其中某个 CompletableFuture 异常了,allOf().join() 也会异常。

比如:

CompletableFuture<String> future1 = CompletableFuture.supplyAsync(() -> {
    return "a.pdf 成功";
});

CompletableFuture<String> future2 = CompletableFuture.supplyAsync(() -> {
    int x = 1 / 0;
    return "b.pdf 成功:" + x;
});

CompletableFuture<Void> allFuture = CompletableFuture.allOf(future1, future2);

allFuture.join();

这里 future2 会失败。

那么:

allFuture.join();

也会抛出 CompletionException

这说明 allOf() 不是“忽略失败,只等结束”。

它会感知到其中有任务异常完成。


批量任务里更推荐返回结果对象

如果我不想让一个 PDF 的异常影响整个 allOf().join(),可以在每个任务内部捕获异常,并返回一个失败结果对象。

比如:

CompletableFuture<PdfTaskResult> future = CompletableFuture.supplyAsync(() -> {
    try {
        String targetPath = addWatermark(file);
        return PdfTaskResult.success(file.getName(), targetPath);
    } catch (Exception e) {
        return PdfTaskResult.fail(file.getName(), e.getMessage());
    }
});

这样对 CompletableFuture 来说,任务是正常完成的。

只是返回的结果对象里标记了失败。

后面汇总时就可以统计:

成功多少个;
失败多少个;
失败原因分别是什么。

这个比让异常直接打断整个 allOf() 更适合批量文件处理。


PDF 结果对象版本

可以定义一个简单的结果对象:

static class PdfTaskResult {

    private boolean success;

    private String fileName;

    private String targetPath;

    private String message;

    public static PdfTaskResult success(String fileName, String targetPath) {
        PdfTaskResult result = new PdfTaskResult();
        result.success = true;
        result.fileName = fileName;
        result.targetPath = targetPath;
        result.message = "处理成功";
        return result;
    }

    public static PdfTaskResult fail(String fileName, String message) {
        PdfTaskResult result = new PdfTaskResult();
        result.success = false;
        result.fileName = fileName;
        result.message = message;
        return result;
    }
}

然后任务返回:

CompletableFuture<PdfTaskResult> future = CompletableFuture.supplyAsync(() -> {
    try {
        sleep(3000);

        String targetPath = "output/"
                + file.getName().replace(".pdf", "-watermark.pdf");

        return PdfTaskResult.success(file.getName(), targetPath);

    } catch (Exception e) {
        return PdfTaskResult.fail(file.getName(), e.getMessage());
    }
});

最后统一收集:

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

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

这样结果会比较规整。


allOf 不是控制并发数量

还有一点要分清楚。

allOf() 只是等待一批任务全部完成。

它不负责控制同时执行多少个任务。

如果我这样写:

for (File file : files) {
    CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
        return processPdf(file);
    });
}

又没有传自定义线程池,那任务会使用默认公共线程池。

allOf() 并不会限制“最多同时处理 3 个 PDF”。

如果要控制并发数量,应该通过自定义线程池来控制:

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

也就是说:

线程池:控制任务在哪里执行、最多多少线程执行;
allOf:等待这一批任务全部完成。

这两个解决的是不同问题。


这一节小结

这一节我主要记住几点:

1. allOf 用来等待一批 CompletableFuture 全部完成;
2. allOf 返回 CompletableFuture<Void>,本身不直接返回每个任务结果;
3. 要拿结果,还是要对原来的 future 调用 join;
4. allOf 后再 join 单个 future,一般不会重复阻塞很久;
5. allOf 可以接 thenRun,表达“全部完成后再做某事”;
6. 如果某个任务异常,allOf().join() 也可能抛 CompletionException;
7. 批量 PDF 更适合每个任务返回结果对象,方便统一汇总;
8. allOf 只负责等待,不负责控制并发数量。

用一句话总结:

allOf 解决的是“这一批异步任务什么时候全部结束”的问题,而不是“每个任务结果是什么”的问题。

下一节继续看 CompletableFuture 使用自定义线程池。

因为 PDF 处理这种任务不能一直丢给默认公共线程池,应该放到自己的 PDF 线程池里执行。