跳到正文
hello world

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>

即使三个任务返回的都是 StringanyOf() 的结果类型仍然是 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。

下一节继续看 orTimeoutcompleteOnTimeout

这两个方法都和异步任务超时有关,但处理方式不同:一个在超时后抛异常,一个在超时后返回默认结果。