概述

线程池故名思意就是存放线程的池子。线程是JVM的稀缺资源,使用线程池可以减少线程的创建和销毁的次数,降低系统的资源消耗。当有任务到达的时候,无效等待创建新的线程便能立即执行,提高了响应速度。使用线程池还可以方便对线程进行管理

ThreadPoolExecuter是线程池框架的一个核心类,下面就对ThreadPoolExecuter进行分析。

属性

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
// 线程池的控制状态,用高3位来表示线程池的运行状态,低29位来表示线程池中工作线程的数量
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
//值为29,用来表示偏移量
private static final int COUNT_BITS = Integer.SIZE - 3;
//线程池的最大容量,其值的二进制为:00011111111111111111111111111111(29个1)
private static final int CAPACITY = (1 << COUNT_BITS) - 1;

// 线程池的运行状态,总共有5个状态,用高3位来表示
private static final int RUNNING = -1 << COUNT_BITS;
private static final int SHUTDOWN = 0 << COUNT_BITS;
private static final int STOP = 1 << COUNT_BITS;
private static final int TIDYING = 2 << COUNT_BITS;
private static final int TERMINATED = 3 << COUNT_BITS;

//任务缓存队列,用来存放等待执行的任务
private final BlockingQueue<Runnable> workQueue;

//全局锁,对线程池状态等属性修改时需要使用这个锁
private final ReentrantLock mainLock = new ReentrantLock();

//线程池中工作线程的集合,访问和修改需要持有全局锁
private final HashSet<Worker> workers = new HashSet<Worker>();

// 终止条件
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;

//核心池大小,若allowCoreThreadTimeOut被设置,核心线程全部空闲超时被回收的情况下会为0
private volatile int corePoolSize;

//最大池大小,不得超过CAPACITY
private volatile int maximumPoolSize;

//默认的任务拒绝策略
private static final RejectedExecutionHandler defaultHandler =
new AbortPolicy();

private static final RuntimePermission shutdownPerm =
new RuntimePermission("modifyThread");

private final AccessControlContext acc;

线程池的状态

通过ThreadPoolExecutor中的属性,我们可以知道线程池一共有五种状态。

  • RUNNING: 当线程池处于当前状态时,能够接收新任务,并且能够对已添加的任务进行处理。
  • SHUTDOWN:当线程池处于当前状态,不能够接收新任务,但是能够处理已经添加的任务。
  • STOP:处于当前状态的线程池,不能接收新任务,不能处理已经收的任务,并且会中断正在处理的任务。
  • TIDYING:当所有任务已经终止,且ctl记录的任务数量为0时,线程池处于DIDYING状态。处于该状态的线程池会执行钩子函数terminated()
  • TERMINATED:线程池彻底终止。

线程池中ctl变量的高3位用来记录线程池的状态,低29位用来表示线程池中工作线程的个数。

拒绝策略

  • AbortPolicy:丢弃任务,并抛出RejectedExecutionException异常。
  • DiscardPolicy:丢弃任务,不会抛出任务
  • DiscardOldestPolicy:丢弃队列最前面的任务,然后重新提交被拒绝的任务。
  • CallerRunsPolicy:由调用线程处理该任务

当然我们也可以通过RejectedExecutionHandler 接口自定义自己的拒绝策略。

核心方法

提交任务

1
2
3
4
5
6
public Future<?> submit(Runnable task) {
if (task == null) throw new NullPointerException();
RunnableFuture<Void> ftask = newTaskFor(task, null);
execute(ftask);
return ftask;
}

submit方法中最终调用了execute方法。

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
public void execute(Runnable command) {
if (command == null)
throw new NullPointerException();
//获取线程池的状态变量
int c = ctl.get();
//如果工作线程数小于核心线程数
if (workerCountOf(c) < corePoolSize) {
//创建worker
if (addWorker(command, true))
return;
//创建worker失败后再次获取线程池控制状态
c = ctl.get();
}
//如果线程池处于Running状态,就将任务加入阻塞队列
if (isRunning(c) && workQueue.offer(command)) {
// 再次检查,获取线程池控制状态,防止在任务入队的过程中线程池关闭了或者线程池中没有线程了
int recheck = ctl.get();
//线程池不处于RUNNING状态,且将任务从workQueue移除成功
if (! isRunning(recheck) && remove(command))
//采取任务拒绝策略
reject(command);
//worker数量等于0
else if (workerCountOf(recheck) == 0)
//创建worker
addWorker(null, false);
}
//(3)
else if (!addWorker(command, false)) //创建worker
reject(command); //如果创建worker失败,采取任务拒绝策略
}

整个流程如下:

  • 如果线程池中的工作线程数小于corePoolSize,则创建新线程来执行任务
  • 如果工作线程数量大于或等于corePoolSize,则将任务加入BlockingQueue
  • 若无法将任务加入阻塞队列,且工作线程数小于最大线程数,则创建新的线程来执行任务
  • 当工作线程数量达到最大线程数,则创建工作线程失败,采取任务拒绝策略

GjhqSg.png

创建线程

我们在execute方法的实现中可以看出、创建线程并执行任务由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
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
//addWorker有两个参数:Runnable类型的firstTask,用于指定新增的线程执行的第一个任务;boolean类型的core,表示是否创建核心线程
//该方法的返回值代表是否成功新增一个线程
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (;;) {
int c = ctl.get();
//获取线程池的当前运行状态
int rs = runStateOf(c);

// (1)
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;
//CAS操作递增workCount
//如果成功,那么创建线程前的所有条件校验都满足了,准备创建线程执行任务,退出retry循环
//如果失败,说明有其他线程也在尝试往线程池中创建线程(往线程池提交任务可以是并发的),则继续往下执行
if (compareAndIncrementWorkerCount(c))
break retry;
//重新获取线程池控制状态
c = ctl.get();
// 如果线程池的状态发生了变更,如有其他线程关闭了这个线程池,那么需要回到外层的for循环
if (runStateOf(c) != rs)
continue retry;
//如果只是CAS操作失败的话,进入内层的for循环就可以了
}
}

//到这里,创建线程前的所有条件校验都满足了,可以开始创建线程来执行任务
//worker是否已经启动
boolean workerStarted = false;
//是否已将这个worker添加到workers这个HashSet中
boolean workerAdded = false;
Worker w = null;
try {
//创建一个worker,从这里可以看出对线程的包装
w = new Worker(firstTask);
//取出worker中的线程对象,Worker的构造方法会调用ThreadFactory来创建一个新的线程
final Thread t = w.thread;
if (t != null) {
//获取全局锁, 并发的访问线程池workers对象必须加锁,持有锁的期间线程池也不会被关闭
final ReentrantLock mainLock = this.mainLock;
//加锁
mainLock.lock();
try {
//重新获取线程池的运行状态
int rs = runStateOf(ctl.get());

//小于SHUTTDOWN即RUNNING
//等于SHUTDOWN并且firstTask为null,不接受新的任务,但是会继续执行等待队列中的任务
if (rs < SHUTDOWN ||
(rs == SHUTDOWN && firstTask == null)) {
//worker里面的thread不能是已启动的
if (t.isAlive())
throw new IllegalThreadStateException();
//将新创建的线程加入到线程池中
workers.add(w);
int s = workers.size();
// 更新largestPoolSize
if (s > largestPoolSize)
largestPoolSize = s;
workerAdded = true;
}
} finally {
mainLock.unlock();
}
//线程添加线程池成功,则启动新创建的线程
if (workerAdded) {
t.start();
workerStarted = true;
}
}
} finally {
//若线程启动失败,做一些清理工作,例如从workers中移除新添加的worker并递减wokerCount
if (! workerStarted)
addWorkerFailed(w);
}
//返回线程是否启动成功
return workerStarted;
}

该方法主要完成这几件事情:

  • 原子性的增加workerCount
  • 将用户给定的任务封装成为一个worker,并将此worker添加进workers集合中
  • 启动worker对应的线程
  • 若线程启动失败,回滚worker的创建动作,即从worker中溢出新添加的worker,并原子性的减少workerCount。

工作线程

从之前的代码中,我们可以发现,线程池中正在执行任务的是worker对象。

1
2
3
private final class Worker
extends AbstractQueuedSynchronizer
implements Runnable

worker的属性:

1
2
3
4
5
6
//用来封装worker的线程,线程池中真正运行的线程,通过线程工厂创建而来
final Thread thread;
//worker所对应的第一个任务,可能为空
Runnable firstTask;
//记录当前线程完成的任务数
volatile long completedTasks;

worker的构造器:

1
2
3
4
5
6
7
8
9
Worker(Runnable firstTask) {
//设置AQS的state为-1,在执行runWorker()方法之前阻止线程中断
setState(-1);
//初始化第一个任务
this.firstTask = firstTask;
//利用指定的线程工厂创建一个线程,注意,参数是Worker实例本身this
//也就是当执行start方法启动线程thread时,真正执行的是Worker类的run方法
this.thread = getThreadFactory().newThread(this);
}

worker类继承了AQS类,重写了其相应的方法,是实现了一个自定义的同步器,实现了不可重入锁。

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
//是否持有独占锁
protected boolean isHeldExclusively() {
return getState() != 0;
}
//尝试获取锁
protected boolean tryAcquire(int unused) {
if (compareAndSetState(0, 1)) {
//设置独占线程
setExclusiveOwnerThread(Thread.currentThread());
return true;
}
return false;
}
//尝试释放锁
protected boolean tryRelease(int unused) {
//设置独占线程为null
setExclusiveOwnerThread(null);
setState(0);
return true;
}
//获取锁
public void lock() { acquire(1); }
//尝试获取锁
public boolean tryLock() { return tryAcquire(1); }
//释放锁
public void unlock() { release(1); }
//是否持有锁
public boolean isLocked() { return isHeldExclusively(); }

worker类还提供了一个中断线程thread的方法

1
2
3
4
5
6
7
8
9
10
void interruptIfStarted() {
Thread t;
//AQS状态大于等于0,worker对应的线程不为null,且该线程没有被中断
if (getState() >= 0 && (t = thread) != null && !t.isInterrupted()) {
try {
t.interrupt();
} catch (SecurityException ignore) {
}
}
}

线程复用机制

worker中的线程start后,实际上是调用的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
47
48
49
50
51
52
53
54
55
56
57
58
59
final void runWorker(Worker w) {
//获取当前线程
Thread wt = Thread.currentThread();
//获取w的firstTask
Runnable task = w.firstTask;
//设置w的firstTask为null
w.firstTask = null;
// 释放锁,设置AQS的state为0,允许中断
w.unlock();
//用于标识线程是否异常终止,finally中processWorkerExit()方法会有不同逻辑
boolean completedAbruptly = true;
try {
//循环调用getTask()获取任务,不断从任务缓存队列获取任务并执行
while (task != null || (task = getTask()) != null) {
//进入循环内部,代表已经获取到可执行的任务,则对worker对象加锁,保证线程在执行任务过程中不会被中断
w.lock();
if ((runStateAtLeast(ctl.get(), STOP) || //若线程池状态大于等于STOP,那么意味着该线程要中断
(Thread.interrupted() && //线程被中断
runStateAtLeast(ctl.get(), STOP))) && //且是因为线程池内部状态变化而被中断
!wt.isInterrupted()) //确保该线程未被中断
//发出中断请求
wt.interrupt();
try {
//开始执行任务前的Hook方法
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 {
//执行任务后的Hook方法
afterExecute(task, thrown);
}
} finally {
//置空task,准备通过getTask()获取下一个任务
task = null;
//completedTasks递增
w.completedTasks++;
//释放掉worker持有的独占锁
w.unlock();
}
}
completedAbruptly = false;
} finally {
//到这里,线程执行结束,需要执行结束线程的一些清理工作
//线程执行结束可能有两种情况:
//1.getTask()返回null,也就是说,这个worker的使命结束了,线程执行结束
//2.任务执行过程中发生了异常
//第一种情况,getTask()返回null,那么getTask()中会将workerCount递减
//第二种情况,workerCount没有进行处理,这个递减操作会在processWorkerExit()中处理
processWorkerExit(w, completedAbruptly);
}
}

该方法主要完成了这些事情:

  • 运行第一个任务firstTask之后,循环调用getTask方法从同步队列中获取任务并执行。
  • 获取到任务就对worker对象加锁,保证线程在执行任务的过程中不会被中断,任务执行完会释放锁。
  • 在执行任务的前后,可以根据业务场景重写beforeExecuteafterExecute等钩子函数。
  • 线程执行结束后,调用processWorkerExit方法执行结束的一些清理工作。

从队列中取任务

从之前的源码中我们可以看出 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
39
40
41
42
43
44
45
46
47
48
49
50
51
private Runnable getTask() {
//标识当前线程是否超时未能获取到task对象
boolean timedOut = false;

for (;;) {
//获取线程池的控制状态
int c = ctl.get();
//获取线程池的运行状态
int rs = runStateOf(c);

//如果线程池状态大于等于STOP,或者处于SHUTDOWN状态,并且阻塞队列为空,线程池工作线程数量递减,方法返回null,回收线程
if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) {
decrementWorkerCount();
return null;
}

//获取worker数量
int wc = workerCountOf(c);

//标识当前线程在空闲时,是否应该超时回收
// 如果allowCoreThreadTimeOut为ture,或当前线程数大于核心池大小,则需要超时回收
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;

//如果worker数量大于maximumPoolSize(有可能调用了 setMaximumPoolSize(),导致worker数量大于maximumPoolSize)
if ((wc > maximumPoolSize || (timed && timedOut)) //或者获取任务超时
&& (wc > 1 || workQueue.isEmpty())) { //workerCount大于1或者阻塞队列为空(在阻塞队列不为空时,需要保证至少有一个工作线程)
if (compareAndDecrementWorkerCount(c))
//线程池工作线程数量递减,方法返回null,回收线程
return null;
//线程池工作线程数量递减失败,跳过剩余部分,继续循环
continue;
}

try {
//如果允许超时回收,则调用阻塞队列的poll(),只在keepAliveTime时间内等待获取任务,一旦超过则返回null
//否则调用take(),如果队列为空,线程进入阻塞状态,无限时等待任务,直到队列中有可取任务或者响应中断信号退出
Runnable r = timed ?
workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) :
workQueue.take();
//若task不为null,则返回成功获取的task对象
if (r != null)
return r;
// 若返回task为null,表示线程空闲时间超时,则设置timeOut为true
timedOut = true;
} catch (InterruptedException retry) {
//如果此worker发生了中断,采取的方案是重试,没有超时
//在哪些情况下会发生中断?调用setMaximumPoolSize(),shutDown(),shutDownNow()
timedOut = false;
}
}
}

该方法在不同地情况下有不同地返回:

  1. 线程池处于running状态,阻塞队列不为空,该方法返回task对象。
  2. 线程池处于shutdown状态,阻塞队列不为空,返回空间获取地task对象。
  3. 线程池状态大于等于stop状态,返回null,回收线程
  4. 线程池处于shutdown状态,阻塞队列位空,返回null,回收线程
  5. worker数量大于最大线程数,返回null,回收线程。
  6. 线程空闲时间超时,返回null,回收线程。