Java线程池ThreadPoolExecutor
当需要使用大量线程执行任务时,使用线程池可以提供较好的性能;如果不使用线程池,每个任务创建一个线程,而线程的创建和销毁会带来性能开销
一、简介
Java并发包中提供了线程池相关功能和类,使用线程池主要有以下几方面的作用
- 性能,当需要使用大量线程执行任务时,使用线程池可以提供较好的性能;如果不使用线程池,每个任务创建一个线程,而线程的创建和销毁会带来性能开销;使用线程池可以减少线程创建销毁带来的开销,因为线程池里面的线程是可以复用的,不用每次执行任务都重新创建和销毁线程。
- 资源限制和管理,线程池提供了限制线程资源、管理线程的手段,例如限制线程数、动态新增线程、线程执行任务数监控等。
- 可扩展性,线程池提供了多个可配置参数和可扩展接口,以满足不同的使用场景。
二、线程池对象
ThreadPoolExecutor.java表示一个线程池,ThreadPoolExecutor继承自AbstractExecutorService,而AbstractExecutorService实现了ExecutorService接口
ExecutorService.java
public interface ExecutorService extends Executor, AutoCloseable {
<T> Future<T> submit(Callable<T> task);
<T> Future<T> submit(Runnable task, T result);
Future<?> submit(Runnable task);
}
ExecutorService核心方法submit用来提交任务给线程池执行,并且提供了不同的重载满足不同的场景
可以看到ExecutorService继承了Executor
Executor.java
public interface Executor {
/**
* Executes the given command at some time in the future. The command
* may execute in a new thread, in a pooled thread, or in the calling
* thread, at the discretion of the {@code Executor} implementation.
*
* @param command the runnable task
* @throws RejectedExecutionException if this task cannot be
* accepted for execution
* @throws NullPointerException if command is null
*/
void execute(Runnable command);
}
Executor的execute方法用来执行异步任务。
三、线程池参数
ThreadPoolExecutor.java
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory,
RejectedExecutionHandler handler) {
if (corePoolSize < 0 ||
maximumPoolSize <= 0 ||
maximumPoolSize < corePoolSize ||
keepAliveTime < 0)
throw new IllegalArgumentException();
if (workQueue == null || threadFactory == null || handler == null)
throw new NullPointerException();
this.corePoolSize = corePoolSize;
this.maximumPoolSize = maximumPoolSize;
this.workQueue = workQueue;
this.keepAliveTime = unit.toNanos(keepAliveTime);
this.threadFactory = threadFactory;
this.handler = handler;
String name = Objects.toIdentityString(this);
this.container = SharedThreadContainer.create(name);
}
corePoolSize:核心线程数,也就是线程池可以持有的线程个数,核心线程默认情况下就算空闲也不会被销毁,可以通过设置allowCoreThreadTimeOut为true允许空闲的核心线程销毁。默认情况下,线程池创建时并不会立即创建corePoolSize个线程,可以调用prestartCoreThread()方法或者prestartAllCoreThreads()提前创建好线程
maximumPoolSize:线程池能够持有的最大线程个数,核心线程都被使用后如果还有任务提交给线程池,则线程池还可以创建最多maximumPoolSize减corePoolSize个线程来执行任务,给线程池提交任务时,如果所有核心线程都活动中,则任务被添加进阻塞队列,如果队列满了,则判断线程数是否小于maximumPoolSize,是则创建新线程执行该任务
keepAliveTime:超出核心线程的那部分线程存活时间,当线程池持有的线程数大于核心线程数时,超出核心线程的这部分线程在等待执行任务的最大时间,从线程空闲开始等待keepAliveTime时间还没有任务提交给该线程执行,则该线程将会被销毁。如果keepAliveTime为0,则超出corePoolSize的线程空闲时将会立即被销毁
unit:keepAliveTime参数的时间单位,TimeUnit对象
workQueue:工作队列,用来存储提交给线程池的,被线程池执行之前的任务,这个队列只存储被execute方法提交的Runable类型的任务
threadFactory:线程工厂,当线程池需要创建线程时(核心线程都被占用),调用threadFactory的newThread(Runnable r)方法创建线程
handler:任务执行拒绝策略,当线程池的被使用的线程已经达到maximumPoolSize,并且workQueue队列已满的情况下,提交新的任务给线程池执行时使用的处理策略
参数出现以下几种情况构造方法将抛出IllegalArgumentException 异常:
corePoolSize < 0keepAliveTime < 0maximumPoolSize <= 0maximumPoolSize < corePoolSize
参数出现以下情况调用构造方法将抛出NullPointerException异常:
workQueue==nullthreadFactory==nullhandler==null
设置allowCoreThreadTimeOut参数:
默认情况下线程池的核心线程池是不会被销毁的,如果希望空闲的核心线程也被销毁则可以通过设置allowCoreThreadTimeOut参数为true实现。
ThreadPoolExecutor.java
public void allowCoreThreadTimeOut(boolean value) {
if (value && keepAliveTime <= 0) // (1)
throw new IllegalArgumentException("Core threads must have nonzero keep alive times");
if (value != allowCoreThreadTimeOut) {
allowCoreThreadTimeOut = value;
if (value)
interruptIdleWorkers(); // (2)
}
}
核心线程因为没有任务提交而销毁后,在接收到新的任务时线程池会重新创建核心线程。
代码(1),如果allowCoreThreadTimeOut设置为true的情况下,需要确保keepAliveTime>0,因为在allowCoreThreadTimeOut为true时,keepAliveTime对于非核心线程的超时检测作用对核心线程同样适用,为避免核心线程没有任务时直接关闭(没有任务时关闭,有任务时重新创建),将此值设为 true 时,keepAliveTime必须大于零。
代码(2)设置allowCoreThreadTimeOut为true后,将会中断线程池中正在等待任务的线程。
四、线程池状态
RUNNING:接收新任务并处理阻塞队列里的任务
SHUTDOWN:拒绝新任务但处理阻塞队列里的任务
STOP:拒绝新任务并抛弃阻塞队列里的任务,同时会中断正在处理的任务
TIDYING:所有任务(包括阻塞队里里的任务)都执行完后当前线程池的活动线程为0,将要调用terminated方法
TERMINATED:终止状态,terminated方法调用后的状态
线程池状态转换:
| 源状态 | 终状态 | 操作 |
|---|---|---|
| RUNNING | SHUTDOWN | 显式调用shutdown()方法,或者隐式调用finalize()方法里面的shutdown()方法 |
| RUNNING/SHUTDOWN | STOP | 显示调用shotdownNow()方法 |
| SHUTDOWN | TIDYING | 线程池和任务队列都为空 |
| STOP | TIDYING | 当线程池为空 |
| TIDYING | TERMINATED | 当terminated()方法执行完成时 |
线程池状态在ThreadPoolExecutor的表示:
private static final int COUNT_BITS = Integer.SIZE - 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;
可以看到TERMINATED>TIDYING>STOP>SHUTDOWN>RUNNING
五、创建线程池
ThreadPoolExecutor类提供了多个不同参数的构成方法用来创建不同类型的线程池
ThreadPoolExecutor.java
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue);
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory);
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
RejectedExecutionHandler handler);
public ThreadPoolExecutor(int corePoolSize,
int maximumPoolSize,
long keepAliveTime,
TimeUnit unit,
BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory,
RejectedExecutionHandler handler);
除此之外,jdk提供了一个工具类Executors,该类里面提供了很多静态方法用来创建不同类型的线程池,下面我们来看看Executors类提供的创建线程池的静态方法
1、newFixedThreadPool(int nThreads)
该方法创建一个核心线程数和最大线程数都是nThreads的线程池
public static ExecutorService newFixedThreadPool(int nThreads) {
return new ThreadPoolExecutor(nThreads, nThreads,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>());
}
可以看到设置keepAliveTime为0表示如果有线程数大于核心线程数则空闲的线程会被销毁,阻塞队列为LinkedBlockingQueue,其最大长度为Integer.MAX_VALUE。
2、newSingleThreadExecutor()
该方法创建一个核心线程数和最大线程数都是1的线程池
public static ExecutorService newSingleThreadExecutor() {
return newSingleThreadExecutor(defaultThreadFactory());
}
public static ExecutorService newSingleThreadExecutor(ThreadFactory threadFactory) {
return new AutoShutdownDelegatedExecutorService
(new ThreadPoolExecutor(1, 1,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>(),
threadFactory));
}
可以看到设置keepAliveTime为0表示如果有线程数大于核心线程数则空闲的线程会被销毁,阻塞队列为LinkedBlockingQueue,其最大长度为Integer.MAX_VALUE。
3、newCachedThreadPool()
该方法创建一个按需创建线程的线程池,核心线程数为0,最大可创建Integer.MAX_VALUE个线程
public static ExecutorService newCachedThreadPool() {
return new ThreadPoolExecutor(0, Integer.MAX_VALUE,
60L, TimeUnit.SECONDS,
new SynchronousQueue<Runnable>());
}
设置keepAliveTime为60表示如果有线程的空闲时间超过60秒则会被销毁,该线程池使用的是同步阻塞队列SynchronousQueue,向SynchronousQueue里面添加元素时,如果没有其他线程从SynchronousQueue取元素则执行入队的线程会被阻塞,直到有其他线程执行出队操作。
4、newVirtualThreadPerTaskExecutor()
该方法对应jdk-21提供的虚拟线程,使用该方法创建的线程池会为每个任务创建一个虚拟线程,并且创建的线程数没有限制
public static ExecutorService newVirtualThreadPerTaskExecutor() {
ThreadFactory factory = Thread.ofVirtual().factory();
return newThreadPerTaskExecutor(factory);
}
5、newScheduledThreadPool(int corePoolSize)
该方法创建一个线程池,核心线程数为corePoolSize,最大线程数为Integer.MAX_VALUE,该线程池可以在指定延迟时间后执行任务或者定时执行任务
public static ScheduledExecutorService newScheduledThreadPool(int corePoolSize) {
return new ScheduledThreadPoolExecutor(corePoolSize);
}
private static final long DEFAULT_KEEPALIVE_MILLIS = 10L;
public ScheduledThreadPoolExecutor(int corePoolSize) {
super(corePoolSize, Integer.MAX_VALUE,
DEFAULT_KEEPALIVE_MILLIS, MILLISECONDS,
new DelayedWorkQueue());
}
可以看到可以看到设置keepAliveTime为10表示如果有线程数大于核心线程数则空闲的线程在10秒后会被销毁,阻塞队列使用的是DelayedWorkQueue延时队列
六、DefaultThreadFactory
默认情况下,我们创建线程池时如果不指定线程工厂,则都会使用DefaultThreadFactory线程工厂创建线程,该类有一个变量threadNumber用来记录每个工厂创建的线程数,除此之外还有一个static修饰的类变量poolNumber,用来记录创建工厂对象的数量,从代码(1)和代码(2)可以看到这两个变量将会应用到工厂创建出来的线程名称上。
DefaultThreadFactory.java
private static class DefaultThreadFactory implements ThreadFactory {
private static final AtomicInteger poolNumber = new AtomicInteger(1);
private final ThreadGroup group;
private final AtomicInteger threadNumber = new AtomicInteger(1);
private final String namePrefix;
DefaultThreadFactory() {
@SuppressWarnings("removal")
SecurityManager s = System.getSecurityManager();
group = (s != null) ? s.getThreadGroup() :
Thread.currentThread().getThreadGroup();
namePrefix = "pool-" +
poolNumber.getAndIncrement() +
"-thread-"; // (1)
}
public Thread newThread(Runnable r) {
Thread t = new Thread(group, r,
namePrefix + threadNumber.getAndIncrement(), // (2)
0);
if (t.isDaemon())
t.setDaemon(false);
if (t.getPriority() != Thread.NORM_PRIORITY)
t.setPriority(Thread.NORM_PRIORITY);
return t;
}
}
七、execute方法提交任务
调用``execute方法将在未来的时间使用线程池里已存在的线程或者新创建一个线程执行command任务,如果线程池已关闭或者线程池容量已满,则将会调用线程池的RejectedExecutionHandler`处理任务。
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);
}
execute方法的执行步骤如下:
1、当前运行线程小于corePoolSize,则创建一个新的线程执行任务
2、否则把任务添加到队列里,任务成功添加到队列后,再次检查线程池的状态,如果队列不是RUNNING则把任务从队列里删除,然后拒绝任务
3、如果任务不能添加到阻塞队列里,将尝试创建一个新的线程执行该任务,如果失败,则会拒绝任务(注意只有队列满了才会创建核心线程以外的线程)
只有通过execute(Runnable)方法提交的任务才会被线程池的阻塞队列存储
八、创建线程执行任务
首先调用addWorker方法创建一个执行任务的线程,并添加到workers工作列表里:
private boolean addWorker(Runnable firstTask, boolean core) {
// (1)
retry:
for (int c = ctl.get();;) {
// Check if queue empty only if necessary.
if (runStateAtLeast(c, SHUTDOWN)
&& (runStateAtLeast(c, STOP)
|| firstTask != null
|| workQueue.isEmpty()))
return false;
for (;;) {
if (workerCountOf(c)
>= ((core ? corePoolSize : maximumPoolSize) & COUNT_MASK))
return false;
if (compareAndIncrementWorkerCount(c))
break retry;
c = ctl.get(); // Re-read ctl
if (runStateAtLeast(c, SHUTDOWN))
continue retry;
// else CAS failed due to workerCount change; retry inner loop
}
}
// (2)
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 c = ctl.get();
if (isRunning(c) ||
(runStateLessThan(c, STOP) && firstTask == null)) {
if (t.getState() != Thread.State.NEW)
throw new IllegalThreadStateException();
// 新创建的线程添加到workers工作集合里
workers.add(w);
workerAdded = true;
int s = workers.size();
if (s > largestPoolSize)
largestPoolSize = s;
}
} finally {
mainLock.unlock();
}
if (workerAdded) {
container.start(t);
workerStarted = true;
}
}
} finally {
if (! workerStarted)
addWorkerFailed(w);
}
return workerStarted;
}
addWorker代码分为两个部分,首先代码片段(1)循环使用CAS增加线程数,关键代码是compareAndIncrementWorkerCount(c),使用CAS算法增加ctl线程数的值;第二部分代码片段(2),创建新的线程执行任务。接下来我们看看创建新线程的方法Worker()构造方法:
private final class Worker
extends AbstractQueuedSynchronizer
implements Runnable
{
Worker(Runnable firstTask) {
setState(-1); // (1) inhibit interrupts until runWorker
this.firstTask = firstTask;
this.thread = getThreadFactory().newThread(this); // (2)
}
}
Worker类继承AbstractQueuedSynchronizer抽象同步队列和Runnable接口,代码(1)设置状态值state为-1,保证在创建线程过程中不允许中断Worker线程(在调用shutdownNow()方法时会中断state >=0的Worker线程),代码(2)使用创建线程池的时候指定的线程工厂创建一个线程
在创建好执行线程后,调用runWorker(Worker w)方法使用线程执行任务,从addWorker方法中可以看到创建线程成功后会调用container.start(t);方法启动线程,线程t会调用worker的run方法,run方法中再调用runWorker方法
final void runWorker(Worker w) {
Thread wt = Thread.currentThread();
Runnable task = w.firstTask;
w.firstTask = null;
w.unlock(); // (1) allow interrupts
boolean completedAbruptly = true;
try {
// (2)
while (task != null || (task = getTask()) != null) {
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) ||
(Thread.interrupted() &&
runStateAtLeast(ctl.get(), STOP))) &&
!wt.isInterrupted())
wt.interrupt();
try {
beforeExecute(wt, task);
try {
task.run();
afterExecute(task, null);
} catch (Throwable ex) {
afterExecute(task, ex);
throw ex;
}
} finally {
task = null;
w.completedTasks++;
w.unlock();
}
}
completedAbruptly = false;
} finally {
processWorkerExit(w, completedAbruptly);
}
}
可以看到,只要启动Worker线程后,代码(2)中的循环while (task != null || (task = getTask()) != null)会从队列取任务执行,直到队列为空后,会销毁空闲的Worker线程(如果设置了allowCoreThreadTimeOut参数为true,也会销毁空闲的核心线程)
需要注意的是代码(1)处调用w.unlock();方法把Worker的state状态值设置为0,前面我们说过创建Worker时为了避免Worker创建线程被中断所以把状态值state设置为-1,这里设置为0允许线程可以被中断。
九、销毁空闲线程
从前面的runWorker方法和getTask()方法可以看到,如果getTask()放回null,则runWorker方法将会返回,也就是worker线程的run方法返回,对应的worker线程会被销毁(run方法运行结束线程终止),从而达到销毁空闲线程的功能,我们来看什么情况下getTask()方法会返回null:
private Runnable getTask() {
boolean timedOut = false; // Did the last poll() time out?
for (;;) {
int c = ctl.get();
// Check if queue empty only if necessary.
if (runStateAtLeast(c, SHUTDOWN)
&& (runStateAtLeast(c, STOP) || workQueue.isEmpty())) {
decrementWorkerCount();
return null; // (1)
}
int wc = workerCountOf(c);
// Are workers subject to culling?
boolean timed = allowCoreThreadTimeOut || wc > corePoolSize; // (2)
if ((wc > maximumPoolSize || (timed && timedOut))
&& (wc > 1 || workQueue.isEmpty())) {
if (compareAndDecrementWorkerCount(c))
return null; // (3)
continue;
}
try {
Runnable r = timed ?
workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) :
workQueue.take();
if (r != null)
return r;
timedOut = true;
} catch (InterruptedException retry) {
timedOut = false;
}
}
}
可以看到代码(1)线程池状态为SHUTDOWN并队列为空则返回null,代码(3)线程池的线程数大于核心线程数或者allowCoreThreadTimeOut参数为true并且已经超过线程存活时间的情况下,队列为空则返回null。其中代码(2)用来判断是否需要检查是否超过线程存活时间(allowCoreThreadTimeOut为true或者线程数大于核心线程数情况下需要检查线程的空闲时间是否超过了设置的keepAlive)
processWorkerExit(w, completedAbruptly)方法清理指定的w线程:
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;
if (workerCountOf(c) >= min)
return; // replacement not needed
}
addWorker(null, false);
}
}
可以看到线程终止后主要进行一些清理工作,例如把worker线程从工作集合中删除等
十、关闭线程池
1、shutdown方法
调用shutdown方法后,线程池不再接收新的任务,但是会继续执行任务队列里的任务,该方法立刻返回不会等待队列列的任务执行完成。
public void shutdown() {
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
checkShutdownAccess(); // (1)
advanceRunState(SHUTDOWN); // (2)
interruptIdleWorkers(); // (3)
onShutdown(); // hook for ScheduledThreadPoolExecutor
} finally {
mainLock.unlock();
}
tryTerminate(); // (4)
}
代码(1)检查当前线程是否有权限关闭线程池
代码(2)设置线程池的状态为SHUTDOWN
代码(3)中断线程池中的空闲线程
代码(4)尝试设置线程池的状态为TERMINATE
需要注意的是,在tryTerminate()方法中如果成功设置线程池的状态为TERMINATE,还会调用termination.signalAll()方法唤醒因为调用termination.await()方法挂起的线程。
private void interruptIdleWorkers() {
interruptIdleWorkers(false);
}
private void interruptIdleWorkers(boolean onlyOne) {
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
for (Worker w : workers) {
Thread t = w.thread;
if (!t.isInterrupted() && w.tryLock()) {
try {
t.interrupt();
} catch (SecurityException ignore) {
} finally {
w.unlock();
}
}
if (onlyOne)
break;
}
} finally {
mainLock.unlock();
}
}
w.tryLock()获取Worker自定义的锁,因为Worker在执行任务时首先获取Worker自己的锁,所以这里如果能获取到锁说明这个线程是空闲的,将可以被中断。
2、shutdownNow方法
调用shutdownNow方法后,线程池停止接收新的任务,并且会中断正在执行的线程,丢弃任务队列里的任务。该方法立即返回,返回值是队列里被丢弃的任务
public List<Runnable> shutdownNow() {
List<Runnable> tasks;
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
checkShutdownAccess(); // (1)
advanceRunState(STOP); // (2)
interruptWorkers(); // (3)
tasks = drainQueue(); // (4)
} finally {
mainLock.unlock();
}
tryTerminate(); // (5)
return tasks;
}
代码(1)检查当前线程是否有权限关闭线程池
代码(2)设置线程池状态为STOP
代码(3)中断线程池里的所有线程,包括正在执行的线程
代码(4)取出工作队列里的所有任务到tasks
代码(5)尝试设置线程池的状态为TERMINATE
private void interruptWorkers() {
// assert mainLock.isHeldByCurrentThread();
for (Worker w : workers)
w.interruptIfStarted();
}
该方法会中断所有Worker线程
shutdown方法和shutdownNow方法都需要获取独占锁,确保同一时刻只有一个线程能执行关闭操作。
十一、任务拒绝策略
我们再来看execute方法提交任务时什么情况下会拒绝任务
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); // (1)
else if (workerCountOf(recheck) == 0)
addWorker(null, false);
}
else if (!addWorker(command, false))
reject(command); // (2)
}
首先代码(1)如果线程池状态不是RUNNING则会拒绝接收任务并调用拒绝策略
接下来看看addWorker的代码看看哪些情况下会返回false:
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (int c = ctl.get();;) {
// Check if queue empty only if necessary.
if (runStateAtLeast(c, SHUTDOWN)
&& (runStateAtLeast(c, STOP)
|| firstTask != null
|| workQueue.isEmpty()))
return false; // (1)
for (;;) {
if (workerCountOf(c)
>= ((core ? corePoolSize : maximumPoolSize) & COUNT_MASK))
return false; // (2)
if (compareAndIncrementWorkerCount(c))
break retry;
c = ctl.get(); // Re-read ctl
if (runStateAtLeast(c, SHUTDOWN))
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 c = ctl.get();
if (isRunning(c) ||
(runStateLessThan(c, STOP) && firstTask == null)) {
if (t.getState() != Thread.State.NEW)
throw new IllegalThreadStateException();
workers.add(w);
workerAdded = true;
int s = workers.size();
if (s > largestPoolSize)
largestPoolSize = s;
}
} finally {
mainLock.unlock();
}
if (workerAdded) {
container.start(t);
workerStarted = true;
}
}
} finally {
if (! workerStarted)
addWorkerFailed(w);
}
return workerStarted; // (3)
}
其中代码(1)和线程池状态有关,线程池状态为SHUTDOWN的情况下拒绝任务触发拒绝策略,这个前面的情况一致
代码(2)线程池里面的线程数已经超出maximumPoolSize并且工作队列已满的情况下会返回false出发任务拒绝策略
在代码(3)启动新线程失败会拒绝接收任务并调用拒绝策略
1、AbortPolicy
AbortPolicy.java
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
throw new RejectedExecutionException("Task " + r.toString() +
" rejected from " +
e.toString());
}
AbortPolicy策略抛出RejectedExecutionException异常,这是线程池默认策略
2、CallerRunsPolicy
CallerRunsPolicy.java
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
r.run();
}
}
CallerRunsPolicy策略判断如果线程池状态是RUNNING则在调用者线程中执行任务
3、DiscardOldestPolicy
DiscardOldestPolicy.java
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
if (!e.isShutdown()) {
e.getQueue().poll();
e.execute(r);
}
}
DiscardOldestPolicy策略判断如果线程池状态是RUNNING,则丢弃任务队列里最前面的任务,然后重新提交当前任务
4、DiscardPolicy
DiscardPolicy.java
public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
}
DiscardPolicy策略什么也不做,丢弃当前任务
加载评论中...