跳到正文
hello world

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。

用一句话总结:

线程池的核心思路是:线程数量有限,任务可以排队,线程可以复用。

下一节继续看线程池里的任务到底是怎么排队和执行的。

把这个流程搞清楚以后,后面理解核心线程数、队列、最大线程数和拒绝策略就不会乱。