概述

任务是一组逻辑工作单元,而线程则是是任务异步执行的机制。通过线程池就可以简化线程的管理工作。
ThreadPoolExecutor是线程池的核心实现。线程池中预先提供了指定数量的可重用的线程,使用线程池避免了线程创建和终止的开销,节省了系统的资源。并且线程池维护了一些基础的数据统计,方便了线程的监控和管理。

线程池的基本使用

参数解释

线程池的创建的需要指定非常多的参数,我们需要理解每个参数的含义.我们就以参数最多的构造器为例,来解释每个参数的含义。

1
2
3
4
5
6
7
ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory,
RejectedExecutionHandler handler)

corePoolSize 核心线程数。
maximumPoolSize最大线程数。
这里需要注意:当一个新的任务提交给线程池之后:

  • 如果当前运行线程的数量小于核心线程数,无论有无空闲的线程,都换创建新的线程。
  • 如果当前运行的线程数大于核心线程数,小于最大线程数,只有当等待队列满之后,才会新建线程。
  • 如果等待队列已满,且线程数已达到最大线程数,那么就会根据指定的拒绝策略进行处理了。

keepAliveTime线程最大空闲时间:如果当前线程池中多于核心线程数的线程如果超过最大空闲时间就会被终止。

unit:TimeUnit时间单位。
workQueue线程等待队列。
threadFactory线程创建工厂。
RejectedExecutionHandler:拒绝策略。当线程池已经关闭或者已经达到饱和状态,新提交的任务会被拒绝。一共有4种拒绝策略:

  • AbortPolicy:默认策略,在需要拒绝任务时,抛出RejectedExecutionException
  • CallerRunsPolicy:直接在execute方法的调用线程种运行被拒绝的任务,如果线程池已经关闭,任务将被丢弃。
  • DiscardPolicy:直接丢弃任务
  • DiscardOldestPolicy:丢弃队列中等待时间最长的任务,并执行当前提交的任务,如果线程池被关闭,任务将被丢弃。
  • 我们也可以自己继承RejectedExcutionHandler自定义自己的拒绝策略,拒绝策略的运行需要指定线程池和队列的容量。

生命周期

线程池一共有5种状态:

  • Running:可以接收新的任务和队列任务
  • shutdown:不接受新的任务,但是会运行队列任务。
  • stop:不接受新的任务,也不会运行队列任务,并且中断正在运行的任务。
  • tidying:所有任务都已经终止,workCount为0,当前池状态为tidying时会运行运行terminated()方法
  • terminated,terminated()方法执行完毕。

1C7IpT.png

源码分析

1CHiBd.png

重要属性

ThreadPoolExecutor内部有一个非常重要的内部类Worker,它继承自AQS实现了Runnable接口,实现了不可重入的互斥锁。在线程池种持有一个Work集合,一个worker对应一个工作想,当线程池启动时,对应的worker会执行池种的任务,执行任务完毕后会从阻塞列表中获取一个新的任务继续执行。

Worker内部维护了三个变量,用于记录每个工作线程的工作状态。

1
2
3
4
5
6
7
/** 工作线程 */
final Thread thread;
/** 初始运行任务 */
Runnable firstTask;
/**任务完成计数*/
volatile long completedTasks;

内部的属性:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
/*当核心参数数已满,新增任务的存储队列*/
private final BlockingQueue<Runnable> workQueue;

/*线程池运行期间的锁*/
private final ReentrantLock mainLock = new ReentrantLock();

/*工作池线程*/
private final HashSet<Worker> workers = new HashSet<Worker>();

/*awaitTermination的等待队列*/
private final Condition termination = mainLock.newCondition();

/*最大线程池数量*/
private int largestPoolSize;

/*已完成的任务数*/
private long completedTaskCount;

/*线程创建工厂*/
private volatile ThreadFactory threadFactory;

/*拒绝处理器*/
private volatile RejectedExecutionHandler handler;

/*最大闲置时间*/
private volatile long keepAliveTime;

/*是否允许核心线程超时*/
private volatile boolean allowCoreThreadTimeOut;

/*核心线程数*/
private volatile int corePoolSize;

/*最大线程数*/
private volatile int maximumPoolSize;

/*拒绝策略,默认的拒绝策略为AbortPolicy,即直接抛出异常*/
private static final RejectedExecutionHandler defaultHandler =
new AbortPolicy();

/*针对shutdown和shutdownNow的运行权限许可*/
private static final RuntimePermission shutdownPerm =
new RuntimePermission("modifyThread");

/*高3位标识线程池的运行状态,低29位表示线程池中的任务数*/
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));

重要方法解析

execute方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
public void execute(Runnable command) {
if (command == null)
throw new NullPointerException();
int c = ctl.get();
//线程池中的线程数小于核心线程数
if (workerCountOf(c) < corePoolSize) {
if (addWorker(command, true))//添加到工作线程集合
return;
c = ctl.get();
}
//线程正在运行,添加到等待队列
if (isRunning(c) && workQueue.offer(command)) {
int recheck = ctl.get();
//如果池的状态shutdown,移除任务,执行拒绝策略
if (! isRunning(recheck) && remove(command))
reject(command);//执行拒绝策略
else if (workerCountOf(recheck) == 0)//工作线程位空,添加新的工作线程
addWorker(null, false);
}
else if (!addWorker(command, false))//再次尝试添加任务
//再次尝试还是失败,就会根据拒绝策略进行处理了
reject(command);
}

提交一个任务到线程池,任务不一定会立即执行。提交的任务可以在一个新的线程中执行,也可能在已存在线程中执行。如果由于池关闭或池容量已经饱和导致任务无法提交,那么就根据拒绝策略来处理提交过来的任务。

  1. 如果正在运行的线程数少于corePoolSize,那么就会通过addWorker方法尝试开启一个新的线程并把提交的任务作为它的firstTask运行,addWorker会检查ctl的状态来判断是否可以添加新的线程。
  2. 如果addWorker执行失败(返回false),那么就会把任务添加到等待队列。这里需要对ctl进行双重检查。
  3. 如果不能任务不能入队,那么就会再次尝试增加一个新的线程,如果添加失败,就意味着池关闭或已经饱和,这个时候就会根据拒绝策略来进行处理。

addWorker方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (;;) { //自旋,判断可以添加节点的前提条件
int c = ctl.get();
int rs = runStateOf(c);//线程运行的状态码

// 检查线程池的状态
if (rs >= SHUTDOWN &&
! (rs == SHUTDOWN &&
firstTask == null &&
! workQueue.isEmpty()))
return false;

for (;;) {//检查工作线程数是否饱和
int wc = workerCountOf(c);//获取工作线程数
if (wc >= CAPACITY ||
wc >= (core ? corePoolSize : maximumPoolSize))
return false;
if (compareAndIncrementWorkerCount(c))
//可以添加新的线程,更新ctl中关于新线程的表示
//跳出自旋
break retry;
c = ctl.get(); // 更新ctl失败,重新读取ctl
if (runStateOf(c) != rs)
continue retry;
// else CAS failed due to workerCount change; retry inner loop
}
}

boolean workerStarted = false;
boolean workerAdded = false;
Worker w = null;
try {
//创建新的工作线程
w = new Worker(firstTask);
final Thread t = w.thread;
if (t != null) {
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();//加锁,准备添加新的工作线程
try {
// Recheck while holding lock.
// Back out on ThreadFactory failure or if
// shut down before lock acquired.
int rs = runStateOf(ctl.get());

//重新加成sunState
if (rs < SHUTDOWN ||
(rs == SHUTDOWN && firstTask == null)) {
if (t.isAlive()) // precheck that t is startable
throw new IllegalThreadStateException();
workers.add(w);//添加工作线程
int s = workers.size();
if (s > largestPoolSize)
largestPoolSize = s;//更新最大池容量
workerAdded = true;
}
} finally {
mainLock.unlock();
}
if (workerAdded) {
t.start();//添加成功,启动线程
workerStarted = true;
}
}
} finally {
if (! workerStarted)
//添加失败,回滚操作
addWorkerFailed(w);
}
return workerStarted;
}

runWorker方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
final void runWorker(Worker w) {
Thread wt = Thread.currentThread();
Runnable task = w.firstTask;
w.firstTask = null;
w.unlock(); // 允许中断
boolean completedAbruptly = true;
try {
//如果work的firstTask为null,就调用getTask从队列中取任务
while (task != null || (task = getTask()) != null) {
w.lock();//加锁
if ((runStateAtLeast(ctl.get(), STOP) ||
(Thread.interrupted() &&
runStateAtLeast(ctl.get(), STOP))) &&
!wt.isInterrupted())
//中断线程
wt.interrupt();
try {
//执行任务的前逻辑
beforeExecute(wt, task);
Throwable thrown = null;
try {
task.run();//执行任务
} catch (RuntimeException x) {
thrown = x; throw x;
} catch (Error x) {
thrown = x; throw x;
} catch (Throwable x) {
thrown = x; throw new Error(x);
} finally {
//执行任务后逻辑,自定义任务
afterExecute(task, thrown);
}
} finally {
task = null;
w.completedTasks++; //完成的任务数加一
w.unlock(); //解锁
}
}
completedAbruptly = false;
} finally {

//处理工作线程退出逻辑
processWorkerExit(w, completedAbruptly);
}
}

runWorker是工作线程运行的核心方法,循环从队列中获取任务并执行。工作线程启动后,会首先运行内部持有的任务firstTask.如果firstTask为null,那么就会循环调用getTask方法从队列中获取任务执行。在任务执行前后可以调用beforeExecuteafterWxecute处理执行前后的逻辑。如果线程池的状态正在停止,那么需要确保线程被中断,否则需要确保线程没有被中断。

getTask方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
private Runnable getTask() {
boolean timedOut = false; // Did the last poll() time out?

for (;;) {
int c = ctl.get();
int rs = runStateOf(c);//获取runState

// Check if queue empty only if necessary.
if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) {
//线程池已经关闭或等待队列为null
decrementWorkerCount();
return null;
}

int wc = workerCountOf(c);//获取工作线程的数wokerCount

// Are workers subject to culling?
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;

if ((wc > maximumPoolSize || (timed && timedOut))
&& (wc > 1 || workQueue.isEmpty())) {
if (compareAndDecrementWorkerCount(c))//修改ctl工作线程数减一
return null;
continue;
}

try {
Runnable r = timed ?
workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) :
workQueue.take();//出队,直到得带元素可用
if (r != null)
return r;
timedOut = true;
} catch (InterruptedException retry) {
timedOut = false;
}
}
}

processWorkerExit方法

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
private void processWorkerExit(Worker w, boolean completedAbruptly) {
if (completedAbruptly) // If abrupt, then workerCount wasn't adjusted
//如果任务线程被中断,则工作线程数量减一
decrementWorkerCount();

final ReentrantLock mainLock = this.mainLock;
mainLock.lock(); //加锁
try {
//更新完成任务数
completedTaskCount += w.completedTasks;
//移除工作线程
workers.remove(w);
} finally {
mainLock.unlock();//解锁
}

tryTerminate();//尝试终止线程池

int c = ctl.get();
if (runStateLessThan(c, STOP)) {//线程池尚未完全停止
if (!completedAbruptly) {//工作线程非异常退出
//获取当前核心线程数
int min = allowCoreThreadTimeOut ? 0 : corePoolSize;
if (min == 0 && ! workQueue.isEmpty())
//如果允许空闲工作线程等待任务,且任务队列不为空,则min为1
min = 1;
if (workerCountOf(c) >= min)
return; // replacement not needed
}
//继续尝试添加新的工作线程
addWorker(null, false);
}
}

工作线程处理完所有的任务之后,调用池方法处理工作线程退出逻辑,为已经死亡的工作线程执行相关的清除操作。此方法会从线程池中内的工作线程集合中移除当前线程,并会尝试终止线程池。
在下面这几种情况下,可能会替换当前工作线程:

  1. 用户任务执行异常导致线程退出
  2. 工作线程数少于corePoolSize
  3. 等待队列不为空,但是没有工作线程

tryTerminate方法

该方法的主要作用就是尝试终止线程池。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
final void tryTerminate() {
for (;;) {
int c = ctl.get();
if (isRunning(c) ||
runStateAtLeast(c, TIDYING) ||
(runStateOf(c) == SHUTDOWN && ! workQueue.isEmpty()))
return;
if (workerCountOf(c) != 0) { // Eligible to terminate
interruptIdleWorkers(ONLY_ONE);//中断空闲线程
return;
}

final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
//线程池已经关闭,等待队列为空,并且工作线程等于0,更新池状态为TINDYING
if (ctl.compareAndSet(c, ctlOf(TIDYING, 0))) {
try {

//线程池终止操作,需要自定义实现
terminated();
} finally {
ctl.set(ctlOf(TERMINATED, 0));
//唤醒等待池结束的线程
termination.signalAll();
}
return;
}
} finally {
mainLock.unlock();
}
// else retry on failed CAS
}
}

该方法用于尝试终止线程池,shutDown,shutdownNoe,remove中局势通过此方法来终止线程池的。此方法必须在人恶化可能导致终止的行为之后被调用。一如减少工作线程数,移除队列中的任务,或者是在工作线程运行完毕后处理工作线程退出逻辑方法processWorkerExit