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 处理的是另一个场景:第二个异步任务必须依赖第一个异步任务的结果。