19. ThreadPoolExecutor:线程池的基本使用
发布于 • 阅读量 0
前面讲了为什么不建议一直 new Thread()。
简单说就是:线程数量不好控制,任务没有队列,线程不能复用,异常和关闭也不好统一管理。
所以从这一节开始,我把 PDF 批量处理的写法换成线程池。
Java 里创建线程池有很多方式,但我这里不直接用 Executors.newFixedThreadPool(),而是直接用 ThreadPoolExecutor。
原因也很简单:我想把线程池的几个核心参数看清楚。
先看一个最基本的线程池
新建类:
com.succos.threadpool.ThreadPoolExecutorDemo
代码如下:
package com.succos.threadpool;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class ThreadPoolExecutorDemo {
public static void main(String[] args) {
ThreadPoolExecutor executor = new ThreadPoolExecutor(
3,
3,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new ThreadPoolExecutor.CallerRunsPolicy()
);
for (int i = 1; i <= 10; i++) {
int taskId = i;
executor.execute(() -> {
System.out.println(Thread.currentThread().getName()
+ " 开始处理任务 "
+ taskId);
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
System.out.println(Thread.currentThread().getName()
+ " 处理完成任务 "
+ taskId);
});
}
executor.shutdown();
System.out.println("main 提交任务完成");
}
}
这段代码里,我创建了一个线程池,然后提交了 10 个任务。
线程池里核心线程数和最大线程数都设置成 3:
new ThreadPoolExecutor(
3,
3,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new ThreadPoolExecutor.CallerRunsPolicy()
);
所以效果就是:
同一时间最多 3 个线程在执行任务;
剩下的任务在队列里等;
前面的任务执行完以后,线程继续从队列里取任务执行。
这就比每个任务都 new Thread() 稳定多了。
execute 是提交任务
这里提交任务用的是:
executor.execute(() -> {
// 任务逻辑
});
传进去的是一个 Runnable。
也就是说,我没有自己创建线程。
我只是把任务交给线程池:
这个任务要执行,你来安排线程处理。
线程池内部会判断:
现在有没有空闲线程;
需不需要创建新线程;
任务要不要进入队列;
队列满了怎么办。
这些事情就不用我自己手动处理了。
这就是线程池和 new Thread() 最大的区别之一。
运行效果大概是什么样
运行后,控制台可能会看到类似输出:
pool-1-thread-1 开始处理任务 1
pool-1-thread-2 开始处理任务 2
pool-1-thread-3 开始处理任务 3
main 提交任务完成
等待 3 秒...
pool-1-thread-1 处理完成任务 1
pool-1-thread-1 开始处理任务 4
pool-1-thread-2 处理完成任务 2
pool-1-thread-2 开始处理任务 5
pool-1-thread-3 处理完成任务 3
pool-1-thread-3 开始处理任务 6
顺序不一定完全一样。
但能观察到一个现象:
虽然提交了 10 个任务,但同时执行的只有 3 个。
因为线程池里最多只有 3 个工作线程。
后面的任务没有丢,它们在队列里等着。
等某个线程处理完当前任务,就继续取下一个任务。
线程池的几个核心参数
这段代码里最重要的是构造方法:
ThreadPoolExecutor executor = new ThreadPoolExecutor(
3,
3,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new ThreadPoolExecutor.CallerRunsPolicy()
);
这几个参数分别是:
corePoolSize:核心线程数;
maximumPoolSize:最大线程数;
keepAliveTime:非核心线程空闲多久后回收;
unit:时间单位;
workQueue:任务队列;
handler:拒绝策略。
现在先不用把每个细节都展开到源码层面。
我先按使用角度理解。
corePoolSize:核心线程数
第一个参数:
3
是核心线程数。
也就是:
线程池平时主要保留的工作线程数量。
在这个例子里,核心线程数是 3。
所以前 3 个任务提交进来时,线程池会创建 3 个线程来执行。
大概就是:
任务1 -> pool-1-thread-1
任务2 -> pool-1-thread-2
任务3 -> pool-1-thread-3
这 3 个线程执行完任务后,不会马上销毁。
它们会继续留在线程池里,等待后面的任务。
这就是线程复用。
maximumPoolSize:最大线程数
第二个参数也是:
3
这是最大线程数。
它表示线程池最多能创建多少个线程。
这里我把核心线程数和最大线程数都设置成 3,意思是:
这个线程池最多就 3 个线程;
不会再扩容到 4 个、5 个。
这种写法比较适合我现在的 PDF 处理练习。
因为我就是想控制同一时间最多处理 3 个 PDF。
不过后面会讲到,如果最大线程数大于核心线程数,线程池在队列满了以后,可能会继续创建非核心线程。
这一节先不展开,先把固定 3 个线程跑通。
workQueue:任务队列
这里的任务队列是:
new ArrayBlockingQueue<>(100)
意思是创建一个容量为 100 的阻塞队列。
当 3 个核心线程都在忙时,新提交的任务就会进入队列等待。
比如我提交 10 个任务:
任务1、2、3 被 3 个线程执行;
任务4 到任务10 进入队列等待。
等线程处理完任务1,就会从队列里取任务4继续执行。
所以线程池不是把所有任务都同时执行。
它是:
线程有限;
任务排队;
线程空了再取任务。
这比手动 new Thread() 更接近真实项目的处理方式。
拒绝策略是什么
最后一个参数是:
new ThreadPoolExecutor.CallerRunsPolicy()
这是拒绝策略。
它的作用是:当线程池已经满了,队列也满了,再来任务怎么办?
这节代码里队列容量是 100,任务只有 10 个,所以暂时不会触发拒绝策略。
但真实项目里必须考虑这个问题。
比如:
线程池最多 3 个线程;
队列最多放 100 个任务;
如果第 104 个任务来了,就已经放不下了。
这时候就要有拒绝策略。
CallerRunsPolicy 的意思是:
如果线程池处理不了,就让提交任务的线程自己执行这个任务。
比如是 main 线程提交任务,那就由 main 线程自己执行。
它的好处是任务不容易丢,同时会让提交速度慢下来。
后面会单独整理拒绝策略,这里先知道它是线程池满了以后的处理方式。
shutdown 是什么意思
代码最后有一行:
executor.shutdown();
这个不是立即停止线程池。
它的意思是:
不再接收新任务;
已经提交的任务继续执行;
等队列里的任务都执行完以后,线程池再关闭。
所以调用 shutdown() 后,前面提交的 10 个任务还是会继续执行。
如果不调用 shutdown(),这个普通 Java 程序可能不会马上结束。
因为线程池里的线程还在,JVM 会认为程序还没有完全结束。
所以在普通 main 方法练习里,用完线程池后要记得关闭。
main 提交完任务不代表任务执行完
代码里有这句:
System.out.println("main 提交任务完成");
这个日志可能很快就打印出来。
但这不代表 10 个任务都执行完了。
它只代表:
main 线程已经把任务提交给线程池了。
任务真正执行,是线程池里的工作线程在做。
所以这里要分清楚两个概念:
提交完成;
执行完成。
提交完成只是任务进入线程池。
执行完成是线程真正跑完任务逻辑。
如果我想让 main 等所有任务执行完,还需要配合:
executor.awaitTermination(...)
下一段就加上这个。
等待线程池任务执行完
我可以把代码改得完整一点。
package com.succos.threadpool;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class ThreadPoolExecutorWaitDemo {
public static void main(String[] args) {
ThreadPoolExecutor executor = new ThreadPoolExecutor(
3,
3,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new ThreadPoolExecutor.CallerRunsPolicy()
);
long start = System.currentTimeMillis();
for (int i = 1; i <= 10; i++) {
int taskId = i;
executor.execute(() -> {
System.out.println(Thread.currentThread().getName()
+ " 开始处理任务 "
+ taskId);
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
System.out.println(Thread.currentThread().getName()
+ " 处理完成任务 "
+ taskId);
});
}
executor.shutdown();
try {
boolean finished = executor.awaitTermination(1, TimeUnit.HOURS);
long end = System.currentTimeMillis();
if (finished) {
System.out.println("全部任务处理完成");
} else {
System.out.println("等待超时,还有任务没处理完");
}
System.out.println("总耗时:" + (end - start) + " ms");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
这里多了:
executor.awaitTermination(1, TimeUnit.HOURS);
它的意思是:
最多等 1 个小时;
如果线程池里的任务都执行完了,就返回 true;
如果超过 1 个小时还没结束,就返回 false。
注意,awaitTermination() 一般要在 shutdown() 后调用。
因为不调用 shutdown(),线程池还会继续接收任务,就谈不上“终止”。
shutdown 和 awaitTermination 的关系
这两个方法经常一起出现。
我现在这样理解:
shutdown():告诉线程池,不要再接新任务了;
awaitTermination():当前线程等待线程池里的任务执行完。
它们解决的问题不一样。
只调用 shutdown(),不会让当前线程等待任务完成。
只调用 awaitTermination(),如果前面没有 shutdown(),线程池可能一直不进入终止流程。
所以比较常见的写法是:
executor.shutdown();
try {
executor.awaitTermination(1, TimeUnit.HOURS);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
在这个练习项目里,这样写就够了。
用线程池处理 PDF 的思路
现在把模拟任务换成 PDF 水印任务,整体思路就是:
遍历 input 目录;
每个 PDF 封装成一个 Runnable;
提交给线程池;
线程池最多 3 个线程同时处理;
其他任务排队;
最后 shutdown 并等待完成。
代码结构大概是这样:
for (File file : files) {
executor.execute(() -> {
processPdf(file);
});
}
这里的重点是:我不再自己创建线程。
我只是提交任务。
线程池内部的线程会复用。
PDF 版本示例
新建类:
com.succos.threadpool.ThreadPoolPdfDemo
代码如下:
package com.succos.threadpool;
import java.io.File;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
public class ThreadPoolPdfDemo {
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 executor = new ThreadPoolExecutor(
3,
3,
60,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new ThreadPoolExecutor.CallerRunsPolicy()
);
long start = System.currentTimeMillis();
for (File file : files) {
executor.execute(() -> {
System.out.println(Thread.currentThread().getName()
+ " 开始处理:"
+ file.getName());
try {
Thread.sleep(3000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
return;
}
System.out.println(Thread.currentThread().getName()
+ " 处理完成:"
+ file.getName());
});
}
executor.shutdown();
try {
boolean finished = executor.awaitTermination(1, TimeUnit.HOURS);
long end = System.currentTimeMillis();
if (finished) {
System.out.println("全部 PDF 处理完成");
} else {
System.out.println("等待超时,还有 PDF 没处理完");
}
System.out.println("总文件数:" + files.length);
System.out.println("总耗时:" + (end - start) + " ms");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
这个版本虽然还是用 sleep 模拟处理,但结构已经接近真实批量任务了。
后面只要把 Thread.sleep(3000) 替换成真正的 PdfWatermarkService.addWaterMakerOfPDF(...) 就行。
线程池和 Semaphore 的区别再看一眼
前面用 Semaphore(3) 时,虽然能限制同一时间最多 3 个线程处理 PDF,但我还是可能创建很多线程。
现在用线程池,情况不一样。
比如线程池是 3 个线程:
pool-1-thread-1
pool-1-thread-2
pool-1-thread-3
提交 100 个 PDF 任务,也不是创建 100 个线程。
而是:
3 个线程反复从队列里取任务;
处理完一个,再处理下一个。
这就是线程池比 Thread + Semaphore 更适合批量任务的原因。
Semaphore 是限制入口。
线程池是管理线程和任务队列。
两者不是完全替代关系,但在 PDF 批量处理这种场景里,线程池更自然。
这一节小结
这一节我主要记住几点:
1. ThreadPoolExecutor 用来管理线程和任务队列;
2. execute() 是提交 Runnable 任务;
3. corePoolSize 是核心线程数;
4. maximumPoolSize 是最大线程数;
5. workQueue 是任务等待队列;
6. shutdown() 表示不再接收新任务;
7. awaitTermination() 用来等待线程池任务执行完成;
8. PDF 批量处理更适合提交任务给线程池,而不是每个 PDF 都 new Thread。
用一句话总结:
线程池的核心思路是:线程数量有限,任务可以排队,线程可以复用。
下一节继续看线程池里的任务到底是怎么排队和执行的。
把这个流程搞清楚以后,后面理解核心线程数、队列、最大线程数和拒绝策略就不会乱。