線程池
[TOC]
線程池概述
什么是線程池
為什么使用線程池
-
線程池的優(yōu)勢
第一:降低資源消耗。通過重復(fù)利用已創(chuàng)建的線程降低線程創(chuàng)建和銷毀造成的消耗。
第二:提高響應(yīng)速度。當(dāng)任務(wù)到達(dá)時(shí),任務(wù)可以不需要的等到線程創(chuàng)建就能立即執(zhí)行。
第三:提高線程的可管理性。線程是稀缺資源,如果無限制的創(chuàng)建,不僅會消耗系統(tǒng)資源,還會降低系統(tǒng)的穩(wěn)定性,使用線程池可以進(jìn)行統(tǒng)一的分配,調(diào)優(yōu)和監(jiān)控。但是要做到合理的利用線程池,必須對其原理了如指掌。
創(chuàng)建一個(gè)線程池并提交線程任務(wù)
線程池源碼解析
參數(shù)認(rèn)識
corePoolSize : 線程池的基本大小,當(dāng)提交一個(gè)任務(wù)到線程池時(shí),線程池會創(chuàng)建一個(gè)線程來執(zhí)行任務(wù),即使其他空閑的基本線程能夠執(zhí)行新任務(wù)也會創(chuàng)建線程,等到需要執(zhí)行的任務(wù)數(shù)大于線程池基本大小時(shí)就不再創(chuàng)建。如果調(diào)用了線程池的prestartAllCoreThreads方法,線程池會提前創(chuàng)建并啟動所有基本線程。
runnableTaskQueue:任務(wù)對列,用于保存等待執(zhí)行的任務(wù)的阻塞隊(duì)列。可以選擇以下幾個(gè)阻塞隊(duì)列。
ArrayBlockingQueue:是一個(gè)基于數(shù)組結(jié)構(gòu)的有界阻塞隊(duì)列,此隊(duì)列按 FIFO(先進(jìn)先出)原則對元素進(jìn)行排序。
LinkedBlockingQueue:一個(gè)基于鏈表結(jié)構(gòu)的阻塞隊(duì)列,此隊(duì)列按FIFO (先進(jìn)先出) 排序元素,吞吐量通常要高于ArrayBlockingQueue。靜態(tài)工廠方法Executors.newFixedThreadPool()使用了這個(gè)隊(duì)列。
SynchronousQueue:一個(gè)不存儲元素的阻塞隊(duì)列。每個(gè)插入操作必須等到另一個(gè)線程調(diào)用移除操作,否則插入操作一直處于阻塞狀態(tài),吞吐量通常要高于LinkedBlockingQueue,靜態(tài)工廠方法Executors.newCachedThreadPool使用了這個(gè)隊(duì)列。
PriorityBlockingQueue:一個(gè)具有優(yōu)先級得無限阻塞隊(duì)列。
maximumPoolSize:線程池最大大小,線程池允許創(chuàng)建的最大線程數(shù)。如果隊(duì)列滿了,并且已創(chuàng)建的線程數(shù)小于最大線程數(shù),則線程池會再創(chuàng)建新的線程執(zhí)行任務(wù)。值得注意的是如果使用了無界的任務(wù)隊(duì)列這個(gè)參數(shù)就沒什么效果。
ThreadFactory:用于設(shè)置創(chuàng)建線程的工廠,可以通過線程工廠給每個(gè)創(chuàng)建出來的線程設(shè)置更有意義的名字,Debug和定位問題時(shí)非常又幫助。
RejectedExecutionHandler(飽和策略):當(dāng)隊(duì)列和線程池都滿了,說明線程池處于飽和狀態(tài),那么必須采取一種策略處理提交的新任務(wù)。這個(gè)策略默認(rèn)情況下是AbortPolicy,表示無法處理新任務(wù)時(shí)拋出異常。
CallerRunsPolicy:只用調(diào)用者所在線程來運(yùn)行任務(wù)。
DiscardOldestPolicy:丟棄隊(duì)列里最近的一個(gè)任務(wù),并執(zhí)行當(dāng)前任務(wù)。
DiscardPolicy:不處理,丟棄掉。
當(dāng)然也可以根據(jù)應(yīng)用場景需要來實(shí)現(xiàn)RejectedExecutionHandler接口自定義策略。如記錄日志或持久化不能處理的任務(wù)。
keepAliveTime :線程活動保持時(shí)間,線程池的工作線程空閑后,保持存活的時(shí)間。所以如果任務(wù)很多,并且每個(gè)任務(wù)執(zhí)行的時(shí)間比較短,可以調(diào)大這個(gè)時(shí)間,提高線程的利用率。
TimeUnit:線程活動保持時(shí)間的單位,可選的單位有天(DAYS),小時(shí)(HOURS),分鐘(MINUTES),毫秒(MILLISECONDS),微秒(MICROSECONDS, 千分之一毫秒)和毫微秒(NANOSECONDS, 千分之一微秒)。
類中其他屬性
// 線程池的控制狀態(tài):用來表示線程池的運(yùn)行狀態(tài)(整型的高3位)和運(yùn)行的worker數(shù)量(低29位)
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
// 29位的偏移量
private static final int COUNT_BITS = Integer.SIZE - 3;
// 最大容量(2^29 - 1)
private static final int CAPACITY = (1 << COUNT_BITS) - 1;
// runState is stored in the high-order bits
// 線程運(yùn)行狀態(tài),總共有5個(gè)狀態(tài),需要3位來表示(所以偏移量的29 = 32 - 3)
/**
* RUNNING : 接受新任務(wù)并且處理已經(jīng)進(jìn)入阻塞隊(duì)列的任務(wù)
* SHUTDOWN : 不接受新任務(wù),但是處理已經(jīng)進(jìn)入阻塞隊(duì)列的任務(wù)
* STOP : 不接受新任務(wù),不處理已經(jīng)進(jìn)入阻塞隊(duì)列的任務(wù)并且中斷正在運(yùn)行的任務(wù)
* TIDYING : 所有的任務(wù)都已經(jīng)終止,workerCount為0, 線程轉(zhuǎn)化為TIDYING狀態(tài)并且調(diào)用terminated鉤子函數(shù)
* TERMINATED: terminated鉤子函數(shù)已經(jīng)運(yùn)行完成
**/
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;
// 阻塞隊(duì)列
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;
// 已完成任務(wù)數(shù)量
private long completedTaskCount;
// 線程工廠
private volatile ThreadFactory threadFactory;
// 拒絕執(zhí)行處理器
private volatile RejectedExecutionHandler handler;
// 線程等待運(yùn)行時(shí)間
private volatile long keepAliveTime;
// 是否運(yùn)行核心線程超時(shí)
private volatile boolean allowCoreThreadTimeOut;
// 核心池的大小
private volatile int corePoolSize;
// 最大線程池大小
private volatile int maximumPoolSize;
// 默認(rèn)拒絕執(zhí)行處理器
private static final RejectedExecutionHandler defaultHandler =
new AbortPolicy();
構(gòu)造方法
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory,
RejectedExecutionHandler handler) {
if (corePoolSize < 0 || // 核心大小不能小于0
maximumPoolSize <= 0 || // 線程池的初始最大容量不能小于0
maximumPoolSize < corePoolSize || // 初始最大容量不能小于核心大小
keepAliveTime < 0) // keepAliveTime不能小于0
throw new IllegalArgumentException();
if (workQueue == null || threadFactory == null || handler == null)
throw new NullPointerException();
// 初始化相應(yīng)的域
this.corePoolSize = corePoolSize;
this.maximumPoolSize = maximumPoolSize;
this.workQueue = workQueue;
this.keepAliveTime = unit.toNanos(keepAliveTime);
this.threadFactory = threadFactory;
this.handler = handler;
}
提交任務(wù)
/*
* 進(jìn)行下面三步
*
* 1. 如果運(yùn)行的線程小于corePoolSize,則嘗試使用用戶定義的Runnalbe對象創(chuàng)建一個(gè)新的線程
* 調(diào)用addWorker函數(shù)會原子性的檢查runState和workCount,通過返回false來防止在不應(yīng)
* 該添加線程時(shí)添加了線程
* 2. 如果一個(gè)任務(wù)能夠成功入隊(duì)列,在添加一個(gè)線城時(shí)仍需要進(jìn)行雙重檢查(因?yàn)樵谇耙淮螜z查后
* 該線程死亡了),或者當(dāng)進(jìn)入到此方法時(shí),線程池已經(jīng)shutdown了,所以需要再次檢查狀態(tài),
* 若有必要,當(dāng)停止時(shí)還需要回滾入隊(duì)列操作,或者當(dāng)線程池沒有線程時(shí)需要創(chuàng)建一個(gè)新線程
* 3. 如果無法入隊(duì)列,那么需要增加一個(gè)新線程,如果此操作失敗,那么就意味著線程池已經(jīng)shut
* down或者已經(jīng)飽和了,所以拒絕任務(wù)
*/
public void execute(Runnable command) {
if (command == null)
throw new NullPointerException();
// 獲取線程池控制狀態(tài)
int c = ctl.get();
if (workerCountOf(c) < corePoolSize) { // worker數(shù)量小于corePoolSize
if (addWorker(command, true)) // 添加worker
// 成功則返回
return;
// 不成功則再次獲取線程池控制狀態(tài)
c = ctl.get();
}
// 線程池處于RUNNING狀態(tài),將用戶自定義的Runnable對象添加進(jìn)workQueue隊(duì)列
if (isRunning(c) && workQueue.offer(command)) {
// 再次檢查,獲取線程池控制狀態(tài)
int recheck = ctl.get();
// 線程池不處于RUNNING狀態(tài),將自定義任務(wù)從workQueue隊(duì)列中移除
if (! isRunning(recheck) && remove(command))
// 拒絕執(zhí)行命令
reject(command);
else if (workerCountOf(recheck) == 0) // worker數(shù)量等于0
// 添加worker
addWorker(null, false);
}
else if (!addWorker(command, false)) // 添加worker失敗
// 拒絕執(zhí)行命令
reject(command);
}
addWorker
原子性的增加workerCount。
將用戶給定的任務(wù)封裝成為一個(gè)worker,并將此worker添加進(jìn)workers集合中。
啟動worker對應(yīng)的線程,并啟動該線程,運(yùn)行worker的run方法。
回滾worker的創(chuàng)建動作,即將worker從workers集合中刪除,并原子性的減少workerCount。
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (;;) { // 外層無限循環(huán)
// 獲取線程池控制狀態(tài)
int c = ctl.get();
// 獲取狀態(tài)
int rs = runStateOf(c);
// Check if queue empty only if necessary.
if (rs >= SHUTDOWN && // 狀態(tài)大于等于SHUTDOWN,初始的ctl為RUNNING,小于SHUTDOWN
! (rs == SHUTDOWN && // 狀態(tài)為SHUTDOWN
firstTask == null && // 第一個(gè)任務(wù)為null
! workQueue.isEmpty())) // worker隊(duì)列不為空
// 返回
return false;
for (;;) {
// worker數(shù)量
int wc = workerCountOf(c);
if (wc >= CAPACITY || // worker數(shù)量大于等于最大容量
wc >= (core ? corePoolSize : maximumPoolSize)) // worker數(shù)量大于等于核心線程池大小或者最大線程池大小
return false;
if (compareAndIncrementWorkerCount(c)) // 比較并增加worker的數(shù)量
// 跳出外層循環(huán)
break retry;
// 獲取線程池控制狀態(tài)
c = ctl.get(); // Re-read ctl
if (runStateOf(c) != rs) // 此次的狀態(tài)與上次獲取的狀態(tài)不相同
// 跳過剩余部分,繼續(xù)循環(huán)
continue retry;
// else CAS failed due to workerCount change; retry inner loop
}
}
// worker開始標(biāo)識
boolean workerStarted = false;
// worker被添加標(biāo)識
boolean workerAdded = false;
//
Worker w = null;
try {
// 初始化worker
w = new Worker(firstTask);
// 獲取worker對應(yīng)的線程
final Thread t = w.thread;
if (t != null) { // 線程不為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.
// 線程池的運(yùn)行狀態(tài)
int rs = runStateOf(ctl.get());
if (rs < SHUTDOWN || // 小于SHUTDOWN
(rs == SHUTDOWN && firstTask == null)) { // 等于SHUTDOWN并且firstTask為null
if (t.isAlive()) // precheck that t is startable // 線程剛添加進(jìn)來,還未啟動就存活
// 拋出線程狀態(tài)異常
throw new IllegalThreadStateException();
// 將worker添加到worker集合
workers.add(w);
// 獲取worker集合的大小
int s = workers.size();
if (s > largestPoolSize) // 隊(duì)列大小大于largestPoolSize
// 重新設(shè)置largestPoolSize
largestPoolSize = s;
// 設(shè)置worker已被添加標(biāo)識
workerAdded = true;
}
} finally {
// 釋放鎖
mainLock.unlock();
}
if (workerAdded) { // worker被添加
// 開始執(zhí)行worker的run方法
t.start();
// 設(shè)置worker已開始標(biāo)識
workerStarted = true;
}
}
} finally {
if (! workerStarted) // worker沒有開始
// 添加worker失敗
addWorkerFailed(w);
}
return workerStarted;
}
執(zhí)行任務(wù)
runWorker函數(shù)中會實(shí)際執(zhí)行給定任務(wù)(即調(diào)用用戶重寫的run方法),并且當(dāng)給定任務(wù)完成后,會繼續(xù)從阻塞隊(duì)列中取任務(wù),直到阻塞隊(duì)列為空(即任務(wù)全部完成)。在執(zhí)行給定任務(wù)時(shí),會調(diào)用鉤子函數(shù),利用鉤子函數(shù)可以完成用戶自定義的一些邏輯。在runWorker中會調(diào)用到getTask函數(shù)和processWorkerExit鉤子函數(shù)
final void runWorker(Worker w) {
// 獲取當(dāng)前線程
Thread wt = Thread.currentThread();
// 獲取w的firstTask
Runnable task = w.firstTask;
// 設(shè)置w的firstTask為null
w.firstTask = null;
// 釋放鎖(設(shè)置state為0,允許中斷)
w.unlock(); // allow interrupts
boolean completedAbruptly = true;
try {
while (task != null || (task = getTask()) != null) { // 任務(wù)不為null或者阻塞隊(duì)列還存在任務(wù)
// 獲取鎖
w.lock();
// If pool is stopping, ensure thread is interrupted;
// if not, ensure thread is not interrupted. This
// requires a recheck in second case to deal with
// shutdownNow race while clearing interrupt
if ((runStateAtLeast(ctl.get(), STOP) || // 線程池的運(yùn)行狀態(tài)至少應(yīng)該高于STOP
(Thread.interrupted() && // 線程被中斷
runStateAtLeast(ctl.get(), STOP))) && // 再次檢查,線程池的運(yùn)行狀態(tài)至少應(yīng)該高于STOP
!wt.isInterrupted()) // wt線程(當(dāng)前線程)沒有被中斷
wt.interrupt(); // 中斷wt線程(當(dāng)前線程)
try {
// 在執(zhí)行之前調(diào)用鉤子函數(shù)
beforeExecute(wt, task);
Throwable thrown = null;
try {
// 運(yùn)行給定的任務(wù)
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 {
// 執(zhí)行完后調(diào)用鉤子函數(shù)
afterExecute(task, thrown);
}
} finally {
task = null;
// 增加給worker完成的任務(wù)數(shù)量
w.completedTasks++;
// 釋放鎖
w.unlock();
}
}
completedAbruptly = false;
} finally {
// 處理完成后,調(diào)用鉤子函數(shù)
processWorkerExit(w, completedAbruptly);
}
}
此函數(shù)用于從workerQueue阻塞隊(duì)列中獲取Runnable對象,由于是阻塞隊(duì)列,所以支持有限時(shí)間等待(poll)和無限時(shí)間等待(take)。在該函數(shù)中還會響應(yīng)shutDown和、shutDownNow函數(shù)的操作,若檢測到線程池處于SHUTDOWN或STOP狀態(tài),則會返回null,而不再返回阻塞隊(duì)列中的Runnalbe對象。
private Runnable getTask() {
boolean timedOut = false; // Did the last poll() time out?
for (;;) { // 無限循環(huán),確保操作成功
// 獲取線程池控制狀態(tài)
int c = ctl.get();
// 運(yùn)行的狀態(tài)
int rs = runStateOf(c);
// Check if queue empty only if necessary.
if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) { // 大于等于SHUTDOWN(表示調(diào)用了shutDown)并且(大于等于STOP(調(diào)用了shutDownNow)或者worker阻塞隊(duì)列為空)
// 減少worker的數(shù)量
decrementWorkerCount();
// 返回null,不執(zhí)行任務(wù)
return null;
}
// 獲取worker數(shù)量
int wc = workerCountOf(c);
// Are workers subject to culling?
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize; // 是否允許coreThread超時(shí)或者workerCount大于核心大小
if ((wc > maximumPoolSize || (timed && timedOut)) // worker數(shù)量大于maximumPoolSize
&& (wc > 1 || workQueue.isEmpty())) { // workerCount大于1或者worker阻塞隊(duì)列為空(在阻塞隊(duì)列不為空時(shí),需要保證至少有一個(gè)wc)
if (compareAndDecrementWorkerCount(c)) // 比較并減少workerCount
// 返回null,不執(zhí)行任務(wù),該worker會退出
return null;
// 跳過剩余部分,繼續(xù)循環(huán)
continue;
}
try {
Runnable r = timed ?
workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) : // 等待指定時(shí)間
workQueue.take(); // 一直等待,直到有元素
if (r != null)
return r;
// 等待指定時(shí)間后,沒有獲取元素,則超時(shí)
timedOut = true;
} catch (InterruptedException retry) {
// 拋出了被中斷異常,重試,沒有超時(shí)
timedOut = false;
}
}
}
processWorkerExit函數(shù)是在worker退出時(shí)調(diào)用到的鉤子函數(shù),而引起worker退出的主要因素如下
阻塞隊(duì)列已經(jīng)為空,即沒有任務(wù)可以運(yùn)行了。
調(diào)用了shutDown或shutDownNow函數(shù)
此函數(shù)會根據(jù)是否中斷了空閑線程來確定是否減少workerCount的值,并且將worker從workers集合中移除并且會嘗試終止線程池。
private void processWorkerExit(Worker w, boolean completedAbruptly) {
if (completedAbruptly) // 如果被中斷,則需要減少workCount // If abrupt, then workerCount wasn't adjusted
decrementWorkerCount();
// 獲取可重入鎖
final ReentrantLock mainLock = this.mainLock;
// 獲取鎖
mainLock.lock();
try {
// 將worker完成的任務(wù)添加到總的完成任務(wù)中
completedTaskCount += w.completedTasks;
// 從workers集合中移除該worker
workers.remove(w);
} finally {
// 釋放鎖
mainLock.unlock();
}
// 嘗試終止
tryTerminate();
// 獲取線程池控制狀態(tài)
int c = ctl.get();
if (runStateLessThan(c, STOP)) { // 小于STOP的運(yùn)行狀態(tài)
if (!completedAbruptly) {
int min = allowCoreThreadTimeOut ? 0 : corePoolSize;
if (min == 0 && ! workQueue.isEmpty()) // 允許核心超時(shí)并且workQueue阻塞隊(duì)列不為空
min = 1;
if (workerCountOf(c) >= min) // workerCount大于等于min
// 直接返回
return; // replacement not needed
}
// 添加worker
addWorker(null, false);
}
}
關(guān)閉線程池
public void shutdown() {
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
// 檢查shutdown權(quán)限
checkShutdownAccess();
// 設(shè)置線程池控制狀態(tài)為SHUTDOWN
advanceRunState(SHUTDOWN);
// 中斷空閑worker
interruptIdleWorkers();
// 調(diào)用shutdown鉤子函數(shù)
onShutdown(); // hook for ScheduledThreadPoolExecutor
} finally {
mainLock.unlock();
}
// 嘗試終止
tryTerminate();
}
final void tryTerminate() {
for (;;) { // 無限循環(huán),確保操作成功
// 獲取線程池控制狀態(tài)
int c = ctl.get();
if (isRunning(c) || // 線程池的運(yùn)行狀態(tài)為RUNNING
runStateAtLeast(c, TIDYING) || // 線程池的運(yùn)行狀態(tài)最小要大于TIDYING
(runStateOf(c) == SHUTDOWN && ! workQueue.isEmpty())) // 線程池的運(yùn)行狀態(tài)為SHUTDOWN并且workQueue隊(duì)列不為null
// 不能終止,直接返回
return;
if (workerCountOf(c) != 0) { // 線程池正在運(yùn)行的worker數(shù)量不為0 // Eligible to terminate
// 僅僅中斷一個(gè)空閑的worker
interruptIdleWorkers(ONLY_ONE);
return;
}
// 獲取線程池的鎖
final ReentrantLock mainLock = this.mainLock;
// 獲取鎖
mainLock.lock();
try {
if (ctl.compareAndSet(c, ctlOf(TIDYING, 0))) { // 比較并設(shè)置線程池控制狀態(tài)為TIDYING
try {
// 終止,鉤子函數(shù)
terminated();
} finally {
// 設(shè)置線程池控制狀態(tài)為TERMINATED
ctl.set(ctlOf(TERMINATED, 0));
// 釋放在termination條件上等待的所有線程
termination.signalAll();
}
return;
}
} finally {
// 釋放鎖
mainLock.unlock();
}
// else retry on failed CAS
}
}
private void interruptIdleWorkers(boolean onlyOne) {
// 線程池的鎖
final ReentrantLock mainLock = this.mainLock;
// 獲取鎖
mainLock.lock();
try {
for (Worker w : workers) { // 遍歷workers隊(duì)列
// worker對應(yīng)的線程
Thread t = w.thread;
if (!t.isInterrupted() && w.tryLock()) { // 線程未被中斷并且成功獲得鎖
try {
// 中斷線程
t.interrupt();
} catch (SecurityException ignore) {
} finally {
// 釋放鎖
w.unlock();
}
}
if (onlyOne) // 若只中斷一個(gè),則跳出循環(huán)
break;
}
} finally {
// 釋放鎖
mainLock.unlock();
}
}