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 这种重任务,一定要有自己的线程池。
下一节继续看 thenApply、thenAccept、thenRun 怎么选。
因为异步任务执行完以后,下一步到底是“转换结果”“消费结果”,还是“只做一个动作”,这三个方法刚好对应三种情况。