38. anyOf:多个任务谁先完成就用谁
发布于 • 阅读量 0
38. anyOf:多个任务谁先完成就用谁
前面学过:
CompletableFuture.allOf(...)
它的作用是等待一批异步任务全部完成。
这一节看另一个方法:
CompletableFuture.anyOf(...)
它解决的是另一种情况:
我同时启动多个异步任务;
不需要等待所有任务结束;
只要其中任意一个任务先完成,就先使用它的结果。
比如我有三个文件服务,都能返回同一个 PDF 的下载地址:
服务一:阿里云对象存储;
服务二:腾讯云对象存储;
服务三:本地备用存储。
我可以同时请求三个服务。
谁先返回,就先使用谁的结果。
这种场景适合使用:
anyOf(...)
anyOf 的基本写法
先看一个简单例子:
CompletableFuture<String> future1 =
CompletableFuture.supplyAsync(() -> {
sleep(3000);
return "任务一完成";
});
CompletableFuture<String> future2 =
CompletableFuture.supplyAsync(() -> {
sleep(1000);
return "任务二完成";
});
CompletableFuture<String> future3 =
CompletableFuture.supplyAsync(() -> {
sleep(2000);
return "任务三完成";
});
CompletableFuture<Object> anyFuture =
CompletableFuture.anyOf(
future1,
future2,
future3
);
Object result = anyFuture.join();
System.out.println("最先完成的结果:" + result);
三个任务的耗时分别是:
任务一:3 秒;
任务二:1 秒;
任务三:2 秒。
所以 future2 大概率最先完成。
最终输出:
最先完成的结果:任务二完成
完整示例
新建类:
com.succos.completablefuture.AnyOfDemo
代码如下:
package com.succos.completablefuture;
import java.util.concurrent.CompletableFuture;
public class AnyOfDemo {
public static void main(String[] args) {
CompletableFuture<String> future1 =
CompletableFuture.supplyAsync(() -> {
System.out.println("任务一开始:"
+ Thread.currentThread().getName());
sleep(3000);
System.out.println("任务一完成");
return "任务一结果";
});
CompletableFuture<String> future2 =
CompletableFuture.supplyAsync(() -> {
System.out.println("任务二开始:"
+ Thread.currentThread().getName());
sleep(1000);
System.out.println("任务二完成");
return "任务二结果";
});
CompletableFuture<String> future3 =
CompletableFuture.supplyAsync(() -> {
System.out.println("任务三开始:"
+ Thread.currentThread().getName());
sleep(2000);
System.out.println("任务三完成");
return "任务三结果";
});
CompletableFuture<Object> anyFuture =
CompletableFuture.anyOf(
future1,
future2,
future3
);
Object result = anyFuture.join();
System.out.println("--------------------------------");
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-3
任务二完成
最先完成的结果:任务二结果
anyOf().join() 在任务二完成以后就能返回。
它不会继续等待任务一和任务三。
anyOf 返回 CompletableFuture
这里有一个比较明显的地方:
CompletableFuture<Object> anyFuture =
CompletableFuture.anyOf(
future1,
future2,
future3
);
返回类型是:
CompletableFuture<Object>
不是:
CompletableFuture<String>
即使三个任务返回的都是 String,anyOf() 的结果类型仍然是 Object。
所以拿到结果时是:
Object result = anyFuture.join();
如果确定所有任务返回的都是字符串,可以强制转换:
String result = (String) anyFuture.join();
完整写法:
String result = (String) CompletableFuture
.anyOf(
future1,
future2,
future3
)
.join();
不过强制转换要建立在类型确定的前提下。
为什么 anyOf 返回 Object
anyOf() 允许传入不同结果类型的任务。
比如:
CompletableFuture<String> future1 =
CompletableFuture.completedFuture("PDF 处理完成");
CompletableFuture<Integer> future2 =
CompletableFuture.completedFuture(100);
CompletableFuture<Boolean> future3 =
CompletableFuture.completedFuture(true);
这些任务的结果类型分别是:
String;
Integer;
Boolean。
anyOf() 不知道最终最先完成的是哪个类型。
所以只能统一返回:
Object
这也是为什么实际使用时,我更倾向于让传给 anyOf() 的任务返回相同类型。
这样后续处理会简单一些。
anyOf 和 allOf 的区别
这两个方法名字很像,但等待条件完全不同。
allOf():
CompletableFuture.allOf(
future1,
future2,
future3
).join();
表示:
三个任务全部完成以后才继续。
anyOf():
CompletableFuture.anyOf(
future1,
future2,
future3
).join();
表示:
三个任务里只要有一个完成,就可以继续。
可以简单记成:
| 方法 | 等待条件 | 返回结果 |
|---|---|---|
allOf |
所有任务都完成 | CompletableFuture<Void> |
anyOf |
任意一个任务完成 | CompletableFuture<Object> |
PDF 批量水印处理需要等所有文件完成,适合 allOf()。
多个服务返回相同结果,谁快用谁,适合 anyOf()。
PDF 项目里有什么使用场景
普通的批量 PDF 水印处理一般不需要 anyOf()。
因为每一个 PDF 都应该被处理,不能只等其中一个文件完成。
例如有 10 个 PDF:
a.pdf;
b.pdf;
c.pdf;
……
业务要求通常是:
10 个文件全部处理结束以后,再统计结果。
这种情况应该使用:
allOf(...)
而不是:
anyOf(...)
否则只要第一个 PDF 处理完成,等待就结束了,其他文件的结果还没有收集。
场景一:多个存储服务谁先返回就用谁
假设同一个 PDF 同时上传到多个存储服务:
服务一:主对象存储;
服务二:备用对象存储;
服务三:本地文件服务。
我只需要尽快拿到一个可用下载地址。
可以这样写:
CompletableFuture<String> storageFuture1 =
uploadToStorage1(file);
CompletableFuture<String> storageFuture2 =
uploadToStorage2(file);
CompletableFuture<String> storageFuture3 =
uploadToStorage3(file);
String downloadUrl = (String) CompletableFuture
.anyOf(
storageFuture1,
storageFuture2,
storageFuture3
)
.join();
谁先上传完成,就先返回谁的下载地址。
场景二:多个接口查询同一份数据
假设要查询 PDF 的识别结果,系统有两个可用服务:
识别服务 A;
识别服务 B。
两个服务都能完成相同任务,但响应速度不稳定。
可以同时调用:
CompletableFuture<PdfRecognizeResult> future1 =
recognizeByServiceA(file);
CompletableFuture<PdfRecognizeResult> future2 =
recognizeByServiceB(file);
PdfRecognizeResult result =
(PdfRecognizeResult) CompletableFuture
.anyOf(
future1,
future2
)
.join();
哪个服务先返回,就先使用哪个结果。
这种模式有时也叫“竞速请求”。
不过它会增加服务调用次数和成本,不能随便使用。
使用自定义线程池
PDF、上传、网络请求这类任务,最好不要全部丢到默认公共线程池。
可以创建一个专用线程池:
ThreadPoolExecutor storageExecutor =
new ThreadPoolExecutor(
3,
3,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new NamedThreadFactory("storage-worker"),
new ThreadPoolExecutor.CallerRunsPolicy()
);
然后提交任务:
CompletableFuture<String> future1 =
CompletableFuture.supplyAsync(() -> {
return uploadToStorage1(file);
}, storageExecutor);
CompletableFuture<String> future2 =
CompletableFuture.supplyAsync(() -> {
return uploadToStorage2(file);
}, storageExecutor);
CompletableFuture<String> future3 =
CompletableFuture.supplyAsync(() -> {
return uploadToStorage3(file);
}, storageExecutor);
最后:
String result = (String) CompletableFuture
.anyOf(
future1,
future2,
future3
)
.join();
任务结束后记得关闭自己创建的线程池:
storageExecutor.shutdown();
如果线程池是 Spring 管理的公共 Bean,就不要在每次业务调用结束后关闭。
anyOf 返回以后,其他任务会自动停止吗
不会。
这是使用 anyOf() 时很容易误解的地方。
假设:
任务一耗时 3 秒;
任务二耗时 1 秒;
任务三耗时 5 秒。
任务二先完成后:
anyFuture.join();
会立刻返回任务二的结果。
但任务一和任务三通常仍然会继续执行。
anyOf() 只是表示:
我不再等待其他任务了。
它不表示:
其他任务自动取消。
所以可能出现:
任务二已经返回结果;
任务一还在处理;
任务三也还在处理。
演示其他任务仍然继续执行
CompletableFuture<String> future1 =
CompletableFuture.supplyAsync(() -> {
sleep(3000);
System.out.println("任务一完成");
return "任务一";
});
CompletableFuture<String> future2 =
CompletableFuture.supplyAsync(() -> {
sleep(1000);
System.out.println("任务二完成");
return "任务二";
});
String result = (String) CompletableFuture
.anyOf(future1, future2)
.join();
System.out.println("先拿到:" + result);
sleep(4000);
可能输出:
任务二完成
先拿到:任务二
任务一完成
可以看到,任务二先返回后,任务一仍然继续执行。
是否需要取消剩余任务
有些竞速场景中,第一个结果回来以后,其他任务已经没有必要继续执行。
这时可以尝试取消剩余任务:
CompletableFuture<Object> anyFuture =
CompletableFuture.anyOf(
future1,
future2,
future3
);
Object result = anyFuture.join();
future1.cancel(true);
future2.cancel(true);
future3.cancel(true);
不过这样写会把已经完成的任务也调用一次 cancel(),通常不会有严重问题,但可以先判断:
if (!future1.isDone()) {
future1.cancel(true);
}
if (!future2.isDone()) {
future2.cancel(true);
}
if (!future3.isDone()) {
future3.cancel(true);
}
完整一点:
Object result = CompletableFuture
.anyOf(
future1,
future2,
future3
)
.join();
for (CompletableFuture<?> future : futureList) {
if (!future.isDone()) {
future.cancel(true);
}
}
但仍然要记住:
cancel(true) 只是尝试中断;
不是强制停止任务。
任务是否能及时停止,取决于任务代码是否响应中断。
网络请求不一定能被立即取消
假设异步任务正在调用外部接口:
CompletableFuture.supplyAsync(() -> {
return httpClient.callRemoteApi();
});
调用:
future.cancel(true);
不一定能立即停止底层网络请求。
这取决于:
HTTP 客户端是否支持中断;
请求方法是否响应线程中断;
是否配置了连接超时和读取超时。
所以真实网络请求不能只依赖 cancel(true)。
还应该给客户端配置:
连接超时;
读取超时;
整体调用超时。
如果最先完成的任务失败怎么办
anyOf() 等待的是:
最先完成的任务。
这里的“完成”包括:
正常完成;
异常完成。
假设:
任务一 500 毫秒后失败;
任务二 2 秒后成功;
任务三 3 秒后成功。
任务一最先结束,虽然它是异常结束。
这时候:
CompletableFuture.anyOf(
future1,
future2,
future3
).join();
可能直接抛出 CompletionException。
它不会自动跳过失败任务,再等待下一个成功任务。
这个点非常重要。
anyOf() 不是:
谁先成功用谁。
它是:
谁先完成用谁,包括异常完成。
一个失败示例
CompletableFuture<String> future1 =
CompletableFuture.supplyAsync(() -> {
sleep(500);
throw new RuntimeException("服务一调用失败");
});
CompletableFuture<String> future2 =
CompletableFuture.supplyAsync(() -> {
sleep(2000);
return "服务二调用成功";
});
CompletableFuture<Object> anyFuture =
CompletableFuture.anyOf(
future1,
future2
);
Object result = anyFuture.join();
虽然 future2 最后能成功,但 future1 先异常完成。
因此 join() 可能直接抛异常。
如果想要“第一个成功结果”怎么办
anyOf() 本身不能直接表达:
忽略失败任务,只取第一个成功结果。
一种简单做法是,先把每个任务的异常转换成统一结果对象。
例如:
CompletableFuture<ServiceResult> future1 =
callService1()
.handle((result, ex) -> {
if (ex != null) {
return ServiceResult.fail(
"服务一失败"
);
}
return ServiceResult.success(result);
});
CompletableFuture<ServiceResult> future2 =
callService2()
.handle((result, ex) -> {
if (ex != null) {
return ServiceResult.fail(
"服务二失败"
);
}
return ServiceResult.success(result);
});
然后:
ServiceResult result =
(ServiceResult) CompletableFuture
.anyOf(
future1,
future2
)
.join();
但这个方案仍然只是拿第一个完成的结果。
如果第一个结果是失败对象,还是不会继续自动等待第二个成功结果。
真正实现“第一个成功结果”,需要额外编排逻辑,不能简单依赖一次 anyOf()。
什么时候可以接受第一个完成结果
anyOf() 比较适合这些情况:
多个任务都具有相同价值;
最先完成的结果,无论来自哪个任务,都可以直接使用;
或者每个任务内部已经处理了异常,保证返回的是可判断结果。
如果业务明确要求:
必须拿到成功结果;
失败任务要跳过;
所有任务都失败后才返回失败。
就需要单独设计,不能认为 anyOf() 自动具备这个能力。
一个完整的多存储服务示例
新建类:
com.succos.completablefuture.AnyOfStorageDemo
代码如下:
package com.succos.completablefuture;
import java.util.List;
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 AnyOfStorageDemo {
public static void main(String[] args) {
ThreadPoolExecutor storageExecutor =
new ThreadPoolExecutor(
3,
3,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new NamedThreadFactory(
"storage-worker"
),
new ThreadPoolExecutor.CallerRunsPolicy()
);
CompletableFuture<UploadResult> storage1 =
uploadAsync(
"存储服务一",
3000,
storageExecutor
);
CompletableFuture<UploadResult> storage2 =
uploadAsync(
"存储服务二",
1000,
storageExecutor
);
CompletableFuture<UploadResult> storage3 =
uploadAsync(
"存储服务三",
2000,
storageExecutor
);
List<CompletableFuture<UploadResult>> futureList =
List.of(
storage1,
storage2,
storage3
);
UploadResult result =
(UploadResult) CompletableFuture
.anyOf(
futureList.toArray(
new CompletableFuture[0]
)
)
.join();
System.out.println("--------------------------------");
System.out.println("最先完成的服务:"
+ result.getServiceName());
System.out.println("下载地址:"
+ result.getDownloadUrl());
for (CompletableFuture<UploadResult> future
: futureList) {
if (!future.isDone()) {
future.cancel(true);
}
}
storageExecutor.shutdown();
}
private static CompletableFuture<UploadResult> uploadAsync(
String serviceName,
long millis,
ThreadPoolExecutor executor
) {
return CompletableFuture.supplyAsync(() -> {
System.out.println(Thread.currentThread().getName()
+ " 开始调用:"
+ serviceName);
sleep(millis);
String downloadUrl =
"https://file.example.com/test.pdf";
System.out.println(Thread.currentThread().getName()
+ " 调用完成:"
+ serviceName);
return new UploadResult(
serviceName,
downloadUrl
);
}, executor);
}
private static void sleep(long millis) {
try {
Thread.sleep(millis);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new RuntimeException("任务被中断", e);
}
}
static class UploadResult {
private final String serviceName;
private final String downloadUrl;
public UploadResult(
String serviceName,
String downloadUrl
) {
this.serviceName = serviceName;
this.downloadUrl = downloadUrl;
}
public String getServiceName() {
return serviceName;
}
public String getDownloadUrl() {
return downloadUrl;
}
}
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;
}
}
}
这个例子里三个存储任务同时执行。
服务二耗时最短,所以大概率最先返回。
拿到结果以后,再尝试取消其他尚未完成的任务。
anyOf 不适合普通批量处理
如果业务是:
批量处理 100 个 PDF;
最后统计 100 个文件的成功和失败。
不要使用 anyOf()。
因为它只关心第一个完成的任务。
应该使用:
CompletableFuture.allOf(...)
然后收集全部结果。
我现在会这样判断:
所有任务都重要:allOf;
只需要最先完成的一个:anyOf。
这一节小结
这一节我主要记住几点:
1. anyOf 用来等待多个任务中的任意一个先完成;
2. anyOf 返回 CompletableFuture<Object>;
3. 即使所有任务返回同一种类型,也通常需要进行类型转换;
4. anyOf 返回后,其他异步任务不会自动停止;
5. 可以尝试取消其他任务,但 cancel(true) 不是强制停止;
6. anyOf 取的是第一个完成的任务,不一定是第一个成功的任务;
7. 如果最先完成的任务异常,anyOf().join() 也可能直接抛异常;
8. 普通批量 PDF 处理应该用 allOf,不应该用 anyOf。
用一句话总结:
所有任务都要结果,用 allOf;谁先完成就用谁,用 anyOf。
下一节继续看 orTimeout 和 completeOnTimeout。
这两个方法都和异步任务超时有关,但处理方式不同:一个在超时后抛异常,一个在超时后返回默认结果。