ThreadPooleexecutor, как следует из названия, представляет собой класс инструментов управления пулом потоков, который в основном обеспечивает управление задачами, планирование потоков и соответствующие методы ловушек для управления состоянием пула потоков.
1. Описание метода
Основные методы управления задачами следующие:
public void execute(Runnable command);
public <T> Future<T> submit(Callable<T> task);
public <T> Future<T> submit(Runnable task, T result);
public Future<?> submit(Runnable task);
public void shutdown();
public List<Runnable> shutdownNow();
Среди вышеперечисленных методов методы execute() и submit() будут немедленно вызывать поток для выполнения задачи, если есть незанятый поток, разница в том, что метод execute() игнорирует результат выполнения задачи, а метод submit() метод может получить результат. Кроме того, ThreadPoolExecutor также предоставляет методы shutdown() и shutdownNow() для закрытия пула потоков. Разница в том, что метод shutdown() закроет пул потоков после того, как все задачи в очереди задач будут выполнены после вызова, тогда как shutdownNow ( ) напрямую закроет пул потоков и экспортирует задачи из очереди задач в список для возврата.
Помимо вышеперечисленных методов для выполнения задач, ThreadPoolExecutor также предоставляет следующие методы ловушек:
protected void beforeExecute(Thread t, Runnable r);
protected void afterExecute(Runnable r, Throwable t);
protected void terminated();
В ThreadPoolExecutor эти методы по умолчанию пусты, beforeExecute() будет вызываться перед выполнением каждой задачи, afterExecute() будет вызываться после завершения каждой задачи, а метод terminated() будет вызываться при завершении пула потоков. Способ использования этих методов состоит в том, чтобы объявить подкласс для наследования ThreadPoolExecutor и переписать метод ловушки, который необходимо настроить в подклассе, и, наконец, использовать экземпляр подкласса при создании пула потоков.
2. Планирование задач
А. Связанные параметры
Для создания экземпляра ThreadPoolExecutor он в основном имеет следующие важные параметры:
public ThreadPoolExecutor(int corePoolSize, int maximumPoolSize, long keepAliveTime,
TimeUnit unit, BlockingQueue<Runnable> workQueue,
ThreadFactory threadFactory, RejectedExecutionHandler handler);
- corePoolSize: количество основных потоков в пуле потоков;
- maxPoolSize: максимальное количество потоков, которое может создать пул потоков;
- keepAliveTime: когда количество потоков превышает количество потоков, указанное в параметре corePoolSize, и время простоя бездействующего потока достигает времени, указанного текущим параметром, поток будет уничтожен.Если метод allowCoreThreadTimeOut(boolean value) вызывается для позволить основному потоку истечь, то эта стратегия также действительна для основных потоков;
- unit: указывает единицу измерения keepAliveTime, которая может быть миллисекундами, секундами, минутами, часами и т. д.;
- workQueue: очередь, в которой хранятся невыполненные задачи;
- threadFactory: фабрика для создания потоков, если не указано, используется фабрика потоков по умолчанию;
- обработчик: указывает стратегию обработки вновь добавленной задачи, когда очередь задач заполнена и нет доступных потоков для выполнения задачи;
Б. Стратегия планирования
При инициализации пула потоков в пуле нет активных потоков для выполнения задач пользователями.При поступлении новой задачи ее основные задачи выполнения следующие в соответствии с настроенными параметрами:
- Если количество потоков в пуле потоков меньше, чем количество потоков, заданное параметром corePoolSize, каждый раз, когда приходит задача, будет создаваться новый поток для выполнения задачи, независимо от того, есть ли какие-либо простаивающие потоки в пуле потоков;
- Если текущая выполняемая задача достигает заданного числа потоков corePoolSize, что все основные потоки выполняют свои обязанности, на этот раз новая задача будет сохранена в указанной очереди задач workQueue;
- Когда все основные потоки выполняют задачи и очередь задач заполнена задачами, если в это время прибывает новая задача, пул потоков создает новый поток для выполнения задачи;
- Если все потоки (количество потоков, указанное maxPoolSize) выполняют задачи, а очередь задач заполнена задачами, вновь добавленные задачи будут обрабатываться способом, указанным обработчиком.
c. Примечания по стратегии планирования
- На втором этапе все текущие основные потоки выполняют задачи, и когда очередь задач будет заполнена, будет создан новый поток для выполнения задачи.Здесь следует отметить, что общее количество задач, которые необходимо выполнить, когда создание нового потока — это (corePoolSize + workQueueSize), а не только задачи corePoolSize;
- На третьем шаге есть три основных типа workQueue: ArrayBlockingQueue, LinkedBlockingQueue, SynchronousQueue, первый - это очередь с ограниченной блокировкой, второй - очередь с неограниченной блокировкой, конечно, для нее также можно указать размер границы, и третья очередь синхронизации, для ArrayBlockingQueue необходимо указать размер очереди.Когда очередь заполнена пулами потоков задач, будет создан новый поток для выполнения задач.Для LinkedBlockingQueue, если он указывает ограничение, то это мало чем отличается от ArrayBlockingQueue.Если в ней не указано limit , то теоретически она может хранить неограниченное количество задач, на самом деле она может хранить задачи Integer.MAX_VALUE(или эквивалентно неограниченному количеству задач).На данный момент, поскольку LinkedBlockingQueue никогда не может быть заполнена задачами, параметр maxPoolSize будет бессмысленным. Как правило, для него будет установлено то же значение, что и для corePoolSize. Для SynchronousQueue нет внутренней структуры для хранения задач. Когда задача добавляется в Текущий поток и поток последующих задач будут добавлены.Блокировка до тех пор, пока поток не возьмет задачу из очереди, текущий поток будет освобожден, поэтому, если пул потоков использует очередь, corePoolSize обычно предназначен для быть относительно небольшим, а maxPoolSize будет спроектирован так, чтобы быть относительно большим, потому что очередь больше подходит для выполнения большого количества задач с коротким временем выполнения;
- На четвертом шаге DiscardPolicy и DiscardOldestPolicy обычно не используются с SynchronousQueue, потому что, когда очередь синхронизации блокирует задачу, задача будет отброшена; для AbortPolicy, потому что если очередь заполнена, она выдаст исключение, поэтому при использовании Be осторожно; для CallerRunsPolicy, поскольку вызывающий поток будет использоваться для выполнения текущей задачи при поступлении новой задачи, при ее использовании необходимо учитывать влияние на ответ сервера, а также следует отметить, что по сравнению с несколькими другими политиками , эта политика не будет отбрасывать задачи, которые прибывают к задачам, потому что, если прибывающие задачи заполняют очередь и только вызывающий поток может использоваться для выполнения задачи, это означает, что структура пула потоков недостаточно разумна. разрешено развиваться, могут потребоваться все вызывающие потоки.Выполнение задачи блокируется, что приводит к сбою сервера.
3. Объяснение исходного кода
А. Основные атрибуты
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
private static final int COUNT_BITS = Integer.SIZE - 3; // 32
private static final int CAPACITY = (1 << COUNT_BITS) - 1;
// 00011111 11111111 11111111 11111111
private static final int RUNNING = -1 << COUNT_BITS; // 11100000 00000000 00000000 00000000
private static final int SHUTDOWN = 0 << COUNT_BITS; // 00000000 00000000 00000000 00000000
private static final int STOP = 1 << COUNT_BITS; // 00100000 00000000 00000000 00000000
private static final int TIDYING = 2 << COUNT_BITS; // 01000000 00000000 00000000 00000000
private static final int TERMINATED = 3 << COUNT_BITS; // 01100000 00000000 00000000 00000000
Поскольку ThreadPoolExecutor должен управлять несколькими состояниями, а также записывать количество потоков, выполняющих в данный момент задачи, при использовании нескольких переменных управление во время одновременных обновлений будет очень сложным.Здесь ThreadPoolExecutor в основном использует переменную типа AtomicInteger ctl для хранения всей основной информации. . ctl — 32-разрядное целое число, начальное значение равно 0, а старшие три бита используются для хранения информации о состоянии текущего пула потоков, в основном RUNNING, SHUTDOWN, STOP, TIDING и TERMINATED, которые представляют рабочий статус, состояние выключения и состояние завершения, соответственно, состояние завершения и состояние завершения. Конкретная числовая информация, соответствующая этим состояниям, показана в приведенном выше коде.Здесь следует отметить, что в ThreadPoolExecutor числовые значения этих состояний увеличиваются от малого к большому, и поток состояний также последовательно идет вниз. что обеспечивает более удобный способ оценки информации о состоянии.Например, когда необходимо определить, находится ли состояние пула потоков в состоянии SHUTDOWN, необходимо только определить, равно ли значение репрезентативного бита состояния SHUTDOWN. . В ctl, за исключением того, что три самых старших бита используются для представления состояния, значение, представленное оставшимися битами, указывает количество потоков, выполняющих задачи в текущем пуле потоков. Ниже приведены связанные методы для управления атрибутами ctl:
private static int runStateOf(int c) { return c & ~CAPACITY; }
private static int workerCountOf(int c) { return c & CAPACITY; }
private static int ctlOf(int rs, int wc) { return rs | wc; }
private static boolean runStateLessThan(int c, int s) {
return c < s;
}
private static boolean runStateAtLeast(int c, int s) {
return c >= s;
}
private static boolean isRunning(int c) {
return c < SHUTDOWN;
}
private boolean compareAndIncrementWorkerCount(int expect) {
return ctl.compareAndSet(expect, expect + 1);
}
private boolean compareAndDecrementWorkerCount(int expect) {
return ctl.compareAndSet(expect, expect - 1);
}
- runStateOf(int c): используется для получения состояния текущего пула потоков, c — значение атрибута ctl, когда текущий пул потоков работает;
- workerCountOf(int c): используется для получения количества потоков, работающих в текущем пуле потоков, c — значение атрибута ctl, когда работает текущий пул потоков;
- ctlOf(int rs, int wc): здесь rs представляет рабочее состояние текущего потока, а wc представляет количество работающих потоков.Этот метод используется для объединения этих двух параметров в значение атрибута ctl;
- runStateLessThan(int c, int s): Определить, не достигло ли текущее состояние пула потоков заданного состояния.Как упоминалось выше, значение потока состояний последовательно увеличивается, поэтому здесь необходимо судить только о его размере;
- runStateAtLeast(int c, int s): используется для определения того, находится ли текущее состояние пула потоков хотя бы в определенном состоянии;
- isRunning(int c): используется для оценки того, находится ли текущий пул потоков в нормальном рабочем состоянии;
- compareAndIncrementWorkerCount(int expect): увеличить количество рабочих потоков в текущем пуле потоков;
- compareAndDecrementWorkerCount(int expect): уменьшить количество рабочих потоков в текущем пуле потоков.
б. Основной метод
Фактически, для методов execute() и submit() пула потоков базовый метод submit() будет инкапсулировать входящую задачу как объект FutureTask. Поскольку объект FutureTask реализует интерфейс Runnable, его также можно выполнять как task. , где инкапсулированный объект FutureTask передается методу execute() для выполнения. Здесь мы в основном объясняем реализацию метода execute() Ниже приведен код метода execute():
public void execute(Runnable command) {
if (command == null)
throw new NullPointerException();
int c = ctl.get(); // 获取当前线程池状态
if (workerCountOf(c) < corePoolSize) {
// 当工作线程数小于核心线程数时,则调用addWorker()方法创建线程并执行任务
if (addWorker(command, true))
return;
c = ctl.get(); // 若添加失败,则更新当前线程池状态
}
// 执行到此处,则说明线程池中的工作线程要么大于等于核心线程数,要么当前线程池已经被命令关闭了(addWorker方法添加失败的原因),因而这里判断线程池是否为RUNNING状态,是则将任务添加到任务队列中
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); // 线程池至少处于SHUTDOWN状态,拒绝当前任务的执行
}
В методе execute() он сначала определяет, меньше ли количество рабочих потоков в пуле потоков, чем количество основных потоков, и если да, то создает основной поток для выполнения задачи, и если добавление не удается или количество рабочих потоков больше или равно количеству основных потоков, задача добавляется в очередь задач, после успешного добавления будет выполнена двойная проверка, чтобы убедиться, что текущий пул потоков находится в правильном состоянии и что в настоящее время доступны потоки для выполнения вновь добавленных задач. Видно, что для реализации метода execute() основным методом является метод addWorker(), ниже приведена реализация метода addWorker():
private boolean addWorker(Runnable firstTask, boolean core) {
retry:
for (;;) {
int c = ctl.get();
int rs = runStateOf(c); // 获取当前运行状态
// 判断当前线程池是否至少为SHUTDOWN状态,并且firstTask和任务队列中没有任务,是则直接返回
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;
// 走到这一步说明当前线程池处于RUNNING状态,或者任务队列存在任务,并且工作线程数不超过
// 指定的线程数量,那么就增加工作线程数量,成功则继续往下执行,失败则重复上述添加步骤
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());
// 重新检查线程池状态,或者是判断当前是SHUTDOWN状态,而firstTask为空,这说明任务队列此时不为空
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;
}
В методе addWorker() он сначала проверяет, находится ли текущий пул потоков в состоянии RUNNING или в состоянии SHUTDOWN, но в очереди задач все еще есть задачи, затем он создаст новый объект Worker и добавит его в worker Object Collection, а затем вызовите поток, поддерживаемый рабочим объектом, для выполнения задачи.Ниже приведен код реализации рабочего объекта:
private final class Worker extends AbstractQueuedSynchronizer implements Runnable {
private static final long serialVersionUID = 6138294804551838833L;
final Thread thread; // 当前工作者中执行任务的线程
Runnable firstTask; // 第一个需要执行的任务
volatile long completedTasks; // 当前工作者完成的任务数
Worker(Runnable firstTask) {
// 默认设置为-1,那么如果不调用当前工作者的run()方法,那么其状态是不会改变的,
// 其他的线程也无法使用当前工作者执行任务,在run()方法调用的runWorker()方法中会
// 调用unlock()方法使当前工作者处于正常状态
setState(-1);
this.firstTask = firstTask;
this.thread = getThreadFactory().newThread(this); // 使用线程工厂创建线程
}
public void run() {
runWorker(this); // 使用当前工作者执行任务
}
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(); }
// 如果当前线程已经在执行任务,那么将其标记为打断状态,待其任务执行完毕则终止任务的执行
void interruptIfStarted() {
Thread t;
if (getState() >= 0 && (t = thread) != null && !t.isInterrupted()) {
try {
t.interrupt();
} catch (SecurityException ignore) {
}
}
}
}
В рабочем объекте он в основном поддерживает рабочий поток для выполнения задач. Рабочий объект наследует AbstractQueuedSynchronizer для управления получением текущего рабочего статуса рабочего процесса, а также реализует интерфейс Runnable для инкапсуляции выполнения основной задачи в метод run(). Ниже приведена конкретная реализация метода runWorker():
final void runWorker(Worker w) {
Thread wt = Thread.currentThread();
Runnable task = w.firstTask;
w.firstTask = null;
w.unlock(); // 重置Worker对象的状态
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);
}
}
Видно, что в методе runWorker() он сначала выполнит задачу инициализации рабочего объекта, а после завершения выполнения будет непрерывно получать выполнение задачи из очереди задач через бесконечный цикл. Ниже приведен исходный код метода getTask():
private Runnable getTask() {
boolean timedOut = false;
for (;;) {
int c = ctl.get();
int rs = runStateOf(c);
// 判断当前线程是否处于STOP状态,或者处于SHUTDOWN状态,并且工作队列是空的,是则不返回任务
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 {
// 允许核心线程过期或者工作线程数大于corePoolSize时,从任务队列获取任务时会指定等待时间,
// 否则会一直等待任务队列中新的任务
Runnable r = timed ? workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) : workQueue.take();
if (r != null)
return r;
timedOut = true;
} catch (InterruptedException retry) {
timedOut = false;
}
}
}
Видно, что метод GetTask сначала определит, является ли текущее состояние пула потоков состоянием STOP или состоянием выключения, а очередь задач пуста, не возвращена ли задача, иначе задача получена из очереди задач по соответствующему параметру.
Основные этапы реализации описанного выше метода execute(), еще один важный метод в ThreadPoolExecutor — это метод shutdown(). Ниже приведен основной код метода shutdown():
public void shutdown() {
final ReentrantLock mainLock = this.mainLock;
mainLock.lock();
try {
checkShutdownAccess(); // 检查对线程状态的控制权限
advanceRunState(SHUTDOWN); // 更新当前线程池状态为SHUTDOWN
interruptIdleWorkers(); // 打断空闲的工作者
onShutdown(); // 钩子方法,但是没有对外公开,因为该方法只有包访问权限
} finally {
mainLock.unlock();
}
tryTerminate();
}
В методе shutdown() он сначала проверяет, имеет ли текущий поток разрешение на изменение состояния потока, затем изменяет состояние текущего пула потоков на SHUTDOWN, затем вызывает метод interruptIdleWorkers() для прерывания всех бездействующих потоков и, наконец, вызывает tryTerminate. Метод() пытается изменить состояние текущего пула потоков с SHUTDOWN на TERMINATED, где метод interruptIdleWorkers() в итоге вызовет свой перегруженный метод interruptIdleWorkers(boolean), код которого выглядит следующим образом:
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();
}
}
Как видите, этот метод обходит все рабочие объекты и завершает их работу, если они находятся в состоянии ожидания. Для потока в рабочем состоянии, поскольку текущее состояние пула потоков было установлено как SHUTDOWN в методе shutdown(), поток в рабочем состоянии автоматически уничтожит задачи в очереди задач после выполнения всех задач.
В этой статье в основном объясняется основной метод ThreadPoolExecutor, метод планирования пула потоков и принцип реализации его основных функций.Если в этой статье есть какие-либо неуместности, пожалуйста, поправьте меня, спасибо!