跳到正文
hello world

30. join:等待一个 CompletableFuture 完成

发布于阅读量 0

30. join:等待一个 CompletableFuture 完成

前面用了 runAsyncsupplyAsync

这两个方法都会把任务提交到后台异步执行。

但是有一点要记住:

提交异步任务,不代表任务已经执行完。

所以如果我需要在某个位置等待异步任务完成,就会用到:

future.join();

这一节就专门整理 join()


join 最基本的作用

join() 的作用很直接:

等待 CompletableFuture 执行完成。

如果这个 CompletableFuture 有返回值,join() 会返回结果。

如果这个 CompletableFuture 没有返回值,join() 就只是等待它结束。

比如:

CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    return "PDF 处理成功";
});

String result = future.join();

这里的 future 是:

CompletableFuture<String>

所以 join() 拿到的就是一个 String


supplyAsync 里使用 join

新建类:

com.succos.completablefuture.JoinSupplyDemo

代码如下:

package com.succos.completablefuture;

import java.util.concurrent.CompletableFuture;

public class JoinSupplyDemo {

    public static void main(String[] args) {

        System.out.println("main 开始:"
                + Thread.currentThread().getName());

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

            sleep(3000);

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

            return "output/test-watermark.pdf";
        });

        System.out.println("main 提交异步任务后,继续往下执行");

        String result = future.join();

        System.out.println("main 拿到结果:" + result);
        System.out.println("main 结束:"
                + Thread.currentThread().getName());
    }

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

这里最关键的是:

String result = future.join();

如果异步任务还没执行完,main 线程会在这里等待。

等后台任务返回:

return "output/test-watermark.pdf";

join() 才会继续往下走。


join 会阻塞当前线程

这个点很重要。

join() 会阻塞的是“调用它的线程”。

比如:

String result = future.join();

这行代码写在 main 方法里,那么阻塞的就是 main 线程。

如果它写在 Web 接口线程里,那么阻塞的就是当前请求线程。

所以 join() 不能乱放。

它虽然写起来简单,但本质上还是等待。

只要是等待,就可能影响当前线程继续执行。


runAsync 里使用 join

如果用的是 runAsync,任务没有返回值。

比如:

CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
    System.out.println("处理 PDF");
});

这时候 join() 不会返回业务结果。

它只是等待任务执行完。

新建类:

com.succos.completablefuture.JoinRunDemo

代码如下:

package com.succos.completablefuture;

import java.util.concurrent.CompletableFuture;

public class JoinRunDemo {

    public static void main(String[] args) {

        CompletableFuture<Void> future = CompletableFuture.runAsync(() -> {
            System.out.println(Thread.currentThread().getName()
                    + " 开始处理 PDF");

            sleep(3000);

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

        System.out.println("main 提交任务完成");

        future.join();

        System.out.println("main 确认 PDF 任务已经结束");
    }

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

这里的:

future.join();

只是为了保证:

异步任务执行完以后,main 再继续打印最后一句。

join 和 Future.get 有点像

前面学 Future 的时候,用过:

String result = future.get();

现在 CompletableFuture 里用:

String result = future.join();

这两个方法都可以等待任务完成,也都可以拿结果。

从使用感觉上看,它们有点像。

但异常处理不一样。

get() 需要处理受检异常:

try {
    String result = future.get();
} catch (InterruptedException e) {
    Thread.currentThread().interrupt();
} catch (ExecutionException e) {
    e.printStackTrace();
}

join() 不要求强制写 try-catch

比如:

String result = future.join();

代码看起来更简洁。


join 的异常会包装成 CompletionException

虽然 join() 不强制处理异常,但不代表异常消失了。

如果异步任务里抛异常:

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

String result = future.join();

执行到 join() 时,还是会抛异常。

只不过这个异常通常会被包装成:

CompletionException

可以写一个例子看一下。

新建类:

com.succos.completablefuture.JoinExceptionDemo

代码如下:

package com.succos.completablefuture;

import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;

public class JoinExceptionDemo {

    public static void main(String[] args) {

        CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
            System.out.println("异步任务开始");

            int result = 1 / 0;

            return "处理成功:" + result;
        });

        try {
            String result = future.join();

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

        } catch (CompletionException e) {
            System.out.println("join 获取结果失败");
            System.out.println("真正原因:" + e.getCause());
        }
    }
}

这里真正的异常是:

ArithmeticException: / by zero

join() 抛出来时,外层会包一层 CompletionException

所以排查时可以看:

e.getCause()

join 放在哪里很关键

join() 放的位置,会直接影响并发效果。

比如我有多个 PDF,要异步处理。

如果这样写:

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

    String result = future.join();

    System.out.println(result);
}

这段代码看起来用了 CompletableFuture,但并发效果很差。

因为每提交一个任务,就马上 join() 等它完成。

执行流程会变成:

提交 a.pdf;
等待 a.pdf 处理完;

提交 b.pdf;
等待 b.pdf 处理完;

提交 c.pdf;
等待 c.pdf 处理完。

这和顺序处理差不多。


批量任务不要提交一个就 join 一个

更合理的写法是:

先把所有任务都提交出去;
把 CompletableFuture 保存起来;
最后再统一 join。

示例:

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

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

    futureList.add(future);
}

for (CompletableFuture<String> future : futureList) {
    String result = future.join();
    System.out.println(result);
}

这个结构才是真正的批量异步。

它和前面学过的思路是一样的:

Thread:先全部 start,再统一 join;
Future:先全部 submit,再统一 get;
CompletableFuture:先全部 supplyAsync,再统一 join。

这个习惯很重要。


批量 PDF 的 join 示例

新建类:

com.succos.completablefuture.JoinBatchPdfDemo

代码如下:

package com.succos.completablefuture;

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

public class JoinBatchPdfDemo {

    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);
        }

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

        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();
        }
    }
}

这段代码先把所有 PDF 的异步任务提交出去。

然后再统一遍历 futureList 调用 join()

这样多个 PDF 才能并发处理。


这里还没有指定自定义线程池

上面的例子里,我写的是:

CompletableFuture.supplyAsync(() -> {
    return targetPath;
});

没有传线程池。

所以它会使用默认公共线程池。

学习阶段可以先这样观察。

但真实 PDF 处理里,我更倾向于使用自己的线程池:

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

这样 PDF 任务不会占用公共线程池。

后面会单独整理 CompletableFuture 使用自定义线程池。


join 不等于关闭线程池

还有一个点要分清楚。

join() 等待的是某个 CompletableFuture 完成。

它不负责关闭线程池。

如果我用的是默认公共线程池,这个问题不明显。

但如果我自己创建了线程池:

ThreadPoolExecutor pdfExecutor = new ThreadPoolExecutor(...);

任务完成后,还是要根据情况关闭:

pdfExecutor.shutdown();

所以:

join:等任务完成;
shutdown:关闭线程池。

这两个不是一回事。


join 适合用在哪里

我现在会这样用 join()

学习示例里:用 join 保证 main 等异步任务结束;
批量任务里:先收集 CompletableFuture,最后统一 join;
流程最后:在确实需要结果的位置 join。

但我不会在每个异步步骤后面立刻 join()

因为那样会把异步流程又写回同步流程。

比如这种写法就不太好:

String path = CompletableFuture.supplyAsync(() -> {
    return processPdf(file);
}).join();

String url = CompletableFuture.supplyAsync(() -> {
    return buildUrl(path);
}).join();

这样每一步都在等。

如果只是一个连续流程,更适合用:

CompletableFuture<String> future = CompletableFuture
        .supplyAsync(() -> processPdf(file))
        .thenApply(path -> buildUrl(path));

最后需要结果时,再 join()


这一节小结

这一节我主要记住几点:

1. join 用来等待 CompletableFuture 完成;
2. CompletableFuture<T> 调用 join 后,可以拿到 T 类型结果;
3. CompletableFuture<Void> 调用 join,只是等待任务结束;
4. join 会阻塞调用它的当前线程;
5. join 不强制处理受检异常,但异常会包装成 CompletionException;
6. 批量异步任务不要提交一个就 join 一个;
7. 更好的写法是先全部提交,再统一 join;
8. join 只负责等任务,不负责关闭线程池。

用一句话总结:

join 是 CompletableFuture 的等待点,但放早了会把异步写回同步。

下一节继续看 allOf

因为批量任务里一个个遍历 join() 能用,但如果我想表达“等待这一批异步任务全部完成”,CompletableFuture.allOf() 会更自然。