概述
线程池故名思意就是存放线程的池子。线程是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
| private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0)); private static final int COUNT_BITS = Integer.SIZE - 3; private static final int CAPACITY = (1 << COUNT_BITS) - 1;
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;
private volatile int corePoolSize;
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) { if (addWorker(command, true)) return; c = ctl.get(); } if (isRunning(c) && workQueue.offer(command)) { int recheck = ctl.get(); if (! isRunning(recheck) && remove(command)) reject(command); else if (workerCountOf(recheck) == 0) addWorker(null, false); } else if (!addWorker(command, false)) reject(command); }
|
整个流程如下:
- 如果线程池中的工作线程数小于corePoolSize,则创建新线程来执行任务
- 如果工作线程数量大于或等于corePoolSize,则将任务加入BlockingQueue
- 若无法将任务加入阻塞队列,且工作线程数小于最大线程数,则创建新的线程来执行任务
- 当工作线程数量达到最大线程数,则创建工作线程失败,采取任务拒绝策略

创建线程
我们在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
|
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)) break retry; c = ctl.get(); if (runStateOf(c) != rs) continue retry; } }
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 { int rs = runStateOf(ctl.get());
if (rs < SHUTDOWN || (rs == SHUTDOWN && firstTask == null)) { if (t.isAlive()) 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; }
|
该方法主要完成这几件事情:
- 原子性的增加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
| final Thread thread; Runnable firstTask; volatile long completedTasks;
|
worker的构造器:
1 2 3 4 5 6 7 8 9
| Worker(Runnable firstTask) { setState(-1); this.firstTask = firstTask; 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) { 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; 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(); Runnable task = w.firstTask; w.firstTask = null; w.unlock(); boolean completedAbruptly = true; try { 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); } }
|
该方法主要完成了这些事情:
- 运行第一个任务
firstTask之后,循环调用getTask方法从同步队列中获取任务并执行。
- 获取到任务就对worker对象加锁,保证线程在执行任务的过程中不会被中断,任务执行完会释放锁。
- 在执行任务的前后,可以根据业务场景重写
beforeExecute和afterExecute等钩子函数。
- 线程执行结束后,调用
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() { boolean timedOut = false;
for (;;) { int c = ctl.get(); int rs = runStateOf(c);
if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) { decrementWorkerCount(); return null; }
int wc = workerCountOf(c);
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;
if ((wc > maximumPoolSize || (timed && timedOut)) && (wc > 1 || workQueue.isEmpty())) { if (compareAndDecrementWorkerCount(c)) 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; } } }
|
该方法在不同地情况下有不同地返回:
- 线程池处于running状态,阻塞队列不为空,该方法返回task对象。
- 线程池处于shutdown状态,阻塞队列不为空,返回空间获取地task对象。
- 线程池状态大于等于stop状态,返回null,回收线程
- 线程池处于shutdown状态,阻塞队列位空,返回null,回收线程
- worker数量大于最大线程数,返回null,回收线程。
- 线程空闲时间超时,返回null,回收线程。