跳到正文
hello world

36. thenCombine:两个异步任务结果怎么合并

发布于阅读量 0

36. thenCombine:两个异步任务结果怎么合并

前面写的异步流程,大多是一条单线:

处理 PDF
↓
生成下载地址
↓
保存结果

后一步依赖前一步的结果,所以可以一直使用:

thenApply(...)
thenAccept(...)

但有些场景不是前后依赖,而是两个任务可以同时执行。

比如生成 PDF 水印结果时,我还需要读取一份用户信息:

任务一:处理 PDF,返回输出路径;

任务二:查询用户信息,返回用户名。

这两个任务之间没有依赖关系。

处理 PDF 不需要先知道用户名,查询用户信息也不需要等待 PDF 处理完成。

所以可以让它们同时执行。

等两个任务都完成以后,再把结果合并成一句话:

张三的 PDF 已处理完成,文件地址为 output/test-watermark.pdf

这种场景就适合使用:

thenCombine(...)

thenCombine 解决什么问题

thenCombine() 解决的是:

两个互相独立的 CompletableFuture 同时执行;

等两个任务都完成以后;

把两个任务的结果合并成一个新结果。

它的基本写法是:

CompletableFuture<String> resultFuture =
        future1.thenCombine(future2, (result1, result2) -> {
            return result1 + result2;
        });

这里有两个异步任务:

future1
future2

等它们都完成以后,thenCombine() 会拿到:

result1:future1 的结果;

result2:future2 的结果。

然后返回一个新的结果。


一个最简单的 thenCombine 示例

新建类:

com.succos.completablefuture.ThenCombineDemo

代码如下:

package com.succos.completablefuture;

import java.util.concurrent.CompletableFuture;

public class ThenCombineDemo {

    public static void main(String[] args) {

        CompletableFuture<String> future1 =
                CompletableFuture.supplyAsync(() -> {

                    System.out.println("任务一开始:"
                            + Thread.currentThread().getName());

                    sleep(3000);

                    System.out.println("任务一完成");

                    return "PDF 处理成功";
                });

        CompletableFuture<String> future2 =
                CompletableFuture.supplyAsync(() -> {

                    System.out.println("任务二开始:"
                            + Thread.currentThread().getName());

                    sleep(2000);

                    System.out.println("任务二完成");

                    return "用户:张三";
                });

        CompletableFuture<String> resultFuture =
                future1.thenCombine(
                        future2,
                        (pdfResult, userResult) -> {

                            System.out.println("开始合并结果:"
                                    + Thread.currentThread().getName());

                            return userResult + "," + pdfResult;
                        }
                );

        String result = resultFuture.join();

        System.out.println("最终结果:" + result);
    }

    private static void sleep(long millis) {

        try {
            Thread.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("任务被中断", e);
        }
    }
}

运行结果可能类似:

任务一开始:ForkJoinPool.commonPool-worker-1
任务二开始:ForkJoinPool.commonPool-worker-2
任务二完成
任务一完成
开始合并结果:ForkJoinPool.commonPool-worker-1
最终结果:用户:张三,PDF 处理成功

可以看到,任务一和任务二是同时执行的。

thenCombine() 不会在其中一个任务完成后马上执行。

它会等待两个任务都完成。


thenCombine 的两个任务谁先完成都可以

上面的例子中:

任务一耗时 3 秒;

任务二耗时 2 秒。

任务二会先完成。

thenCombine() 还是要等任务一。

因为它需要同时拿到两个任务的结果:

pdfResult
userResult

所以总耗时大约取决于较慢的那个任务。

这里大约是 3 秒,而不是:

3 秒 + 2 秒 = 5 秒

因为两个任务是并发执行的。

可以简单理解成:

串行执行时间:任务一耗时 + 任务二耗时;

并行执行时间:两个任务中较长的耗时。

当然,还要加上一点线程调度和结果合并的开销。


用 PDF 项目理解 thenCombine

假设我现在要生成一份 PDF 处理结果说明。

需要两个数据:

PDF 处理结果;

用户姓名。

两个任务可以分别执行:

CompletableFuture<String> pdfFuture =
        CompletableFuture.supplyAsync(() -> {
            return addWatermark(file);
        }, pdfExecutor);
CompletableFuture<String> userFuture =
        CompletableFuture.supplyAsync(() -> {
            return queryUserName(userId);
        }, queryExecutor);

然后合并:

CompletableFuture<String> resultFuture =
        pdfFuture.thenCombine(
                userFuture,
                (targetPath, userName) -> {
                    return userName
                            + " 的 PDF 已处理完成,输出路径:"
                            + targetPath;
                }
        );

最后:

String result = resultFuture.join();

这样就得到了完整结果。


完整的 PDF 合并示例

新建类:

com.succos.completablefuture.ThenCombinePdfDemo

代码如下:

package com.succos.completablefuture;

import java.io.File;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ThreadFactory;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public class ThenCombinePdfDemo {

    public static void main(String[] args) {

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

        ThreadPoolExecutor queryExecutor = new ThreadPoolExecutor(
                2,
                2,
                60,
                TimeUnit.SECONDS,
                new ArrayBlockingQueue<>(100),
                new NamedThreadFactory("query-worker"),
                new ThreadPoolExecutor.CallerRunsPolicy()
        );

        File file = new File("input/test.pdf");

        CompletableFuture<String> pdfFuture =
                CompletableFuture.supplyAsync(() -> {

                    System.out.println(Thread.currentThread().getName()
                            + " 开始处理 PDF");

                    sleep(3000);

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

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

                    return targetPath;

                }, pdfExecutor);

        CompletableFuture<String> userFuture =
                CompletableFuture.supplyAsync(() -> {

                    System.out.println(Thread.currentThread().getName()
                            + " 开始查询用户");

                    sleep(2000);

                    System.out.println(Thread.currentThread().getName()
                            + " 用户查询完成");

                    return "张三";

                }, queryExecutor);

        CompletableFuture<PdfMessage> resultFuture =
                pdfFuture.thenCombine(
                        userFuture,
                        (targetPath, userName) -> {

                            PdfMessage message = new PdfMessage();

                            message.setUserName(userName);
                            message.setTargetPath(targetPath);
                            message.setMessage(
                                    userName
                                            + " 的 PDF 已处理完成"
                            );

                            return message;
                        }
                );

        PdfMessage result = resultFuture.join();

        System.out.println("--------------------------------");
        System.out.println("用户:" + result.getUserName());
        System.out.println("文件路径:" + result.getTargetPath());
        System.out.println("处理信息:" + result.getMessage());

        pdfExecutor.shutdown();
        queryExecutor.shutdown();
    }

    private static void sleep(long millis) {

        try {
            Thread.sleep(millis);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new RuntimeException("任务被中断", e);
        }
    }

    static class PdfMessage {

        private String userName;

        private String targetPath;

        private String message;

        public String getUserName() {
            return userName;
        }

        public void setUserName(String userName) {
            this.userName = userName;
        }

        public String getTargetPath() {
            return targetPath;
        }

        public void setTargetPath(String targetPath) {
            this.targetPath = targetPath;
        }

        public String getMessage() {
            return message;
        }

        public void setMessage(String message) {
            this.message = message;
        }
    }

    static class NamedThreadFactory implements ThreadFactory {

        private final String prefix;

        private final AtomicInteger threadNumber =
                new AtomicInteger(1);

        public NamedThreadFactory(String prefix) {
            this.prefix = prefix;
        }

        @Override
        public Thread newThread(Runnable r) {

            Thread thread = new Thread(r);

            thread.setName(
                    prefix
                            + "-"
                            + threadNumber.getAndIncrement()
            );

            return thread;
        }
    }
}

这个例子里:

pdfFuture 使用 PDF 线程池;

userFuture 使用查询线程池;

两个任务同时执行;

thenCombine 等两个任务都结束后组装 PdfMessage。

这就是比较典型的多任务结果合并。


thenCombine 和 thenApply 的区别

thenApply() 处理的是一个前后依赖关系。

例如:

CompletableFuture<String> future =
        CompletableFuture
                .supplyAsync(() -> addWatermark(file))
                .thenApply(targetPath -> {
                    return buildDownloadUrl(targetPath);
                });

这里必须先得到 targetPath,才能生成下载地址。

所以是:

任务一完成
↓
任务二才能开始

thenCombine() 处理的是两个独立任务:

pdfFuture.thenCombine(userFuture, (path, user) -> {
    return user + ":" + path;
});

这里 PDF 处理和用户查询可以同时开始。

所以是:

任务一 ─┐
        ├─ 两个都完成以后合并
任务二 ─┘

可以这样区分:

后一步依赖前一步:thenApply;

两个任务互不依赖,最后合并:thenCombine。

thenCombine 和 allOf 的区别

allOf() 也能等待多个任务完成。

例如:

CompletableFuture.allOf(
        pdfFuture,
        userFuture
).join();

但它只负责等待,不直接帮我合并结果。

后面还要自己取:

String targetPath = pdfFuture.join();
String userName = userFuture.join();

再手动组装。

thenCombine() 可以直接把两个结果传进 Lambda:

pdfFuture.thenCombine(
        userFuture,
        (targetPath, userName) -> {
            return userName + ":" + targetPath;
        }
);

所以:

等待多个任务完成:allOf;

合并两个任务结果:thenCombine。

如果只有两个结果需要合并,thenCombine() 通常更自然。

如果有十几个任务,allOf() 更适合统一等待。


thenCombine 返回什么类型

假设:

CompletableFuture<String> future1
CompletableFuture<Integer> future2

合并时可以返回任何类型。

例如:

CompletableFuture<Boolean> resultFuture =
        future1.thenCombine(
                future2,
                (name, count) -> {
                    return name.length() > count;
                }
        );

所以 thenCombine() 的最终类型,取决于合并函数返回什么。

PDF 例子里:

CompletableFuture<String> pdfFuture
CompletableFuture<String> userFuture

不代表合并结果也必须是 String

完全可以返回:

CompletableFuture<PdfMessage>

这也是为什么真实项目里更推荐返回结果对象,而不是不断拼字符串。


thenCombineAsync 有什么区别

thenCombine() 也有带 Async 的版本:

thenCombineAsync(...)

以及可以指定线程池的版本:

thenCombineAsync(
        otherFuture,
        (result1, result2) -> {
            return merge(result1, result2);
        },
        executor
)

区别和前面讲的一样。

不带 Async

两个任务都完成后,通常由完成最后一个任务的线程继续执行合并逻辑。

Async

两个任务完成后,把合并逻辑重新提交给执行器调度。

例如:

CompletableFuture<PdfMessage> resultFuture =
        pdfFuture.thenCombineAsync(
                userFuture,
                (targetPath, userName) -> {
                    return buildPdfMessage(
                            targetPath,
                            userName
                    );
                },
                resultExecutor
        );

如果合并逻辑只是创建一个对象、拼几个字段,一般没必要加 Async

直接用:

thenCombine(...)

就够了。

如果合并之后还有比较重的计算,才考虑 thenCombineAsync()


thenCombine 的合并逻辑不要太重

例如:

pdfFuture.thenCombine(
        userFuture,
        (targetPath, userName) -> {
            return new PdfMessage(
                    userName,
                    targetPath
            );
        }
);

这里只是组装对象,很轻。

不需要专门切换线程。

但如果合并逻辑里还要:

压缩文件;
上传远程存储;
执行复杂计算;
调用外部接口。

那它已经不只是“合并结果”了。

这时候可以考虑:

thenCombineAsync(...)

或者先合并成一个参数对象,再用后续的异步步骤处理。

例如:

pdfFuture
        .thenCombine(
                userFuture,
                PdfContext::new
        )
        .thenApplyAsync(
                context -> upload(context),
                uploadExecutor
        );

这样职责会更清楚。


如果其中一个任务失败会怎么样

假设 PDF 处理任务失败:

CompletableFuture<String> pdfFuture =
        CompletableFuture.supplyAsync(() -> {
            throw new RuntimeException("PDF 处理失败");
        });

用户查询任务成功:

CompletableFuture<String> userFuture =
        CompletableFuture.supplyAsync(() -> "张三");

再执行:

CompletableFuture<String> resultFuture =
        pdfFuture.thenCombine(
                userFuture,
                (path, userName) -> {
                    return userName + ":" + path;
                }
        );

因为 pdfFuture 没有正常结果,所以合并函数不会正常执行。

最后调用:

resultFuture.join();

会抛出 CompletionException

也就是说:

thenCombine 要求两个任务都正常完成,才能执行合并逻辑。

只要其中一个任务异常,最终结果就会异常完成。


给每个任务单独做异常兜底

如果希望某个任务失败后还能继续合并,可以先给任务兜底。

例如:

CompletableFuture<String> pdfFuture =
        CompletableFuture
                .supplyAsync(() -> addWatermark(file), pdfExecutor)
                .exceptionally(ex -> "PDF 处理失败");
CompletableFuture<String> userFuture =
        CompletableFuture
                .supplyAsync(() -> queryUserName(userId), queryExecutor)
                .exceptionally(ex -> "未知用户");

然后再合并:

CompletableFuture<String> resultFuture =
        pdfFuture.thenCombine(
                userFuture,
                (pdfResult, userName) -> {
                    return userName + "," + pdfResult;
                }
        );

这样即使某一个任务失败,也有一个兜底结果参与合并。

不过这种写法是否合理,要看业务。

如果 PDF 处理失败,可能就不应该继续生成成功提示。

所以更稳的方式通常是返回结果对象。


使用结果对象处理成功和失败

可以定义两个结果:

class PdfTaskResult {

    private boolean success;

    private String targetPath;

    private String message;
}
class UserQueryResult {

    private boolean success;

    private String userName;

    private String message;
}

然后两个异步任务都不直接抛异常,而是返回明确结果。

合并时判断:

CompletableFuture<PdfMessage> resultFuture =
        pdfFuture.thenCombine(
                userFuture,
                (pdfResult, userResult) -> {

                    if (!pdfResult.isSuccess()) {
                        return PdfMessage.fail(
                                pdfResult.getMessage()
                        );
                    }

                    if (!userResult.isSuccess()) {
                        return PdfMessage.fail(
                                userResult.getMessage()
                        );
                    }

                    return PdfMessage.success(
                            userResult.getUserName(),
                            pdfResult.getTargetPath()
                    );
                }
        );

这样业务状态会更清楚。

不会把:

未知用户,PDF 处理失败

简单拼成一个看似正常的结果。


thenCombine 不适合有依赖关系的任务

假设上传 PDF 必须依赖水印处理后的路径:

先添加水印;

再把水印 PDF 上传到服务器。

这种情况不能为了并发而使用 thenCombine()

因为上传任务一开始还不知道文件路径。

应该使用:

CompletableFuture
        .supplyAsync(() -> addWatermark(file), pdfExecutor)
        .thenApplyAsync(
                targetPath -> uploadFile(targetPath),
                uploadExecutor
        );

这里两个任务存在明确依赖:

没有水印文件,就不能上传。

所以适合链式调用,而不是并行合并。

我判断时会先问一句:

这两个任务能不能在同一时间开始?

如果能,最后还要合并结果,可以用 thenCombine()

如果不能,后一个必须等前一个结果,就用 thenApply()thenCompose() 等链式方法。


thenCombine 的典型场景

我觉得它比较适合这些业务:

同时查询用户信息和订单信息,最后组装页面数据;

同时查询产品数据和库存数据,最后生成产品详情;

同时处理 PDF 和查询任务配置,最后生成任务结果;

同时查询两个不同接口,最后合并返回值;

同时计算两个独立指标,最后生成统计结果。

共同特点都是:

两个任务互不依赖;

可以同时执行;

最后需要两个任务的结果。

这一节小结

这一节我主要记住几点:

1. thenCombine 用来合并两个独立异步任务的结果;
2. 两个任务可以同时执行,不需要互相等待;
3. thenCombine 会等两个任务都完成后再执行合并逻辑;
4. 总耗时通常取决于两个任务中较慢的一个;
5. thenApply 适合前后依赖,thenCombine 适合并行后合并;
6. allOf 主要负责等待,thenCombine 可以直接拿到两个任务的结果;
7. 两个任务中只要有一个异常,合并结果通常也会异常;
8. 合并逻辑很轻时用 thenCombine,逻辑较重或需要线程隔离时考虑 thenCombineAsync。

用一句话总结:

两个任务可以同时做,最后又需要把两个结果放在一起,就用 thenCombine。

下一节继续看 thenCompose

thenCombine 处理的是两个互相独立的任务,而 thenCompose 处理的是另一个场景:第二个异步任务必须依赖第一个异步任务的结果。