跳到正文
hello world

32. CompletableFuture 使用自定义线程池

发布于阅读量 0

32. CompletableFuture 使用自定义线程池

前面几节写 CompletableFuture 时,大多数代码都是这样:

CompletableFuture.supplyAsync(() -> {
    return "output/test-watermark.pdf";
});

或者:

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

这两种写法都没有指定线程池。

学习阶段可以这么写,因为代码短,也容易观察。

但如果放到真实项目里,尤其是 PDF 水印这种比较重的任务,我一般不会直接用默认线程池,而是会给它指定一个自己的线程池。

这一节就专门整理这个问题。


不指定线程池时会用哪里

先看这段代码:

CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    System.out.println(Thread.currentThread().getName());
    return "PDF 处理成功";
});

因为没有传线程池,所以它默认会使用公共线程池。

线程名一般类似:

ForkJoinPool.commonPool-worker-1

这个公共线程池不是我专门为 PDF 任务创建的。

它是 JDK 提供的公共资源。

如果只是简单学习、写个小 demo,没什么问题。

但如果是业务系统,我会比较谨慎。


默认公共线程池的问题

公共线程池的问题在于:它可能被很多地方共用。

比如项目里其他地方也写了:

CompletableFuture.supplyAsync(() -> {
    // 查询数据
});

或者:

CompletableFuture.runAsync(() -> {
    // 发送通知
});

如果这些地方都不指定线程池,它们就可能都跑到公共线程池里。

这时候 PDF 水印任务如果很重,就可能影响其他异步任务。

比如:

PDF 处理占用了大量线程;
其他异步任务排队变慢;
系统里一些看似无关的功能也受到影响。

这不是我想看到的。

所以我的习惯是:

轻量 demo 可以用默认线程池;
真实业务任务尽量用自己的线程池。

PDF 任务为什么应该单独建线程池

PDF 水印处理不是特别轻的任务。

它可能涉及:

读取文件;
解析 PDF;
遍历页面;
绘制水印;
写出新文件;
占用内存;
占用磁盘 IO。

如果文件比较大,或者一次处理很多个 PDF,压力会更明显。

所以我希望 PDF 任务有自己的边界:

最多几个线程同时处理;
最多允许多少任务排队;
任务太多时怎么拒绝;
线程名能不能看出来是 PDF 任务。

这些都应该由 PDF 专用线程池来控制。


自定义线程池的基本写法

先创建一个线程池:

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

然后在 CompletableFuture 里传进去:

CompletableFuture<String> future = CompletableFuture.supplyAsync(() -> {
    return "output/test-watermark.pdf";
}, pdfExecutor);

注意这里多了第二个参数:

pdfExecutor

这就表示:

这个异步任务不要使用默认公共线程池;
请使用我指定的 PDF 线程池执行。

完整示例:supplyAsync 使用自定义线程池

新建类:

com.succos.completablefuture.CustomExecutorDemo

代码如下:

package com.succos.completablefuture;

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 CustomExecutorDemo {

    public static void main(String[] args) {

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

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

            sleep(3000);

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

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

        String result = future.join();

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

        pdfExecutor.shutdown();
    }

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

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

运行后,线程名会类似:

pdf-watermark-thread-1 开始处理 PDF
pdf-watermark-thread-1 PDF 处理完成

这就说明任务已经不是跑在默认公共线程池里了,而是跑在我自己的 PDF 线程池里。


runAsync 也可以传自定义线程池

runAsync 也一样。

默认写法是:

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

指定线程池后:

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

所以无论是:

runAsync;
supplyAsync。

都可以传入自定义线程池。

区别还是原来的区别:

runAsync 没有返回值;
supplyAsync 有返回值。

线程池只是决定任务在哪里执行。


批量 PDF 使用自定义线程池

接下来把它放到批量 PDF 场景里。

新建类:

com.succos.completablefuture.CustomExecutorPdfDemo

代码如下:

package com.succos.completablefuture;

import java.io.File;
import java.util.ArrayList;
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;
import java.util.stream.Collectors;

public class CustomExecutorPdfDemo {

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

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

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

            futureList.add(future);
        }

        CompletableFuture.allOf(
                futureList.toArray(new CompletableFuture[0])
        ).join();

        List<String> targetPathList = futureList.stream()
                .map(CompletableFuture::join)
                .collect(Collectors.toList());

        for (String targetPath : targetPathList) {
            System.out.println("输出路径:" + targetPath);
        }

        long end = System.currentTimeMillis();

        System.out.println("--------------------------------");
        System.out.println("全部 PDF 处理完成");
        System.out.println("总文件数:" + files.length);
        System.out.println("总耗时:" + (end - start) + " ms");

        pdfExecutor.shutdown();
    }

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

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

这版代码就比较完整了。

它做到了:

每个 PDF 一个 CompletableFuture;
每个任务都提交到 pdfExecutor;
线程池最多 3 个线程同时处理;
其他任务进入队列;
allOf 等待全部完成;
最后统一收集输出路径。

线程池控制并发,allOf 只负责等待

这里再强调一次。

这段代码里,控制并发数量的是:

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

因为核心线程数和最大线程数都是 3,所以同一时间最多 3 个线程处理 PDF。

而这段代码:

CompletableFuture.allOf(
        futureList.toArray(new CompletableFuture[0])
).join();

只是等待所有任务完成。

所以我不会把 allOf 理解成控制并发。

它只是等待。

真正控制任务执行节奏的,是自定义线程池。


为什么最后要 shutdown

因为这里的线程池是我自己创建的:

ThreadPoolExecutor pdfExecutor = new ThreadPoolExecutor(...);

所以任务处理完以后,要记得关闭:

pdfExecutor.shutdown();

否则普通 main 程序可能不会退出。

因为线程池里的工作线程还在。

这里要区分:

future.join():等待任务完成;
allOf().join():等待一批任务完成;
pdfExecutor.shutdown():关闭线程池。

它们不是一回事。


shutdown 放在哪里比较合适

在这个示例里,我是所有任务完成后再调用:

pdfExecutor.shutdown();

也就是:

先提交任务;
allOf 等任务完成;
收集结果;
关闭线程池。

也可以在任务提交完成后就调用 shutdown(),然后再等待任务完成。

比如:

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

    futureList.add(future);
}

pdfExecutor.shutdown();

CompletableFuture.allOf(
        futureList.toArray(new CompletableFuture[0])
).join();

因为 shutdown() 的意思不是马上停止任务,而是不再接收新任务。

已经提交的任务会继续执行。

两种写法都能用。

我个人在学习示例里更喜欢先 allOf().join(),最后再 shutdown(),看起来更直观。

但真实项目里,只要线程池是长期共用的 Spring Bean,就不应该每次任务结束都关掉它。

这个后面会讲。


在 Spring Boot 里不要每次都 new 线程池

如果是普通 main 方法练习,在线程里写:

ThreadPoolExecutor pdfExecutor = new ThreadPoolExecutor(...);

没问题。

但如果是 Spring Boot 项目,我不会在每个接口方法里都这样创建线程池。

比如这种写法不太好:

@PostMapping("/pdf/watermark")
public String watermark() {

    ThreadPoolExecutor pdfExecutor = new ThreadPoolExecutor(...);

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

    return "提交成功";
}

因为每次请求都创建一个线程池,会带来很多问题:

线程池数量不可控;
线程资源浪费;
线程池可能没有关闭;
系统压力变大;
排查问题困难。

更合理的方式是:

把 PDF 线程池定义成一个 Bean;
整个应用共用这一套线程池;
应用关闭时再统一关闭。

后面会专门写“线程池应该定义在哪里”。


自定义线程池也要注意拒绝策略

指定线程池以后,不代表万事大吉。

比如我现在配置的是:

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

这表示:

最多 3 个线程处理;
最多 100 个任务排队;
超过后触发 CallerRunsPolicy。

如果是在本地批处理,这个策略还可以。

但如果是在 Web 请求里,CallerRunsPolicy 可能会让请求线程自己处理 PDF,导致接口变慢。

所以 Web 场景可能更适合:

new ThreadPoolExecutor.AbortPolicy()

然后捕获拒绝异常,返回:

系统任务繁忙,请稍后重试。

也就是说,自定义线程池不只是传进去这么简单,参数和拒绝策略仍然要结合场景。


thenApply 默认在哪个线程执行

还有一个容易忽略的问题。

比如:

CompletableFuture<String> future = CompletableFuture
        .supplyAsync(() -> {
            return "output/test-watermark.pdf";
        }, pdfExecutor)
        .thenApply(path -> {
            return "http://localhost:8080/download?file=" + path;
        });

这里 supplyAsync 使用了 pdfExecutor

thenApply 呢?

一般来说,不带 Async 的后续步骤,可能会由完成上一步任务的线程继续执行。

也就是说,thenApply 有可能还是在:

pdf-watermark-thread-1

里执行。

这个在简单转换时没问题。

比如只是拼一个下载地址,耗时很短。

但如果后续步骤也很重,就要考虑用带 Async 的方法,并指定线程池。

比如:

thenApplyAsync(path -> {
    return buildDownloadUrl(path);
}, anotherExecutor)

这个问题后面会单独整理“带 Async 和不带 Async 的区别”。


PDF 处理任务适合返回结果对象

上面的批量示例返回的是字符串路径。

真实项目里,我更倾向于返回结果对象。

比如:

CompletableFuture<PdfTaskResult> future = CompletableFuture.supplyAsync(() -> {
    try {
        String targetPath = addWatermark(file);
        return PdfTaskResult.success(file.getName(), targetPath);
    } catch (Exception e) {
        return PdfTaskResult.fail(file.getName(), e.getMessage());
    }
}, pdfExecutor);

这样可以避免某个文件失败时,整个 allOf().join() 直接异常。

每个任务都返回一个明确结果:

成功;
失败;
文件名;
输出路径;
失败原因。

最后统一汇总就很方便。


这一节小结

这一节我主要记住几点:

1. CompletableFuture 不指定线程池时,会使用默认公共线程池;
2. PDF 处理这种重任务,不建议长期使用默认公共线程池;
3. runAsync 和 supplyAsync 都可以传入自定义线程池;
4. 自定义线程池可以控制线程数量、队列大小、拒绝策略和线程名;
5. 线程池负责控制并发,allOf 只负责等待全部完成;
6. 自己创建的线程池,用完要考虑 shutdown;
7. Spring Boot 项目里不要每次请求都 new 线程池,应该统一管理;
8. 后续 thenApply 是否继续使用同一线程,要看是否带 Async。

用一句话总结:

CompletableFuture 负责异步流程,线程池负责执行边界;PDF 这种重任务,一定要有自己的线程池。

下一节继续看 thenApplythenAcceptthenRun 怎么选。

因为异步任务执行完以后,下一步到底是“转换结果”“消费结果”,还是“只做一个动作”,这三个方法刚好对应三种情况。

32. CompletableFuture 使用自定义线程池 - Java并发之从Thread到CompletableFuture - 上下文网