Поклонник параллелизма @ThreadPoolExecutor All-in-one

задняя часть

Поклонник параллелизма @ThreadPoolExecutor All-in-one

JAVA 并发 1.8版


Уведомление:Автор откроет.并发番@Future一文通Представлено далее вс возвращаемым значениемОперация пула потоков будет подробно описана позже.AbstractExecutorService, так что эта статья не потребует слишком много знаний;并发番@Future一文通, автор также дополнительно представит多样化的线程池配置使用а такжеTomcat的线程池配置, Следите за обновлениями

1. Вступительная глава

1.1 Как создать больше тем

В Java вы можете сделать это, настроив-XssПараметры для настройки размера стека каждого потока (64-битная система по умолчанию 1024 КБ), когда уменьшение этого значения означает, что можно создать больше потоков, но проблема заключается в следующем.JVM资源是有限的,线程不能无限创建!

ты можешь пройтиПул потоков контролирует количество потоков, пул потоков аналогичен пулу соединений, который может бытьПовторное использование ограниченного числа потоков эффективно снижает накладные расходы на частое создание и уничтожение потоков.; в то же время пул потоков прекрасно с этим справляется生产者-消费者模式, отправка задач эквивалентна производству, а выполнение задач эквивалентно потреблению.

1.2 Обзор пулов потоков

Исходя из опыта разработки автора, пул потоков должен быть параллельной структурой с наибольшей частотой и сценариями приложений в параллельных пакетах.Почти все задачи параллельного/асинхронного выполнения требуют использования пулов потоков, повторного использования потоков в виде внутренних пулов потоков. задачи внешнего управления, планирование потоков, управление пулом потоков и другие службы. Разумное использование пулов потоков может принести следующие три преимущества:

1. Сокращение потребления ресурсов:Уменьшите стоимость создания и уничтожения потоков за счет повторного использования уже созданных потоков.

2. Улучшить скорость отклика:При поступлении задачи ее можно выполнить немедленно, не дожидаясь создания потока.

3. Улучшить управляемость потоков:Пулы потоков можно единообразно управлять, распределять, настраивать и отслеживать.

2. Состав пула потоков

2.1 Определение класса

public class ThreadPoolExecutor extends AbstractExecutorService

2.2 Конструкторы

/**
 * 线程工厂默认为DefaultThreadFactory
 * 饱和策略默认为AbortPolicy
 */
public ThreadPoolExecutor(int corePoolSize,
                          int maximumPoolSize,
                          long keepAliveTime,
                          TimeUnit unit,
                          BlockingQueue<Runnable> workQueue) {
    this(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue,
        Executors.defaultThreadFactory(), defaultHandler);
}

/**
 * 线程工厂可配置
 * 饱和策略默认为AbortPolicy
 */    
public ThreadPoolExecutor(int corePoolSize,
                          int maximumPoolSize,
                          long keepAliveTime,
                          TimeUnit unit,
                          BlockingQueue<Runnable> workQueue,
                          ThreadFactory threadFactory) {
    this(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue,
        threadFactory, defaultHandler);
}   

/**
 * 线程工厂默认为DefaultThreadFactory
 * 饱和策略可配置
 */
public ThreadPoolExecutor(int corePoolSize,
                          int maximumPoolSize,
                          long keepAliveTime,
                          TimeUnit unit,
                          BlockingQueue<Runnable> workQueue,
                          RejectedExecutionHandler handler) {
    this(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue,
        Executors.defaultThreadFactory(), handler);
}

/**
 * 线程工厂可配置
 * 饱和策略可配置
 */
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.acc = System.getSecurityManager() == null ?
                null :  AccessController.getContext();
    this.corePoolSize = corePoolSize;
    this.maximumPoolSize = maximumPoolSize;
    this.workQueue = workQueue;
    this.keepAliveTime = unit.toNanos(keepAliveTime);
    this.threadFactory = threadFactory;
    this.handler = handler;
}

2.3 Важные переменные

//线程池控制器
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));
//任务队列
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;
/**
 * 允许核心工作线程响应超时回收
 * false:核心工作线程即使空闲超时依旧存活
 * true:核心工作线程一旦超过keepAliveTime仍然空闲就被回收
 */
private volatile boolean allowCoreThreadTimeOut;
//核心工作线程数
private volatile int corePoolSize;
//最大工作线程数
private volatile int maximumPoolSize;
//默认饱和策略执行器 - AbortPolicy -> 直接抛出异常
private static final RejectedExecutionHandler defaultHandler =
        new AbortPolicy();

3. Использование пула потоков

3.1 Создание пула потоков

Создание пула потоков фактически создает экземпляр объекта пула потоков.Здесь мы используем наиболее полный конструктор для описания наиболее полного процесса создания:

1.corePoolSize (количество основных рабочих потоков):Когда задач нет, пул потоков позволяет (поддерживает) минимальное количество простаивающих пулов потоков; когда задача отправляется в пул потоков, создается новый рабочий поток для выполнения задачи (даже если есть простаивающие основные рабочие потоки). в это время) пока(实际工作线程数 >= 核心工作线程数)пока; звонитьprestartAllCoreThreads()метод создаст и запустит все основные рабочие потоки заранее

2.workQueue (очередь задач):Блокирующая очередь, используемая для хранения задач, ожидающих выполнения; когда(实际工作线程数 >= 核心工作线程数) && (任务数 < 任务队列长度), задача будетoffer()Ожидание постановки в очередь; подробности об очередях задач см. ниже.任务队列与排队策略

3.maximumPoolSize (максимальное количество рабочих потоков):Максимальное количество рабочих потоков, которое может создавать пул потоков; когда(队列已满 && 实际工作线程数 < 最大工作线程数), пул потоков будет создавать новые рабочие потоки (даже если все еще есть незанятые рабочие потоки) для выполнения задач до максимального количества рабочих потоков; этот параметр фактически недействителен при установке неограниченной очереди

4.keepAliveTime (максимальное время простоя рабочих потоков):Через наносекунды бездействующие рабочие потоки, соответствующие условию тайм-аута, будут перезапущены; неосновные рабочие потоки, у которых истекло время ожидания, будут перезапущены, но основные рабочие потоки не будут перезапущены;allowCoreThreadTimeOut=trueЕсли значение не установлено, поток будет жить вечно; рекомендуется увеличить время для улучшения использования потока, когда сцена короткая и много задач

5.unit (единица времени удержания активности потока):Единица времени удержания активности потока, необязательно включаетNANOSECONDS纳秒,MICROSECONDS微秒,MILLISECONDS毫秒,SECONDS秒,MINUTES分,HOURS时,DAYS天

6.threadFactory (фабрика создания потоков):Как следует из названия, это фабрика для создания потоков, позволяющая настраивать фабрики и начальную настройку потоков, например имена, потоки демонов, обработку исключений и т. д.

7.handler (исполнитель политики насыщения):Когда пул потоков и очередь заполнены, это означает, что поток не может получить больше задач, то есть количество задач насыщено, и принимать заказы невозможно, в это время необходимо использовать стратегию насыщения. разобраться с этим.新提交的任务, по умолчаниюAbort(直抛Reject异常),Также включаетDiscard(LIFO规则丢弃),DiscardOldest(LRU规则丢弃)а такжеCallerRuns(调用者线程执行), что позволяет использовать настраиваемые приводы

Дополнение: причина, по которой при сравнении количества потоков не отображается только "=", заключается в том, что пул потоков допускает динамическое управление, подробности см. ниже.

public ThreadPoolExecutor(int corePoolSize,
                          int maximumPoolSize,
                          long keepAliveTime,
                          TimeUnit unit,
                          BlockingQueue<Runnable> workQueue,
                          ThreadFactory threadFactory,
                          RejectedExecutionHandler handler) {

    //注意数值条件,否则在初始化时会直接抛出IAE                 
    if (corePoolSize < 0 ||
        maximumPoolSize <= 0 ||
        maximumPoolSize < corePoolSize ||
        keepAliveTime < 0)
        throw new IllegalArgumentException();

    //任务队列、线程工厂、饱和策略执行器都不允许为空,否则在初始化是直接排除NPE    
    if (workQueue == null || threadFactory == null || handler == null)
        throw new NullPointerException();
    this.acc = System.getSecurityManager() == null ?
                null :  AccessController.getContext();
    this.corePoolSize = corePoolSize;
    this.maximumPoolSize = maximumPoolSize;
    this.workQueue = workQueue;
    this.keepAliveTime = unit.toNanos(keepAliveTime);
    this.threadFactory = threadFactory;
    this.handler = handler;
}

3.2 Отправка и выполнение задач

Вы можете выбрать один из двух в зависимости от того, хотите ли вы возвращаемое значение:
1. execute(): подходит для отправки задач, которым не нужно возвращать значение.

- Этот метод не может определить, успешно ли выполняется задача пулом потоков.

2. submit(): подходит для отправки задач, требующих возвращаемого значения.

- можно вернуть черезFutureОбъект знает, успешно ли выполнена задача

-get()метод заблокирует текущий поток, пока задача не будет завершена,Но будьте осторожны, чтобы предотвратить бесконечную блокировку! ! !

-использоватьget(long timeout,TimeUnit unit)Метод будет блокировать текущий поток до тех пор, пока задача не завершится или не истечет время ожидания, бесконечная блокировка не происходит, но требуетОбратите внимание, что задание может быть не выполнено по истечении тайм-аута! ! !

3.3 Закройте пул потоков

Вы можете выбрать один из следующих двух способов закрытия пула потоков, автор рекомендует использовать более безопасный первый:

1. shutdown():Остановите пул потоков в установленном порядке, отправленные задачи будут выполняться (в том числе находящиеся в исполнении и в очереди задач), но новые задачи будут отклонены

2. shutdownNow():Немедленно (попытаться) прекратить выполнение всех задач (как выполняемых, так и находящихся в очереди задач) и вернуться к списку ожидающих задач

Примечание. Доступ ко всем вышеперечисленным методам можно получить, вызвавawaitTermination()Дождитесь завершения задачи, прежде чем завершать работу пула потоков.

3.4 Разумно настроить пул потоков

Рекомендуем вам прочитатьРазумно настроить пул потоков, если будет возможность, автор в дальнейшем поделится реальным боевым опытом

Размер пула потоков рекомендуется определять по результатам стресс-тестов конкретного бизнеса или оценивать по закону Литтла.

Закон Литтла, английское название: Закон Литтла (результат Литтла, теорема, лемма или формула), в устойчивой системе среднее количество клиентов L, наблюдаемое в течение длительного времени, равно эффективной скорости прибытия λ, наблюдаемой в течение длительного времени и произведение среднего времени, которое каждый клиент проводит в системе, т. е. L = λW. (из Байду)

Конечно, реальность всегда жестока, например, сетевые исключения, тайм-ауты ввода-вывода и т. д. могут вызвать нестабильность во времени, в то же время могут возникнуть различные странные проблемы во время использования пулов потоков, так что это также необходимо учитывать обработку тайм-аута задачи и обработку сбоев, обработку изоляции и т. д., особенно при возникновении большого количества ошибок, таких как аномальные торренты, торренты с высоким параллелизмом и т. д.

Рекомендуем вам прочитатьExecutorService - 10 советов и хитростей

Уведомление:Автор откроет.并发番@Future一文通Представлено далее вс возвращаемым значениемОперация пула потоков будет подробно описана позже.AbstractExecutorService, так что эта статья не потребует слишком много знаний;并发番@Future一文通, автор также дополнительно представит多样化的线程池配置使用а такжеTomcat的线程池配置, Следите за обновлениями

4. Принцип реализации пула потоков

4.1 Блок-схема

流程图

4.2 Реализация

По мере увеличения количества задач выполнение пула потоков в основном включает следующие четыре ситуации, которые соответствуют суждениям в потоке обработки:

1. Если фактическое количество рабочих потоковworkerscorePoolSize, создается новый рабочий поток для выполнения новой задачиexecute(Runable)

2. Если фактическое количество рабочих потоковworkers>= количество основных рабочих потоковcorePoolSize(все основные рабочие потоки выполняют задачи) и очередь задачworkQueueЕсли он не заполнен, добавьте задачу в очередь задач.workQueueсередина

3. Если очередь задачworkQueueзаполнен, создайте новый рабочий поток для выполнения задачиexecute()

4. Если фактическое количество рабочих потоковworkers>= максимальное количество рабочих потоковmaximumPoolSize(Все потоки выполняют задачи), количество задач в это время насыщено, и стратегия отклонения должна основываться на насыщении.rejectedExecutionHandlerВыполните соответствующую операцию отклонения насыщения

Общий дизайн пула потоков основан на соображениях производительности, максимально избегая получения глобальных блокировок:

1. Поскольку глобальная блокировка должна быть получена при создании нового потока, поэтому步骤1и步骤3должен быть заблокирован

2. Во избежание многократного получения глобальной блокировки (узкое место масштабирования производительности), когда实际工作线程数>=核心工作线程数, затем выполняет步骤2(Нет необходимости приобретать глобальную блокировку при постановке в очередь)

Примечание: Не будьтеrejectЗапутался, это просто означает, что в пуле потоков нет дополнительных рабочих потоков для выполнения и дополнительного места в очереди для хранения задачи, это не означает, что задача действительно не обрабатывается, способ обработки задачи зависит от политики отклонения насыщения

4.3 Обработка тайм-аута

Если вам нужно обработать основные рабочие потоки с истекшим временем ожидания, выберите второй; если нет, выберите первый:

1. Если фактическое количество рабочих потоковworkers> количество основных рабочих потоковcorePoolSize, время простоя рециркуляции превышаетkeepAliveTimeбездействующих неосновных потоков (уменьшите количество рабочих потоков до

2. Если установленоallowCoreThreadTimeOutКогда правда, болееkeepAliveTimeБездействующие основные рабочие потоки также перерабатываются.

4.4 Статус пула потоков

4.4.1 Контроллер состояний

//线程池状态控制器,用于保证线程池状态和工作线程数 ps:低29位为工作线程数量,高3位为线程池状态
private final AtomicInteger ctl = new AtomicInteger(ctlOf(RUNNING, 0));

//设定偏移量 Integer.SIZE = 32 -> 即COUNT_BITS = 29
private static final int COUNT_BITS = Integer.SIZE - 3;

//确定最大的容量2^29-1
private static final int CAPACITY   = (1 << COUNT_BITS) - 1;

//获取线程池状态,取高3位
private static int runStateOf(int c)     { return c & ~CAPACITY; }

//获取工作线程数量,取低29位
private static int workerCountOf(int c)  { return c & CAPACITY; }

/** 
 * 获取线程池状态控制器
 * @param rs 表示runState 线程池状态
 * @param wc 表示workerCount 工作线程数量
 */
private static int ctlOf(int rs, int wc) { return rs | wc; }
Вот немного базовых знаний о бинарных операторах для читателей, которые забыли понять:

&: оператор И, если тот же бит равен 1, он равен 1, иначе он равен 0

|: оператор ИЛИ, если один из одинаковых битов равен 1, это 1, иначе это 0

~: Не оператор, 0 и 1 меняются местами, то есть, если 0 становится 1, 1 становится 0

^: оператор XOR, 0, если один и тот же бит совпадает, 1, если он отличается

4.4.2 Статус пула потоков

Поток состояний потока следует в следующем порядке, от малого к большему:

RUNNING -> SHUTDOWN -> STOP -> TIDYING -> TERMINATED

Дополнение: Смена ценностей ощущается как наш возраст, чем мы старше, тем ближе мы к Богу

// runState is stored in the high-order bits 用Integer的高三位表示
//高3位111,低29位为0 该状态下线程池会接收新提交任务和执行队列任务
private static final int RUNNING    = -1 << COUNT_BITS;

//高3位000,低29位为0 该状态下线程池不再接收新任务,但还会继续执行队列任务
private static final int SHUTDOWN   =  0 << COUNT_BITS;

//高3位001,低29位为0 该状态下线程池不再接收新任务,不会再执行队列任务,并会中断正在执行中的任务
private static final int STOP       =  1 << COUNT_BITS;

//高3位010,低29位为0 该状态下线程池的所有任务都被终止,工作线程数为0,期间会调用钩子方法terminated()
private static final int TIDYING    =  2 << COUNT_BITS;

//高3位011,低29位为0 该状态下表明线程池terminated()方法已经调用完成
private static final int TERMINATED =  3 << COUNT_BITS;

4.5 Worker

Worker — это внутренний класс пула потоков, который используется для инкапсуляции рабочих потоков и задач и управления статусом прерывания рабочих потоков и других функций.

4.5.1 Состав

Класс Worker инкапсулирует эти три части (блокировка + поток + задача), таким образом становясь универсальным:

1. Наследовать класс AQS:Реализуйте простой неповторный мьютекс, чтобы обеспечить удобные операции блокировки для обработки условий прерывания.

2. Реализуйте интерфейс Runnable:«Конъюнктурный» дизайн, в основном заимствованныйRunnable接口Унифицированный метод написания, преимущество в том, что нет необходимости переписывать интерфейс с той же функцией

3. Рабочий поток:Рабочий пройдетthread变量Привязать рабочий поток (один к одному), который фактически выполняет задачу, которая выделяется фабрикой потоков во время инициализации, и он будет многократно получать и выполнять задачи.

4. Задача:Работник назначает новые задачиfirstTask变量, рабочий поток каждый раз обрабатывает вновь полученную задачу через эту переменную (значение может быть нулевым при инициализации, что имеет особый эффект, который будет подробно описан ниже)

Дополнение: поскольку между рабочими потоками и рабочими экземплярами существует отношение один к одному, легко понять, что рабочие потоки эквивалентны рабочим, особенно когда речь идет о количестве, например, создание рабочего потока фактически создает работник, требующий вашего понимания, для удобства описания автор скроет эту связь в тексте

private final class Worker
        extends AbstractQueuedSynchronizer
        implements Runnable{

    /** 实际上真正的工作线程 - 幕后大佬,但可能因线程工厂创建失败而为null */
    final Thread thread;

    /** 待执行任务,可能为null */
    Runnable firstTask;

    /** 该工作线程已完成的任务数 -- 论KPI的重要性 */
    volatile long completedTasks;        

    Worker(Runnable firstTask) {

        //设置锁状态为-1,目的是为了阻止在runWorker()之前被中断
        setState(-1); 

        /**
         * 新任务,任务来源有两个:
         *   1.调用addWorker()方法新建线程时传入的第一个任务
         *   2.调用runWorker()方法时内部循环调用getTask() -- 这就是线程复用的具现
         */
        this.firstTask = firstTask;

        /**
         * 创建一个新的线程 -> 这个是真正的工作线程
         * 注意Worker本身就是个Runnable对象
         * 因此newThread(this)中的this也是个Runnable对象
         */
        this.thread = getThreadFactory().newThread(this);
    }
}

4.5.2 Выполнение задач

@FunctionalInterface
public interface Runnable {
    public abstract void run();
}
/**
 * 工作线程运行
 * runWorker方法内部会通过轮询的方式
 * 不停地获取任务和执行任务直到线程被回收
 */
public void run() {
    runWorker(this);
}
(выделено) Вот краткое введение в рабочий процесс потоков, выполняющих задачи в пуле потоков:

1. Прежде чем рабочий поток начнет выполняться, рабочий поток должен быть сначала заблокирован, а задача разблокирована после завершения задачи.

2. Выполнять до и после выполнения задачи отдельноbeforeExecute()иafterExecute()метод

3. Когда во время выполнения возникает исключение, оно будет выброшено.Умирает ли поток, зависит от вашей обработки исключения

4. После выполнения каждой задачи количество задач, выполненных текущим рабочим потоком, будет автоматически увеличиваться, при этом будет неоднократно вызываться getTask() для получения задач из очереди задач и их выполнения. выполняется, поток будет заблокирован в этом методе.

5. Когда рабочий поток завершается по разным причинам, он выполняетсяprocessWorkerExit()Перезапустить поток (ядро состоит в том, чтобы удалить воркер из коллекции воркеров, обратите внимание, что воркер ранее вышел из цикла задач, поэтому он больше не работает, и это удобно для gc после удаления из коллекции)

4.5.3 Метод блокировки

// Lock methods
// The value 0 represents the unlocked state. 0表示未锁定
// The value 1 represents the locked state. 1表示已锁定

protected boolean isHeldExclusively() {
    return getState() != 0;
}

protected boolean tryAcquire(int unused) {

    //锁状态非0即1,即不可重入
    //特殊情况:只有初始化时才为-1,目的是防止线程初始化阶段被中断
    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) {
        }
    }
}

Небольшой вопрос: почему бы не выполнить отправленную команду напрямую, но нужно использовать пакет Worker?
友情小提示:这跟worker的作用有关系

Небольшой ответ: в основном для контроля прерываний


Небольшой вопрос: как управлять прерываниями?
友情小提示:Worker继承了AQS从而是一把AQS锁

Небольшой ответ: у Worker есть следующие четыре критерия для обработки прерываний:

1. Когда рабочий поток фактически начинает выполняться, его нельзя прерывать.

2. Когда рабочий поток выполняет задачу, его нельзя прерывать.

3. Когда рабочий поток ожидает получения задачи из очереди задачgetTask()быть прерванным

4. позвонитьinterruptIdleWorkers()Рабочая блокировка должна быть получена первой при прерывании бездействующего потока.


Небольшой вопрос: почему рабочие процессы не предназначены для повторных блокировок?
友情小提示:禁止在动态控制时再次获取锁

Небольшой ответ: поскольку поток может быть прерван в методе динамического управления, таком как вызовinterruptIdleWorkers(), чтобы метод выполнялсяinterrupt()позвонит раньшеworker.tryLock(), если в это время разрешен повторный вход, это приведет к неожиданному прерыванию потока, что аналогично当工作线程正在执行任务时,不允许被中断норма противоречит


4.6 Динамическое управление

Пул потоков предоставляет несколько общедоступных методов для динамического управления информацией о конфигурации пула потоков:

/**
 * 设置核心工作线程数
 * 1.若新值<当前值时,将调用interruptIdleWorkers()处理超出部分线程
 * 2.若新值>当前值时,新创建的线程(若有必要)直接会处理队列中的任务
 */
public void setCorePoolSize(int corePoolSize)

/**
 * 设置是否响应核心工作线程超时处理
 * 1.设置false时,核心工作线程不会因为任务数不足(空闲)而被终止
 * 2.设置true时,核心工作线程和非核心工作线程待遇一样,会因为超时而终止
 * 注意:为了禁止出现持续性的线程替换,当设置true时,超时时间必须>0
 * 注意:该方法通常应在线程池被使用之前调用
 */
public void allowCoreThreadTimeOut(boolean value)

/**
 * 设置最大工作线程数
 * 1.若新值<当前值时,将调用interruptIdleWorkers()处理超出部分线程
 * 注意:当新值>当前值时是无需做任何处理的,跟设置核心工作线程数不一样
 */
public void setMaximumPoolSize(int maximumPoolSize)

/** 
 * 设置超时时间,超时后工作线程将被终止
 * 注意:若实际工作线程数只剩一个,除非线程池被终止,否则无须响应超时
 */
public void setKeepAliveTime(long time, TimeUnit unit) 

5. Подача и выполнение задания

5.1 execute() — отправить задачу

/**
 * 在未来的某个时刻执行给定的任务
 * 这个任务由一个新线程执行,或者用一个线程池中已经存在的线程执行
 * 如果任务无法被提交执行,要么是因为这个Executor已经被shutdown关闭
 * 要么是已经达到其容量上限,任务会被当前的RejectedExecutionHandler处理
 */
public void execute(Runnable command) {

        //新任务不允许为空,空则抛出NPE
        if (command == null)
            throw new NullPointerException();

        /**
         * 1.若实际工作线程数 < 核心工作线程数,会尝试创建一个工作线程去执行该
         * 任务,即该command会作为该线程的第一个任务,即第一个firstTask
         * 
         * 2.若任务入队成功,仍需要执行双重校验,原因有两点:
         *      - 第一个是去确认是否需要新建一个工作线程,因为可能存在
         *        在上次检查后已经死亡died的工作线程
         *      - 第二个是可能在进入该方法后线程池被关闭了,
         *        比如执行shutdown()
         *   因此需要再次检查state状态,并分别处理以上两种情况:
         *      - 若线程池中已无可用工作线程了,则需要新建一个工作线程 
         *      - 若线程池已被关闭,则需要回滚入队列(若有必要)
         * 
         * 3.若任务入队失败(比如队列已满),则需要新建一个工作线程; 
         *   若新建线程失败,说明线程池已停止或者已饱和,必须执行拒绝策略
         */
        int c = ctl.get();

        /**
         * 情况一:当实际工作线程数 < 核心工作线程数时
         * 执行方案:会创建一个新的工作线程去执行该任务
         * 注意:此时即使有其他空闲的工作线程也还是会新增工作线程,
         *      直到达到核心工作线程数为止
         */
        if (workerCountOf(c) < corePoolSize) {

            /**
             * 新增工作线程,true表示要对比的是核心工作线程数
             * 一旦新增成功就开始执行当前任务
             * 期间也会通过自旋获取队列任务进行执行
             */
            if (addWorker(command, true))
                return;

            /**
             * 需要重新获取控制器状态,说明新增线程失败
             * 线程失败的原因可能有两种:
             *  - 1.线程池已被关闭,非RUNNING状态的线程池是不允许接收新任务的
             *  - 2.并发时,假如都通过了workerCountOf(c) < corePoolSize校验,但其他线程
             *      可能会在addWorker先创建出线程,导致workerCountOf(c) >= corePoolSize,
             *      即实际工作线程数 >= 核心工作线程数,此时需要进入情况二
             */
            c = ctl.get();
        }

        /**
         * 情况二:当实际工作线程数>=核心线程数时,新提交任务需要入队
         * 执行方案:一旦入队成功,仍需要处理线程池状态突变和工作线程死亡的情况
         */
        if (isRunning(c) && workQueue.offer(command)) {

            //双重校验
            int recheck = ctl.get();

            /**
             * recheck的目的是为了防止线程池状态的突变 - 即被关闭
             * 一旦线程池非RUNNING状态时,除了从队列中移除该任务(回滚)外
             * 还需要执行任务拒绝策略处理新提交的任务
             */
            if (!isRunning(recheck) && remove(command))

                //执行任务拒绝策略
                reject(command);

            /**
             * 若线程池还是RUNNING状态 或 队列移除失败(可能正好被一个工作线程拿到处理了)
             * 此时需要确保至少有一个工作线程还可以干活
             * 补充一句:之所有无须与核心工作线程数或最大线程数相比,而只是比较0的原因是
             *          只要保证有一个工作线程可以干活就行,它会自动去获取任务
             */
            else if (workerCountOf(recheck) == 0)

                /**
                 * 若工作线程都已死亡,需要新增一个工作线程去干活
                 * 死亡原因可能是线程超时或者异常等等复杂情况
                 *
                 * 第一个参数为null指的是传入一个空任务,
                 * 目的是创建一个新工作线程去处理队列中的剩余任务
                 * 第二个参数为false目的是提示可以扩容到最大工作线程数
                 */
                addWorker(null, false);
        }

        /**
         * 情况三:一旦线程池被关闭 或者 新任务入队失败(队列已满)
         * 执行方案:会尝试创建一个新的工作线程,并允许扩容到最大工作线程数
         * 注意:一旦创建失败,比如超过最大工作线程数,需要执行任务拒绝策略
         */
        else if (!addWorker(command, false))

            //执行任务拒绝策略
            reject(command);
    }

5.2 addWorker() — добавить рабочий поток

/**
 * 新增工作线程需要遵守线程池控制状态规定和边界限制
 *
 * @param core core为true时允许扩容到核心工作线程数,否则为最大工作线程数
 * @return 新增成功返回true,失败返回false
 */
private boolean addWorker(Runnable firstTask, boolean core) {

    //重试标签
    retry:

    /***
     * 外部自旋 -> 目的是确认是否能够新增工作线程
     * 允许新增线程的条件有两个:
     *   1.满足线程池状态条件 -> 条件一
     *   2.实际工作线程满足数量边界条件 -> 条件二
     * 不满足条件时会直接返回false,表示新增工作线程失败
     */
    for (;;) {

        //读取原子控制量 - 包含workerCount(实际工作线程数)和runState(线程池状态)
        int c = ctl.get();

        //读取线程池状态
        int rs = runStateOf(c);

        /**
         * 条件一.判断是否满足线程池状态条件
         *  1.只有两种情况允许新增线程:
         *    1.1 线程池状态==RUNNING
         *    1.2 线程池状态==SHUTDOWN且firstTask为null同时队列非空
         *
         *  2.线程池状态>=SHUTDOWN时不允许接收新任务,具体如下:
         *    2.1 线程池状态>SHUTDOWN,即为STOP、TIDYING、TERMINATED
         *    2.2 线程池状态==SHUTDOWN,但firstTask非空
         *    2.3 线程池状态==SHUTDOWN且firstTask为空,但队列为空
         *  补充:针对1.2、2.2、2.3的情况具体请参加后面的"小问答"环节
         */
        if (rs >= SHUTDOWN &&
            !(rs == SHUTDOWN && firstTask == null && ! workQueue.isEmpty()))
            return false;

        /***
         * 内部自旋 -> 条件二.判断实际工作线程数是否满足数量边界条件
         *   -数量边界条件满足会对尝试workerCount实现CAS自增,否则新增失败
         *   -当CAS失败时会再次重新判断是否满足新增条件:
         *       1.若此期间线程池状态突变(被关闭),重新判断线程池状态条件和数量边界条件
        *        2.若此期间线程池状态一致,则只需重新判断数量边界条件
        */
        for (;;) {

            //读取实际工作线程数
            int wc = workerCountOf(c);

            /**
             * 新增工作线程会因两种实际工作线程数超标情况而失败:
             *  1.实际工作线程数 >= 最大容量
             *  2.实际工作线程数 > 工作线程比较边界数(当前最大扩容数)
             *   -若core = true,比较边界数 = 核心工作线程数
             *   -若core = false,比较边界数 = 最大工作线程数
             */
            if (wc >= CAPACITY || wc >= (core ? corePoolSize : maximumPoolSize))
                return false;

            /**
             * 实际工作线程计数CAS自增:
             *   1.一旦成功直接退出整个retry循环,表明新增条件都满足
             *   2.因并发竞争导致CAS更新失败的原因有三种: 
             *      2.1 线程池刚好已新增一个工作线程
             *        -> 计数增加,只需重新判断数量边界条件
             *      2.2 刚好其他工作线程运行期发生错误或因超时被回收
             *        -> 计数减少,只需重新判断数量边界条件
             *      2.3 刚好线程池被关闭 
             *        -> 计数减少,工作线程被回收,
             *           需重新判断线程池状态条件和数量边界条件
             */
            if (compareAndIncrementWorkerCount(c))
                break retry;

            //重新读取原子控制量 -> 原因是在此期间可能线程池被关闭了
            c = ctl.get();

            /**
             * 快速检测是否发生线程池状态突变
             *  1.若状态突变,重新判断线程池状态条件和数量边界条件
             *  2.若状态一致,则只需重新判断数量边界条件
             */
            if (runStateOf(c) != rs)
                continue retry;
        }
    }

    /**
     * 这里是addWorker方法的一个分割线
     * 前面的代码的作用是决定了线程池接受还是拒绝新增工作线程
     * 后面的代码的作用是真正开始新增工作线程并封装成Worker接着执行后续操作
     * PS:虽然笔者觉得这个方法其实可以拆分成两个方法的(在break retry的位置)
     */

    //记录新增的工作线程是否开始工作
    boolean workerStarted = false;

    //记录新增的worker是否成功添加到workers集合中
    boolean workerAdded = false;

    Worker w = null;
    try {

        //将新提交的任务和当前线程封装成一个Worker
        w = new Worker(firstTask);

        //获取新创建的实际工作线程
        final Thread t = w.thread;

        /**
         * 检测是否有可执行任务的线程,即是否成功创建了新的工作线程
         *   1.若存在,则选择执行任务
         *   2.若不存在,则需要执行addWorkerFailed()方法
         */
        if (t != null) {

            /**
             * 新增工作线程需要加全局锁
             * 目的是为了确保安全更新workers集合和largestPoolSize
             */
            final ReentrantLock mainLock = this.mainLock;

            mainLock.lock();
            try {

                /**
                 * 获得全局锁后,需再次检测当前线程池状态
                 * 原因在于预防两种非法情况:
                 *  1.线程工厂创建线程失败
                 *  2.在锁被获取之前,线程池就被关闭了
                 */
                int rs = runStateOf(ctl.get());

                /**
                 * 只有两种情况是允许添加work进入works集合的
                 * 也只有进入workers集合后才是真正的工作线程,并开始执行任务
                 *  1.线程池状态为RUNNING(即rs<SHUTDOWN)
                 *  2.线程池状态为SHUTDOWN且传入一个空任务
                 *  (理由参见:小问答之快速检测线程池状态?) 
                 */
                if (rs < SHUTDOWN ||
                    (rs == SHUTDOWN && firstTask == null)) {

                    /**
                     * 若线程处于活动状态时,说明线程已启动,需要立即抛出"线程状态非法异常"
                     * 原因是线程是在后面才被start的,已被start的不允许再被添加到workers集合中
                     * 换句话说该方法新增线程时,而线程是新的,本身应该是初始状态(new)
                     * 可能出现的场景:自定义线程工厂newThread有可能会提前启动线程
                     */
                    if (t.isAlive())
                        throw new IllegalThreadStateException();

                    //由于加锁,所以可以放心的加入集合
                    workers.add(w);
                    int s = workers.size();

                    //更新最大工作线程数,由于持有锁,所以无需CAS
                    if (s > largestPoolSize)
                        largestPoolSize = s;

                    //确认新建的worker已被添加到workers集合中  
                    workerAdded = true;
                }
            } finally {

                //千万不要忘记主动解锁
                mainLock.unlock();
            }

            /**
             * 一旦新建工作线程被加入工作线程集合中,就意味着其可以开始干活了
             * 有心的您肯定发现在线程start之前已经释放锁了
             * 原因在于一旦workerAdded为true时,说明锁的目的已经达到
             * 根据最小化锁作用域的原则,线程执行任务无须加锁,这是种优化
             * 也希望您在使用锁时尽量保证锁的作用域最小化
             */
            if (workerAdded) {

                /**
                 * 启动线程,开始干活啦
                 * 若您看过笔者的"并发番@Thread一文通"肯定知道start()后,
                 * 一旦线程初始化完成便会立即调用run()方法
                 */
                t.start();

                //确认该工作线程开始干活了
                workerStarted = true;
            }
        }
    } finally {

        //若新建工作线程失败或新建工作线程后没有成功执行,需要做新增失败处理
        if (!workerStarted)
            addWorkerFailed(w);
    }

    //返回结果表明新建的工作线程是否已启动执行
    return workerStarted;
}

Небольшой вопрос: каково значение случаев 1.2, 2.1 и 2.3 при быстром определении статуса потока?
友情小提示:读者可以反问自己 -> 何时新增Worker才是有意义的呢?传入一个空任务的目的是什么?

Небольшой ответ: Прежде чем прояснить этот вопрос, давайте проясним два момента знания:

1. Целью добавления нового воркера является обработка задач, а источник задач делится на начальные задачи и задачи очереди (т.е. оставшиеся ожидающие задачи)

2. Пул потоков не может получать новые задачи в состоянии НЕРАБОТАЕТ. Другими словами, вы собираетесь уйти с работы. Вы все еще хотите получать новые требования?

Для 2.1 -> статус пула потоков == ВЫКЛЮЧЕНИЕ, но firstTask! = null, новый рабочий не разрешен
  • Когда состояние пула потоков SHUTDOWN, так как ему не разрешено получать новые задачи, после того, как firstTask! = null требует прямого отклонения
Для 2.2 -> статус пула потоков == SHUTDOWN и firstTask == null, но очередь пуста, новый рабочий не разрешен
  • Когда firstTask имеет значение null, это означает, что целью вызова addWorker() не является обработка новых задач.
  • Тогда его целью должна быть обработка оставшихся задач, то есть задач в очереди, и как только очередь опустеет, нет необходимости добавлять нового воркера.
Для 1.2 -> Если статус пула потоков == SHUTDOWN, firstTask должен иметь значение null и очередь не пуста, тогда разрешены новые рабочие процессы
  • Когда состояние пула потоков SHUTDOWN (вызов shutdown()), новые задачи в это время не разрешены, поэтому firstTask должен иметь значение null
  • Но остальные задачи нужно обрабатывать, поэтому очередь должна быть непустой, иначе у нового рабочего потока не будет задач, что бессмысленно.

Вывод: цель передачи пустой задачи — добавить рабочие потоки для обработки оставшихся задач в очереди задач.


Небольшой вопрос: как на самом деле начинает работать поток, то есть когда начинается runWorker()?
友情小提示:结合Thread和Worker的构造器考虑一下

Небольшой ответ: Автор использует некоторые "оппортунистические" (очень умные) методы написания в задаче на выполнение потока.Давайте проанализируем класс Worker.

private final class Worker
        extends AbstractQueuedSynchronizer

        //步骤1:实现Runnable接口,从而自身是个Runnable,可以调用run方法
        implements Runnable{ 
    Worker(Runnable firstTask) {
        setState(-1); 
        this.firstTask = firstTask;

        //步骤2:newThread()的参数传入的是this,即Worker本身,注意Worker是Runnable
        this.thread = getThreadFactory().newThread(this);    
    }    

    /**
     * 步骤3:调用run()最终执行runWorker()
     * - 在addWorker()中会使用 worker.thread.start()启动线程
     * - thread启动后会立即调用run()方法,这就意味着启动调用会经历这样的过程:
     *   worker = new Worker(Runnable) - > thread = newThread(worker) -> thread.start() -> 
     *   thread.run()[JVM自动调用] -> worker.run() -> threadPoolExecuter.runWorker(worker)
     */
    public void run() {
        runWorker(this);
    }
}
Завершение начального вызова будет проходить через следующий процесс:

(1) worker = new Worker(Runnable) --> (2) thread = newThread(worker) --> (3) thread.start() --> (4) thread.run() [автоматический вызов JVM] - -> (5) worker.run() --> (6) threadPoolExecuter.runWorker(worker)


5.3 runWorker() — выполнение задач

final void runWorker(Worker w) {

    //读取当前线程 -即调用execute()方法的线程(一般是主线程)
    Thread wt = Thread.currentThread();

    //读取待执行任务
    Runnable task = w.firstTask;

    //清空任务 -> 目的是用来接收下一个任务
    w.firstTask = null;

    /**
     * 注意Worker本身也是一把不可重入的互斥锁!
     * 由于Worker初始化时state=-1,因此此处的解锁的目的是:
     * 将state-1变成0,因为只有state>=0时才允许中断;
     * 同时也侧面说明在worker调用runWorker()之前是不允许被中断的,
     * 即运行前不允许被中断
     */
    w.unlock();

    //记录是否因异常/错误突然完成,默认有异常/错误发生
    boolean completedAbruptly = true;
    try {

        /**
         * 获取任务并执行任务,取任务分两种情况:
         *   1.初始任务:Worker被初始化时赋予的第一个任务(firstTask)
         *   2.队列任务:当firstTask任务执行好后,线程不会被回收,而是之后自动自旋从任务队列中取任务(getTask)
         *     此时即体现了线程的复用
         */
        while (task != null || (task = getTask()) != null) {

            /**
             * Worker加锁的目的是为了在shutdown()时不要立即终止正在运行的worker,
             * 因为需要先持有锁才能终止,而不是为了处理并发情况(注意不是全局锁)
             * 在shutdownNow()时会立即终止worker,因为其无须持有锁就能终止
             * 关于关闭线程池下文会再具体详述
             */
            w.lock();

            /**
             * 当线程池被关闭且主线程非中断状态时,需要重新中断它
             * 由于调用线程一般是主线程,因此这里是主线程代指调用线程
             */
            if ((runStateAtLeast(ctl.get(), STOP) ||
                 (Thread.interrupted() &&
                    runStateAtLeast(ctl.get(), STOP))) &&
                        !wt.isInterrupted())
                wt.interrupt();
            try {

                /**
                 * 每个任务执行前都会调用"前置方法",
                 * 在"前置方法"可能会抛出异常,
                 * 结果是退出循环且completedAbruptly=true,
                 * 从而线程死亡,任务未执行(并被丢弃)
                 */
                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 {

                //清空任务 help gc
                task = null;

                //无论成功失败任务数都要+1,由于持有锁所以无须CAS
                w.completedTasks++;

                //必须要主动释放锁
                w.unlock();
            }
        }

        //无异常时需要清除异常状态
        completedAbruptly = false;
    } finally {

        /**
         * 工作线程退出循环的原因有两个:
         *  1.因意外的错误/异常退出
         *  2.getTask()返回空 -> 原因有四种,下文会详述
         * 工作线程退出循环后,需要执行相对应的回收处理
         */
        processWorkerExit(w, completedAbruptly);
    }
}

Небольшой вопрос: почему новая задача не помещается напрямую в очередь задач, а выполняется новым потоком?
小提示:主要是为了减少不必要的开销,从而提供性能

Небольшой ответ: цель не помещать новые задачи непосредственно в очередь задач состоит в том, чтобы уменьшить количество операций постановки в очередь и удаления из очереди очереди задач, поскольку сама очередь задач является блокирующей очередью, поэтому ее операции постановки в очередь и удаления из очереди будут включать операции блокировки и параллельные операции. обработка.


5.4 getTask() — получить задачу

Есть 5 причин, по которым метод getTask() возвращает null:

1. Пул потоков закрыт, и статус (STOP || TIDYING || TERMINATED)

2. Пул потоков закрыт, статус SHUTDOWN и очередь задач пуста.

3. Фактическое количество рабочих потоков превышает максимальное количество рабочих потоков.

4. После того, как рабочий поток удовлетворяет условию тайм-аута, он также удовлетворяет любому из следующих условий:

4.1 В пуле потоков есть по крайней мере еще один доступный рабочий поток.

4.2 В пуле потоков нет других доступных рабочих потоков, но очередь задач пуста

private Runnable getTask() {

    // 记录任务队列的poll()是否超时,默认未超时
    boolean timedOut = false; 

    //自旋获取任务
    for (;;) {

        /**
         * 线程池会依次判断五种情况,满足任意一种就返回null:
         *    1.线程池被关闭,状态为(STOP || TIDYING || TERMINATED)
         *    2.线程池被关闭,状态为SHUTDOWN且任务队列为空
         *    3.实际工作线程数超过最大工作线程数
         *    4.工作线程满足超时条件后,同时符合下述的任意一种情况:
         *      4.1 线程池中还存在至少一个其他可用的工作线程
         *      4.2 线程池中已没有其他可用的工作线程但任务队列为空
         */
        int c = ctl.get();
        int rs = runStateOf(c);

        /**
         * 判断线程池状态条件,有两种情况直接返回null
         *  1.线程池状态大于SHUTDOWN(STOP||TIDYING||TERMINATED),说明不允许再执行任务
         *    - 因为>=STOP以上状态时不允许接收新任务同时会中断正在执行中的任务,任务队列的任务也不执行了         
         *  
         *  2.线程池状态为SHUTDOWN且任务队列为空,说明已经无任务可执行
         *    - 因为SHUTDOWN时还需要执行任务队列的剩余任务,只有当无任务才可退出
         */
        if (rs >= SHUTDOWN && (rs >= STOP || workQueue.isEmpty())) {

            /**
             * 减少一个工作线程数
             * 值得注意的是工作线程的回收是放在processWorkerExit()中进行的
             * decrementWorkerCount()方法是内部不断循环执行CAS的,保证最终一定会成功
             * 补充:因线程池被关闭而计数减少可能与addWorker()的
             *      计数CAS自增发生并发竞争
             */
            decrementWorkerCount();
            return null;
        }

        //读取实际工作线程数
        int wc = workerCountOf(c);

        /**
         * 判断是否需要处理超时:
         *   1.allowCoreThreadTimeOut = true 表示需要回收空闲超时的核心工作线程
         *   2.wc > corePoolSize 表示存在空闲超时的非核心工作线程需要回收
         */
        boolean timed = allowCoreThreadTimeOut || wc > corePoolSize;

         /**
          * 有三种情况会实际工作线程计数-1且直接返回null
          *
          *    1.实际工作线程数超过最大线程数
          *    2.该工作线程满足空闲超时条件需要被回收:
          *       2.1 当线程池中还存在至少一个其他可用的工作线程
          *       2.2 线程池中已没有其他可用的工作线程但任务队列为空
          *  
          * 结合2.1和2.2我们可以推导出:
          *
          *   1.当任务队列非空时,线程池至少需要维护一个可用的工作线程,
          *     因此此时即使该工作线程超时也不会被回收掉而是继续获取任务
          *
          *   2.当实际工作线程数超标或获取任务超时时,线程池会因为
          *     一直没有新任务可执行,而逐渐减少线程直到核心线程数为止;
          *     若设置allowCoreThreadTimeOut为true,则减少到1为止;
          *
          * 提示:由于wc > maximumPoolSize时必定wc > 1,因此无须比较
          * (wc > maximumPoolSize && workQueue.isEmpty()) 这种情况
          */
        if ((wc > maximumPoolSize || (timed && timedOut))
            && (wc > 1 || workQueue.isEmpty())) {

            /**
             * CAS失败的原因还是出现并发竞争,具体参考上文
             * 当CAS失败后,说明实际工作线程数已经发生变化,
             * 必须重新判断实际工作线程数和超时情况
             * 因此需要countinue
             */
            if (compareAndDecrementWorkerCount(c))
                return null;
           /**        

            */                
            continue;
        }

        //若满足获取任务条件,根据是否需要超时获取会调用不同方法

        try {

           /**
            * 从任务队列中取任务分两种:
            *  1.timed=true 表明需要处理超时情况
            *   -> 调用poll(),超过keepAliveTime返回null
            *  2.timed=fasle 表明无须处理超时情况
            *   -> 调用take(),无任务则挂起等待
            */
            Runnable r = timed ?
                workQueue.poll(keepAliveTime, TimeUnit.NANOSECONDS) :
                workQueue.take();

            //一旦获取到任务就返回该任务并退出循环
            if (r != null)
                return r;

            //当任务为空时说明poll超时
            timedOut = true;

            /**
             * 关于中断异常获取简单讲一些超出本章范畴的内容
             * take()和poll(long timeout, TimeUnit unit)都会throws InterruptedException
             * 原因在LockSupport.park(this)不会抛出异常但会响应中断;
             * 但ConditionObject的await()会通过reportInterruptAfterWait()响应中断
             * 具体内容笔者会在阻塞队列相关番中进一步介绍
             */
        } catch (InterruptedException retry) {

            /**
             * 一旦该工作线程被中断,需要清除超时标记
             * 这表明当工作线程在获取队列任务时被中断,
             * 若您不对中断异常做任务处理,线程池就默认
             * 您希望线程继续执行,这样就会重置之前的超时标记
             */
            timedOut = false;
        }
    }
}

Небольшой вопрос: почему тайм-аут опроса, когда задача пуста?
友情小提示:可以联想一下阻塞队列操作接口

Небольшой ответ: Для этой задачи нам нужно только посмотреть на рисунок ниже, причина в том, что взятие является блокирующей операцией.
225554_IqGM_2246410.png-17.8kB

Дополнение: Автор помнит "АР бросок", "ОП ткань супер", "ПТ сопротивление"... можно


6. Закройте пул потоков

Существует два основных способа закрытия пула потоков, разница между которыми заключается в следующем:

shutdown() :Все остальные задачи в очереди выполняются, а затем завершаются.

shutdownNow() :Отказаться от оставшихся задач в очереди на выполнение, но вернуть их

Что у них общего:

1. Выполняемая задача будет продолжать выполняться и не будет прекращена или отменена.

2. Новые отправленные задачи будут отклонены напрямую

6.1 shutdown() — корректное завершение работы

Использование shutdown() для закрытия пула потоков в основном выполняет пять операций:

1. Получить глобальную блокировку

2. CAS spin изменяет состояние пула потоков на SHUTDOWN.

3. Прервите все простаивающие рабочие потоки (установите флаг прерывания) -> обратите внимание, что он простаивает

4. Снимите глобальную блокировку

5. Попробуйте завершить пул потоков

/**
 * 有序关闭线程池
 * 在关闭过程中,之前已提交的任务将被执行(包括正在和队列中的),
 * 但新提交的任务会被拒绝
 * 如果线程池已经被关闭,调用该方法不会有任何附加效果
 */
public void shutdown() {

    //1.获取全局锁
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {
        checkShutdownAccess();

        //2.CAS自旋变更线程池状态为SHUTDOWN
        advanceRunState(SHUTDOWN);

        //3.中断所有空闲工作线程
        interruptIdleWorkers();

        //专门提供给ScheduledThreadPoolExecutor的钩子方法
        onShutdown();
    } finally {

        //4.释放全局锁
        mainLock.unlock();
    }

    /**
     * 5.尝试终止线程池,此时线程池满足两个条件:
     *   1.线程池状态为SHUTDOWN
     *   2.所有空闲工作线程已被中断
     */
    tryTerminate();
}

6.2 shutdownNow() — немедленное отключение

Использование shutdownNow() для закрытия пула потоков в основном выполняет шесть операций:

1. Получить глобальную блокировку

2. CAS spin изменяет состояние пула потоков на SHUTDOWN.

3. Прервать все рабочие потоки (установить флаг прерывания)

4. Верните оставшиеся задачи в список и очистите очередь задач.

5. Снимите глобальную блокировку

6. Попробуйте завершить пул потоков

/**
 * 尝试中断所有工作线程,并返回待处理任务列表集合(从任务队列中移除)
 *
 * 1.若想等待执行中的线程完成任务,可使用awaitTermination()
 * 2.由于取消任务操作是通过Thread#interrupt实现,因此
 *   响应中断失败的任务可能永远都不会被终止(谨慎使用!!!)
 *   响应中断失败指的是您选择捕获但不处理该中断异常
 */
public List<Runnable> shutdownNow() {
    List<Runnable> tasks;

    //1.获取全局锁
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {
        checkShutdownAccess();

        //2.CAS自旋更新线程池状态为STOP
        advanceRunState(STOP);

        //3.中断所有工作线程
        interruptWorkers();

        //4.将剩余任务重新放入一个list中并清空任务队列
        tasks = drainQueue();
    } finally {

        //5.释放全局锁
        mainLock.unlock();
    }

    /**
     * 6.尝试终止线程池,此时线程池满足两个条件:
     *   1.线程池状态为STOP
     *   2.任务队列为空
     * 注意:此时不一定所有工作线程都被中断回收,详述见
     *       7.3 tryTerminate
     */
    tryTerminate();

    //5.返回待处理任务列表集合
    return tasks;
}

6.3 awaitTermination() — ожидание завершения работы пула потоков

При закрытии пула потоков awaitTermination() будет блокироваться до тех пор, пока не произойдет одно из следующих событий:

1. Все задания выполнены:Пул потоков вызывает толькоtryTerminated()Попытка завершить работу пула потоков и успешно изменить состояние наTERMINATEDбудет вызван послеtermination.signalAll(), после этого заблокированный поток снова будет оценивать статус после его пробуждения.TERMINATEDвыйдет

2. Достигнуть тайм-аута блокировки: termination.awaitNanos()После достижения овертайма оставшееся время будет возвращено (на этот раз 0), а затем будет снова оценено, чтобы удовлетворитьnano==0Привести кreturn false, то есть ожидание не удается

3. Текущий поток прерывается:Если текущий поток (основной поток) прерывается, поток выбрасываетInterruptExceptionИсключение прерывания, если обработка исключения не выполняется, оно будет разблокировано из-за исключения

public boolean awaitTermination(long timeout, TimeUnit unit)
        throws InterruptedException {
    long nanos = unit.toNanos(timeout);
    final ReentrantLock mainLock = this.mainLock;

    //1.获取全局锁
    mainLock.lock();
    try {
        for (;;) {

            //2.所有任务执行完毕,等待成功而退出
            if (runStateAtLeast(ctl.get(), TERMINATED))
                return true;

            //3.到达阻塞超时时间,等待失败而退出
            if (nanos <= 0)
                return false;

            nanos = termination.awaitNanos(nanos);
        }
    } finally {

        //4.释放全局锁
        mainLock.unlock();
    }
}

Вы можете узнать, действительно ли пул потоков закрыт:

//关闭线程池
threadPoolExecutor.shutdown();
try{
    //循环调用等待任务最终全部完成
    while(!threadPoolExecutor.awaitTermination(300, TimeUnit.MILLISECONDS)) {
        logger.info("task executing...");
    }
    //此时剩余任务全部执行完毕,开始执行终止流程
    logger.info("shutdown completed!")
} catch (InterruptedException e) {
    //中断处理
}

7. Обработка прерываний и завершений

7.1 interruptIdleWorkers() — Прерывает бездействующие потоки

У Worker есть следующие четыре критерия для обработки прерываний (мы еще раз вернемся к предыдущим знаниям):

1. Когда рабочий поток фактически начинает выполняться, его нельзя прерывать.

2. Когда рабочий поток выполняет задачу, его нельзя прерывать.

3. Когда рабочий поток ожидает получения задачи из очереди задачgetTask()быть прерванным

4. позвонитьinterruptIdleWorkers()Рабочая блокировка должна быть получена первой при прерывании бездействующего потока.

/**
 * 中断全部空闲线程
 */
private void interruptIdleWorkers() {
    interruptIdleWorkers(false);
}

/**
 * 中断未上锁且在等待任务的空闲线程
 * 中断的作用在于便于处理终止线程池或动态控制的情况
 *
 * @param onlyOne 为true时为中断一个,为false时为中断全部
 */
private void interruptIdleWorkers(boolean onlyOne) {
    //加全局锁
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {

        /**
         * 循环方式中断工作线程
         * 这里也体现了workers集合的核心作用之一
         */
        for (Worker w : workers) {
            Thread t = w.thread;

            /**
             * 非中断且成功获取到worker锁的工作线程才允许被中断
             * 
             *  1.已被中断的工作线程无须再次标记中断
             *
             *  2.w.tryLock()体现了Worker作为一把锁的核心作用:
             *  即控制线程中断 -> 当线程还在运行中是不允许被中断的
             * 
             *  3.具体可以参见runWorker()方法,运行前都是调用lock()
             * 
             *  4.由于该方法只会在shutdown()中调用,间接也说明
             *  shutdown()只会中断在该方法中获取到worker锁
             *  的空闲线程(此时线程正在获取新任务getTask(),还没上锁)
             */
            if (!t.isInterrupted() && w.tryLock()) {
                try {

                    //中断工作线程
                    t.interrupt();

                } catch (SecurityException ignore) {
                } finally {

                    //注意这里释放的是worker锁,对应tryLock()
                    w.unlock();
                }
            }

            //onlyOne为true时,只随机中断一个空闲线程(Set可是无序的哦)
            if (onlyOne)
                break;
        }
    } finally {

        //释放全局锁
        mainLock.unlock();
    }
}

7.2 interruptWorkers() — прерывает все потоки

/**
 * 中断所有线程,包括正在执行任务的线程
 * 该方法只提供给shutdownNow()使用
 */
private void interruptWorkers() {
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {

        //循环设置中断标志 
        for (Worker w : workers)
            w.interruptIfStarted();
    } finally {
        mainLock.unlock();
    }
}

/**
 * Worker实现的中断方法
 */
void interruptIfStarted() {
    Thread t;

    /**
     * 当线程池非RUNNING状态 && 线程非空 && 线程非中断
     * 三者同时满足时才允许中断
     * 
     * 为什么线程池必须非RUNNING状态才允许中断呢?
     *  因为该方法只提供给interruptWorkers()使用
     *  而interruptWorkers()只提供给shutdownNow()使用
     *  因此此时线程状态应为STOP
     */
    if (getState() >= 0 && (t = thread) != null 
            && !t.isInterrupted()) {
        try {

            //设置中断标志
            t.interrupt();
        } catch (SecurityException ignore) {
        }
    }
}

7.3 tryTerminate() — попытка завершить пул потоков

Прежде чем разбирать tryTerminate(), давайте рассмотрим несколько важных вопросов.


Небольшой вопрос: почему нельзя прерывать рабочий поток, выполняющий задачу?
友情小提示:工作线程执行任务前需加worker锁且该锁非重入

Маленький ответ:обзорinterruptIdleWorkers()Мы находим, что в (1) мы должны сначала вызватьtryLock()Рабочий поток может быть прерван только после успешного получения рабочей блокировки, потому что(2) После того, как рабочий поток получит задачу и перед выполнением задачи, сначала будет добавлена ​​рабочая блокировка, и рабочая блокировка не является повторной., что означает, что рабочий поток, выполняющий задачу, не может быть прерван


Небольшой вопрос: что считается бездействующим рабочим потоком?
友情小提示:需要获取worker锁有两个时机,一个是shutdown(),一个真正执行任务之前

Небольшой ответ: рабочий поток должен удерживать блокировку при выполнении задачи, и только когда задача getTask() получена из очереди задач, не нужно получать рабочую блокировку, поэтому рабочий поток можно разделить на два состояния:

1. Рабочий поток, выполняющий задачу:Рабочий поток, который выполняет task.run() после получения рабочей блокировки.

2. Простой рабочий поток:Рабочий поток, который получает задачу из очереди задач (включая рабочий поток, который только что получил задачу и заблокирован из-за отсутствия задачи)


Небольшой вопрос: как прерывание потока влияет на перезапуск потока?
友情小提示:核心在于当getTask()返回null时会退出runWorker()并执行processWorkerExit()

Небольшой ответ: давайте сначала рассмотрим ситуацию, когда geTask() возвращает null:

Есть 5 причин, по которым метод getTask() возвращает null:

1. Пул потоков закрыт, и статус (STOP || TIDYING || TERMINATED)

2. Пул потоков закрыт, статус SHUTDOWN и очередь задач пуста.

3. Фактическое количество рабочих потоков превышает максимальное количество рабочих потоков.

4. После того, как рабочий поток удовлетворяет условию тайм-аута, он также удовлетворяет любому из следующих условий:

4.1 В пуле потоков есть как минимум еще один доступный рабочий поток.

4.2 В пуле потоков нет других доступных рабочих потоков, но очередь задач пуста

Из приведенного выше видно, что когда пул потоков закрыт, время ожидания потока или поток динамического управления (например, размер пула, время ожидания и т. д.) могут привести к тому, что getTask() вернет значение null, как это делает getTask( ) влияет на переработку?

Возьмем только закрытие пула потоков в качестве примера (остальные случаи — лишь разница между условными суждениями) и опишем логику, которая будет происходить после прерывания:

1. При блокировкеgetTask()Выдает, когда рабочий поток прерываетсяInterruptedExceptionПрервать исключение, затем разблокировать и повторно получить задачу

2. Повторное получение задач все равно требует повторной проверки условий получения задачи.Когда пул потоков закрыт, например, вызовshutdown(), состояние пула потоков становится SHUTDOWN, а поскольку очередь задач в это время пуста,getTask()Возвращает null напрямую; если вызываетсяshutdownNow(), состояние пула потоков становится STOP, а затем напрямую возвращается null

3. ВrunWorker()метод, когдаgetTask()После возврата null он выйдет из цикла, а затем вызоветprocessWorkerExit()Операция повторного использования потока метода

Стоит отметить: механизм прерывания JAVA просто устанавливает флаг прерывания, поэтому вы делаете это самостоятельно в задаче.Thread.currentThread().interrupt()Это не повлияет на поток, чтобы продолжить выполнение задачи и перезапуск потока, и вы не можете получить это в задаче.InterruptedException(ошибка компиляции), причина в том, что он уже находится вgetTask()захвачен


Небольшой вопрос: поскольку состояние пула потоков изменяется после закрытия пула потоков, и прерванный поток будет перезапущен, зачем вам нужно выполнять tryTerminate()?
友情小提示:调用shutdown()后,interruptIdleWorkers()只会中断空闲工作线程,那么当时正在执行任务的工作线程执行完后怎么办呢?

Маленький ответ:перечислитьshutdown()После выполнения задачи рабочие потоки, выполняющие задачу, не будут прерваны. Когда они завершат задачу, при условии, что очередь не пуста, эти рабочие потоки продолжат выполнение оставшихся задач, пока не будут заблокированы. уменьшается, фактическое количество рабочих потоков будет продолжать уменьшаться до минимального количества обслуживания; когда очередь пуста, рабочие потоки с минимальным количеством обслуживания всегда будут блокироваться вworkerQueue.take()Выше он никогда не может быть остановлен, и пул потоков не будет получать новые отправленные задачи после его закрытия.

противПроблема, заключающаяся в том, что рабочий поток, выполняющий задачу, не может быть завершен при его прерывании., решение Дуга Ли:

- для вызова в любом месте, которое может привести к завершению пула потоковtryTerminate(), метод определит, вступил ли пул потоков в процесс завершения, если в это время еще есть потоки, он повторно прервет простаивающий рабочий поток.

Завершить процесс: состояние пула потоков — SHUTDOWN, а очередь задач пуста, или состояние пула потоков — STOP


/**
 * 终止线程池 -> 最终会将线程池状态变更为TERMINATED
 * 只有同时满足下面两个条件才允许做TERMINATED的状态转变:
 *    1.线程池状态为SHUTDOWN且任务队列为空 或状态为STOP
 *    2.线程池中已没有存活的工作线程 -> 实际工作线程为0
 */
final void tryTerminate() {

    //自旋
    for (;;) {

        //获取线程池控制器
        int c = ctl.get();

        /**
         * 有4种情况是不允许执行变更TERMINATED操作
         *
         *   1.线程池仍为运行态RUNNING,说明线程池还在正常运行中,
         *     此时是不允许尝试中断,起码要SHUTDOWN或STOP
         *     规则参见shutdown()和shutdownNow()
         *
         *   2.线程池状态已经是TIDYING或TERMINATED,
         *     前者说明变更TERMINATED正在执行中,后者说明终止已完成
         *     这两种情况都无须重复执行终止
         *
         *   3.线程池状态为SHUTDOWN且任务队列非空,
         *     说明线程池虽然已被要求关闭,但还有任务还没处理完
         *     需要等待任务队列中剩余任务被执行完毕
         */
        if (isRunning(c) ||
            runStateAtLeast(c, TIDYING) ||
            (runStateOf(c) == SHUTDOWN && ! workQueue.isEmpty()))
            return;

        /**
         * 此时线程池状态为SHUTDOWN状态且队列为空,或已是STOP状态
         * 4.若工作线程数非0,说明还有工作线程可能正在执行或等待任务中,
         * 这种情况的原因参见上文中的小问答之`为什么还要执行tryTerminate()`
         * 此时会选择中断一个空闲工作线程以确保SHUTDOWN信号的传播
         */
        if (workerCountOf(c) != 0) { // Eligible to terminate

            /**
             * 此时已经进入终止流程,为了传播SHUTDOWN信号,
             * 每次总是中断一个空闲工作线程以避免所有线程等待
             * 
             * 小问:此时若调用interruptIdleWorkers(false)呢?
             * 小答:注意每个线程的回收都会调用processWorkerExit()
             *       而该方法都会调用tryTerminate(),而此时一旦
             *       设置为true(表示全部)的话,由于中断操作前必须
             *       通过worker.tryLock()加锁,因此就可能因锁竞争
             *       造成不必要的大量等待,还不如一个个执行
             *
             * 小问:那么为什么shutdown()的时候可以为true呢?
             * 小答:那是因为空闲线程都是没有持有worker锁的!
             *      那么就不会出现锁竞争带来的不必要的开销
             */
            interruptIdleWorkers(ONLY_ONE);
            return;
        }

        /**
         * 当进入终止流程且无存活的工作线程时
         * 那么就可以terminate终止线程池了
         */

        //1.获取全局锁
        final ReentrantLock mainLock = this.mainLock;
        mainLock.lock();
        try {

            /**
             * 2.先尝试变成TIDYING状态
             *   1.一旦成功,执行🐶方法terminated()
             *   2.CAS失败后会重试,失败原因可能是线程池刚好
             *     已被设置为TERMINATED,即线程池终止已经完成,
             *     之后在重新循环中会因runStateAtLeast(c, TIDYING)
             *     而退出该方法
             */
            if (ctl.compareAndSet(c, ctlOf(TIDYING, 0))) {
                try {

                    //3.执行终止
                    terminated();
                } finally {

                    //4.设置TERMINATED状态
                    ctl.set(ctlOf(TERMINATED, 0));

                    /**
                     * 5.通过唤醒解除条件阻塞
                     * 当关闭线程池后需要等待剩余任务完成才真正终止线程池,
                     * 会调用awaitTermination()方法,
                     * 此时主线程会被
                     * 
                     */
                    termination.signalAll();
                }
                return;
            }
        } finally {

            //6.释放全局锁
            mainLock.unlock();
        }
        // else retry on failed CAS
    }
}

8. Сбой потока и повторная обработка

8.1 addWorkerFailed() — добавлена ​​обработка ошибок потока

Обработка отказа при добавлении нового рабочего потока выполняется следующим образом:

1. Получите глобальную блокировку

2. Удалить работника из коллекции рабочих

3. Спин CAS уменьшает фактическое количество рабочих потоков.

4. Попробуйте завершить пул потоков

5. Снимите глобальную блокировку

/**
 * 新增工作线程失败处理
 */
private void addWorkerFailed(Worker w) {

    //1.获取全局锁 -> 目的是为了安全更新workers
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {

        //2.从workers集合中移除该worker
        if (w != null)
            workers.remove(w);

        /**
         * 3.CAS自旋减少实际工作线程计数 -> 最终会成功
         * 小问:为何已经加锁还是使用CAS?
         * 小答:workers必须在持有锁环境下使用,ctl无须在持有锁环境下使用
         *   1.workers集合为非线程安全的HashSet,不能使用CAS只能加锁(即外部控制方式)
         *   2.ctl为AtomicInteger原子类型,因此可以直接使用CAS维护(即内部控制方式)
         * 注意:这里说的持有锁指的是持有全局锁mainLock,虽然ReentrantLock底层实现也是CAS         
         */
        decrementWorkerCount();

        /**
         * 4.尝试终止线程池
         * 
         * 小问:那么为什么此时要尝试终止线程池呢?
         * 小答:因为新增线程失败的原因只有一个
         *       -> 线程池被关闭并进入终止流程
         *      具体可参见addWorker()方法
         */
        tryTerminate();
    } finally {

        //5.释放全局锁
        mainLock.unlock();
    }
}

8.2 processWorkerExit() — обработка повторного использования потока

Процесс переработки нитей в основном делится на две части:

1. Перезапустить рабочий поток

2. При необходимости добавьте новые рабочие потоки.

1. Существует 6 основных шагов по переработке рабочего потока:

1. Потоки, которые были внезапно прерваны из-за исключений ошибок, фактическое количество рабочих потоков -1

2. Получить глобальную блокировку

3. Подсчитайте общее количество выполненных задач в пуле потоков

4. Безопасно удалите воркера из коллекции воркеров

5. Снимите глобальную блокировку

6. Попробуйте завершить пул потоков

2. Если пул потоков находится в состоянии RUNNING или SHUTDOWN, существует две ситуации, в которых необходимо добавить новый рабочий поток:

1. Поток был неожиданно уничтожен из-за исключения ошибки

2. Если это не случайная смерть, гарантированно выживет хотя бы минимальное количество доступных рабочих потоков.

private void processWorkerExit(Worker w, boolean completedAbruptly) {

    //1.因错误异常而被意外死亡的线程,实际工作线程计数-1
    if (completedAbruptly) 
        // If abrupt, then workerCount wasn't adjusted 作者大佬的注释真的没写错吗....
        decrementWorkerCount();

    //2.获取全局锁,主要目的是为了安全将worker从workers集合中移除 
    final ReentrantLock mainLock = this.mainLock;
    mainLock.lock();
    try {

        //3.统计线程池总完成任务数
        completedTaskCount += w.completedTasks;

        //4.将该worker从workers集合中安全移除
        workers.remove(w);
    } finally {

        //5.释放全局锁
        mainLock.unlock();
    }

    /**
     * 6.尝试终止线程池
     * 小问:为什么此处需要尝试终止线程池?
     * 小答:由于processWorkerExit()方法只会在
     *      runWorker()中调用,而调用的时机有两个:
     *          1.工作线程因错误异常而被中断退出
     *          2.getTask()返回null
     *      根据tryTerminate()的终止条件可知,
     *      前者实际上并不会终止线程池,但问题是
     *      后者的getTask()是有可能因进入终止流程而返回null
     */
    tryTerminate();

    int c = ctl.get();

    /**
     * 若线程池状态为RUNNING或SHUTDOWN时,有两种情况需要新增工作线程
     *  1.线程因错误异常而被意外死亡
     *    -> 目的是填补这个意外死亡的工作线程造成的线程缺口(填坑)
     *  2.若非意外死亡,则至少保证有最小存活数个可用工作线程存活
     *    -> 目的是保证线程池正常运行或SHUTDOWN时有能力完成队列剩余任务
     */
    if (runStateLessThan(c, STOP)) {
        if (!completedAbruptly) {

            /**
             * 线程最小存活数由allowCoreThreadTimeOut和队列长度共同决定
             * 1.当allowCoreThreadTimeOut为true时,若队列非空,
             *   至少保证一个可用线程存活
             * 2.当allowCoreThreadTimeOut为false时,实际工作线程数
             *   一旦超过核心工作线程数,无须再新增工作线程了
             */
            int min = allowCoreThreadTimeOut ? 0 : corePoolSize;

            //1.若允许响应核心工作线程超市且队列非空时
            if (min == 0 && ! workQueue.isEmpty())
                //至少保证一个可用线程可用
                min = 1;

            //2.实际工作线程数一旦超过核心工作线程数,无须再新增线程了
            if (workerCountOf(c) >= min)

                // replacement not needed
                //"替换"指的是替已死亡的线程继续填坑(完成剩余任务)
                return; 
        }

        /** 
         * 新增工作线程根据原因区分的目的有两个:
         *   1.因意外死亡的:
         *     -> 目的是为了填补线程空缺
         *   2.非意外死亡正常退出且队列非空:
         *     -> 处理任务队列中的剩余任务
         * 虽然目的有区别,但实际上作用是一致的:
         *     -> 都是为了处理队列任务(因为firstTask为null)
         */
        addWorker(null, false);
    }
}

9. Очередь задач и стратегия организации очереди

благодарныйРазговор о параллелизме (7) — блокирующая очередь в Java
Очередь задач — это блокирующая очередь, используемая для хранения задач, ожидающих выполнения (здесь конкретно относится к реализацииBlockingQueueКласс реализации интерфейса очереди блокировки), его целью является реализация кэширования и совместного использования данных; пакет concurrent изначально предоставляет 7 видов очередей блокировки, которые можно разделить на две части в соответствии с границей:

- Ограниченная очередь:Ограниченная очередь относится к очереди с ограниченной емкостью, которая не допускает неограниченного расширения.Integer.MAX_VALUE, как постановка в очередь, так и удаление из очереди могут блокировать

  • Ограниченная очередь (bounded):Должна быть приведена начальная мощность, в том числеArrayBlockingQueue

  • Необязательно ограниченный:Если начальная емкость не установлена, максимальная емкость по умолчанию равнаInteger.MAX_VALUE,включаютLinkedBlockingQueueиLinkedBlockingDeque

- неограниченная очередь: Неограниченная очередь относится к неограниченной, есть два случая: 0 и неограниченная.

  • Неограниченный (0):Емкость равна 0, никакие элементы не сохраняются и нет блокировки, напримерSynchronousQueue

  • Неограниченный:Разрешить неограниченное расширение емкости, пока не будет брошеноOutOfMemoryError, очередь не будет заблокирована, и очередь может быть заблокирована, в том числеDelayQueue,LinkedTransferQueue,PriorityBlockingQueue

Примечание. Если не указано иное, очереди блокировки следуют правилу FIFO «первым пришел — первым обслужен».

9.1 Ограниченные очереди

Ограниченная очередь относится к блокирующей очереди с ограниченной и фиксированной емкостью, которая не допускает бесконечного расширения.По сравнению с неограниченной очередью, когда максимальное количество пулов ограничено, это может эффективно предотвратить исчерпание ресурсов, но это также увеличивает сложность управления -> потребности ограниченной очереди. является взаимным «компромиссом» между размером очереди и максимальным количеством потоков:

- Большая очередь + малый пул: эффективное сокращение накладных расходов на потоки, но может уменьшить пропускную способность.Если задачи часто блокируются, например частый ввод-вывод,
Использование больших очередей и небольших пулов сводит к минимуму использование ЦП, ресурсов операционной системы и накладные расходы на переключение контекста, но может привести к искусственному снижению пропускной способности. Если задачи часто блокируются (например, если они ограничены вводом-выводом), система может запланировать больше потоков, чем вы позволяете.

- Использование небольших очередей обычно требует больших размеров пула, более высокой загрузки ЦП, но может привести к неприемлемым издержкам планирования, что также снижает пропускную способность.

9.1.1 ArrayBlockingQueue

эффект:

- Ограниченная очередь блокировки, состоящая из структур массива

- В дополнение к массиву фиксированной длины он также включает две переменные типа int для определения положения головы и хвоста в массиве.

- Никакие дополнительные объекты не создаются и не перерабатываются при постановке в очередь и удалении из очереди.

- Поддержка честного и несправедливого режима, нечестная блокировка по умолчанию

- Внутренне принимает метод синхронизации один замок + два условия, которые не могут быть действительно одновременными

9.1.2 LinkedBlockingQueue

эффект:

- ограниченная очередь блокировки, состоящая из структуры связанного списка

- По умолчанию и максимальная длина этой очередиInteger.MAX_VALUE

- Каждый раз при постановке/удалении из очереди создается/уничтожается дополнительный объект Node, который используется для реализации структуры связанного списка.

- Связанные списки обычно имеют лучшую пропускную способность, чем списки массивов (теоретически), просто погуглите свои собственные причиныArrayList和LinkedList的区别

- Внутренне принимает метод синхронизации двух замков + два условия, действительно параллельных

-Executors.newFixedThreadPool()Используемая очередь блокировки

Яма:

При использовании емкости по умолчанию, если скорость производства намного превышает скорость потребления, ресурсы памяти могут быть исчерпаны до того, как очередь будет заполнена, а хранилище будет заблокировано.

предложение:

- обычно просто используютLinkedBlockingQueueиArrayListBlockingQueueБольшинство потребностей производства и потребления могут быть удовлетворены

9.1.3 LinkedBlockingDeque

эффект:

- двусторонняя очередь блокировки, состоящая из структуры связанного списка

- Двусторонняя очередь позволяет ставить в очередь и удалять из очереди, что проявляется во многих других методах xxFirst и xxLast.

- Если начальная емкость не задана, максимальная емкость по умолчанию для этой очереди равнаInteger.MAX_VALUE

-такой жеArrayListBlockingQueueТочно так же внутри используется метод синхронизации одна блокировка + два условия, которые не могут быть по-настоящему параллельными.

9.2 Неограниченные очереди

Неограниченная очередь относится к блокирующей очереди с бесконечной емкостью или емкостью 0. При ее использовании необходимо обратить внимание:

1. Когда емкость равна 0, тщательно задайте maxPoolSize, чтобы избежать отклонения вновь отправленных задач.

2. Когда емкость бесконечна, это означает, что maxPoolSize недействителен, установка этого значения бессмысленна, и количество создаваемых потоков не будет превышать corePoolSize.

Применимая сцена:
Неограниченные очереди наиболее подходят, когда каждая задача независима друг от друга и не влияет друг на друга.

9.2.1 SynchronousQueue

Функции:

- Блокирующая очередь, которая не хранит элементы, а также является типом очереди задач по умолчанию для пула потоков.

- Очередь не хранит задачи, а может только передавать элементы между потоками -> то есть прямая подача

- Поддерживает честный и нечестный режим, по умолчанию не честный (о справедливости см.reentrantLock)

-Executors.newCachedThreadPool()Используемая очередь блокировки

Сцены:

Эта стратегия позволяет избежать блокировок при обработке наборов запросов, которые могут иметь внутренние зависимости.

Яма:

- Когда нет задачи, доступной для немедленного выполнения, присоединение к очереди не удастся, и в это время будет добавлен новый поток; однако, если он превышает maxPoolSize, вновь отправленная задача будет отклонена!

- В непубличном режиме, если существует большой разрыв между скоростями производства и потребления, легко может произойти голодание, и некоторые данные могут никогда не выполняться

предложение:

- Прямая отправка обычно требует неограниченных максимальных размеров пула, чтобы избежать отклонения вновь отправленных задач.

- Эта политика позволяет неограниченным потокам потенциально расти, когда команды поступают последовательно в большем количестве, чем может обработать среднее количество очередей.

9.2.2 PriorityBlockingQueue

Функции:

- Неограниченная очередь блокировки с приоритетом, состоящим из структуры массива, с емкостью по умолчанию 11

- Естественный порядок по умолчанию при поддержке пользовательского порядка элементов в очереди (реализацияComparableинтерфейс)

-Алгоритм сортировки — сортировка кучей, а внутренняя синхронизация потоков использует справедливую блокировку.

- Внутренне использует блокировку + условный метод синхронизации: поскольку это неограниченная очередь, требуется только одинnotEmptyнепустое условие

- Стоит отметить, что только головной узел гарантирует порядок приоритета, остальные узлы не гарантируют

Сцены:
Когда требуется отсортированный массив
Яма:

Из-за использования сортировки кучи, когда скорость потребления намного ниже, чем скорость производства, из-за сжатия задачи и необходимости сортировки кучи, вероятно исчерпание всего пространства кучи, то есть легко переполнение памяти

предложение:
Местные самодуры могут добавлять памяти по желанию, иначе надо следить, чтобы не было слишком большого отставания по задачам

9.2.3 DelayQueue

Функции:

- Используйте приоритетную очередь для реализации упорядоченной неограниченной очереди блокировки с задержкой выборки.

- должен быть реализован элемент enqueueDelayedИнтерфейс, учитывая начальное время задержки, элемент может быть получен из очереди только после времени задержки прибытия, и элемент не может быть нулевым

- Внутренне используйте одну блокировку + одно условие + метод синхронизации приоритетной очереди: из-за характеристики задержки требуется только одинavailableУсловие указывает, доступна ли задача или нет

Сцены:

-Он используется для реализации механизма повторных попыток, а множественные отложенные выполнения могут поддерживать ограничение на количество повторных попыток.

-ScheduledThreadPoolExecutorотложенный пул потоковDelayedWorkQueueОчередь блокировки задержки — это ее оптимизированная версия, используемая для планирования времени и других операций.

- Используется для реализации кэширования, хотя рекомендуется NoSQL

-TimerQueueБазовая структура данных

9.2.4 LinkedTransferQueue

Функции:

- Неограниченная очередь блокировки, состоящая из структуры связанного списка

-TransferQueueдаConcurrentLinkedQueue,SynchronousQueue (公平模式下),无界的LinkedBlockingQueuesнадмножество

- относительно других блокирующих очередейLinkedTransferQueueпереборtryTransfer()иtransfer()метод

- когда нет потребителей, ожидающих получения элементов,transfer()Метод сохранит элемент в хвостовом узле очереди и заблокирует его до тех пор, пока потребитель не потребит элемент перед возвратом; в противном случае он будет передан непосредственно потребителю.
, не будет стоять в очереди в это время

-отличный отtransfer(),tryTransfer()Метод немедленно вернет, является ли результат операции успешным или неудачным, независимо от того, есть ли потребители, ожидающие получения элементов.В это время он не будет поставлен в очередь и не будет блокироваться.

- Сложные двойные структуры данных, используемые этим классом, все методы должны быть реализованы в два этапа:
Резерв:Когда потребитель получает элемент из очереди и обнаруживает, что очередь пуста, элемент узла Node с нулевым полем данных будет создан и помещен в очередь, а затем потребитель должен заблокировать (прокрутить ожидание), пока поле данных не будет заполнено. не пустой
перечислить:Когда производитель собирается сохранить элементы в очереди, он обнаруживает, что поле данных самого переднего элемента равно нулю, затем он напрямую присваивает данные элемента самому переднему элементу, то есть завершаетОбмен данными

10. Мониторинг пула потоков

10.1 Собственный мониторинг

Очень важно отслеживать работу пула потоков, особенно при обнаружении проблемы.К счастью, пул потоков изначально предоставляет нам несколько свойств мониторинга, которые можно получить:

1. количество задач:Количество задач, которые должен выполнить пул потоков (приблизительно)

2.completedTaskCount:Количество задач, выполненных пулом потоков во время запущенного процесса, меньше или равно taskCount

3.наибольший размер пула:Максимальное количество потоков, созданных в пуле потоков. Если значение соответствует maxPoolSize, пул потоков заполнен.

4. размер пула:Количество потоков в пуле потоков, включая количество неработающих рабочих потоков; стоит отметить, что если пул потоков не закрыт, потоки в пуле потоков не будут автоматически перезапущены, поэтому для выполняющегося потока бассейн, это значение только увеличивается или уменьшается

5.активкаунт:Количество запущенных рабочих потоков (приблизительно)

Стоит отметить, что хотя методы get этих свойств мониторинга поддерживаются глобальными блокировками, состояние и количество потоков во время работы пула потоков можно динамически регулировать, напримерallowCoreThreadTimeOut()、setMaximumPoolSize()、setCorePoolSize()、shutdown()и т. д., поэтому некоторые значения могут быть только приблизительно

10.2 Расширенный мониторинг

Пул потоков предоставляет три метода ловушек, которые можно использовать для расширения функций, таких как мониторинг среднего времени выполнения, максимального времени выполнения и минимального времени выполнения задач:
перед выполнением():родыrunWorker()метод, вrun()выполнить перед методом
после выполнения():родыrunWorker()метод, вrun()выполнить после метода
прекращено():родыtryTerminate()метод, состояние CASTIDYINGвыполнить после

Примечание. Поскольку все вышеперечисленные методыprotectedИ пул потоков по умолчанию пуст, поэтому вышеуказанные методы можно переопределить только путем наследования пула потоков или построения

11. Стратегия подавления насыщения

Политика отклонения насыщения пула потоков в основном используется для отклонения задач (но это не означает, что задача не будет выполняться).Пул потоков изначально предоставляет четыре политики отклонения насыщения, которые в основном охватывают распространенные сценарии обработки насыщения:

Политика прерывания:Стратегия по умолчанию, прямое исключение

CallerRunsPolic:Используйте только вызывающий поток для выполнения задачи

DiscardPolicy:Просто скиньте задачу

DiscardOldestPolicy:Отмените хвостовую задачу и повторите задачу с пулом потоков.

Все стратегии отклонения должны реализовывать интерфейс обработчика отказа с унифицированным калибром:

/**
 * 用于拒绝线程池任务的处理器
 */
public interface RejectedExecutionHandler {

    /**
     * 该方法用于拒绝接受线程池任务
     * 
     * 有三种情况可能调用该方法:
     *   1.没有更多的工作线程可用
     *   2.任务队列已满
     *   3.关闭线程池
     *
     * 当没有其他处理选择时,该方法会选择抛出RejectedExecutionException异常
     * 该异常会向上抛出直到execute()的调用者
     */
    void rejectedExecution(Runnable r, ThreadPoolExecutor executor);
}

11.1 CallerRunsPolicy

Правило обработки: вновь отправленная задача выполняется непосредственно вызывающим потоком.

Рекомендация: CallerRunsPolicy рекомендуется для стратегии отклонения, так как эта стратегия не отбрасывает задачи и не генерирует исключения, а откатывает задачи к вызывающему потоку для выполнения.

/**
 * 不会直接丢弃,而是直接用调用execute()方法的线程执行该方法
 * 当然一旦线程池已经被关闭,还是要丢弃的
 *
 * 补充:值得注意的是所有策略类都是public的静态内部类,
 *      其目的应该是告知使用者 -> 该类与线程池相关但无需线程池实例便可直接使用
 */
public static class CallerRunsPolicy implements RejectedExecutionHandler {

    public CallerRunsPolicy() { }

    /**
     * 直接使用调用该方法的线程执行任务
     * 除非线程池被关闭时才会丢弃该任务
     */
    public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
        //一旦线程池被关闭,丢弃该任务
        if (!e.isShutdown()) {

            //注意此时不是线程池执行该任务
            r.run();
        }
    }
}

11.2 AbortPolicy

Правила обработки: выбрасывать RejectedExecutionException напрямую

/**
 * 简单、粗暴的直接抛出RejectedExecutionException异常
 */
public static class AbortPolicy implements RejectedExecutionHandler {

    public AbortPolicy() { }

    /**
     * 直接抛出异常,但r.toString()方法会告诉你哪个任务失败了
     * 更人性化的一点是 e.toString()方法还会告诉你:
     * 线程池的状态、工作线程数、队列长度、已完成任务数
     * 建议若是不处理异常起码也要在日志里面打印一下,留个案底
     */
    public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
        throw new RejectedExecutionException(
            "Task " + r.toString() + " rejected from " + e.toString());
    }
}

11.3 DiscardPolicy

Правила обработки: напрямую отбрасывать последние отправленные задачи в соответствии с правилом LIFO (последние поступили — первыми обслужены).

/**
 * 直接丢弃任务
 * 这个太狠了,连个案底都没有,慎用啊
 */
public static class DiscardPolicy implements RejectedExecutionHandler {

    public DiscardPolicy() { }

    /**
     * 无作为即为丢弃
     */
    public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
    }
}

11.4 DiscardOldestPolicy

Правило обработки: отбросить последнюю задачу в соответствии с правилом LRU (наименее недавно использованное), затем попытаться выполнить новую отправленную задачу.

/**
 * 比起直接丢弃,该类会丢弃队列里最后一个但仍未被处理的任务,
 * 然后会重新调用execute()方法处理当前任务
 * 除非线程池被关闭时才会丢弃该任务
 * 此类充分证明了"来得早不如来的巧"
 */
public static class DiscardOldestPolicy implements RejectedExecutionHandler {

    public DiscardOldestPolicy() { }

    /**
     * 丢弃队列里最近的一个任务,并执行当前任务
     * 除非线程池被关闭时才会丢弃该任务
     * 原因是队列是遵循先进先出FIFO原则,poll()会弹出队尾元素
     */
    public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
        //一旦线程池被关闭,直接丢弃
        if (!e.isShutdown()) {

            //弹出队尾元素
            e.getQueue().poll();

            //直接用线程池执行当前任务
            e.execute(r);
        }
    }
}

12. Обработка исключений пулов потоков

12.1 Обработка исключений submit()

использоватьsubmit()При обработке исключений необходимо учитывать четыре фактора:

1. Исключение будет сохранено вFuture对象изExecutionException, вы можете позвонитьget()использоватьtry-catchЕсли есть N задач с исключениями, будет сгенерировано N исключений, но текущий рабочий поток не будет завершен.

2. Отдельные настройкиUncaughtExceptionHandlerЭто бесполезно, но эффективно в сочетании с (3)

3. Разрешено вsubmit()внутри методаtry-catchПерехват этого исключения также не приведет к завершению текущего потока.

4. Если вы хотите обрабатывать исключения внутри, вы также можете переписатьafterExecute()методы, такие как:

static ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(2, 3, 3, TimeUnit.SECONDS, new SynchronousQueue<>()) {
    //构造时直接重写afterExecute()方法
    protected void afterExecute(Runnable r, Throwable t) {
        super.afterExecute(r, t);
        printException(r, t);
    }
};

private static void printException(Runnable r, Throwable t) {
    if (t == null && r instanceof Future<?>) {
        try {
            Future<?> future = (Future<?>) r;
            if (future.isDone())
                future.get();
        } catch (ExecutionException e) {
            t = e.getCause();
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }
    if (t != null) {
        System.out.println(t);
    }
}

12.2. Обработка исключений execute()

использоватьexecute()При обработке исключений необходимо учитывать четыре фактора:

1. По умолчанию он будет вexecute()Исключение генерируется непосредственно внутри метода. Обратите внимание, что это не прерывает работу пула потоков, но завершает текущий рабочий поток и воссоздает новый рабочий поток для выполнения задачи.

2. Разрешено вexecute()внутри методаtry-catchПреимущество перехвата этого исключения состоит в том, что текущий поток не завершается, а создается новый поток.

3. ПереписатьafterExecute()метод

4. Вы также можете установитьUncaughtExceptionHandler,Например:

ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(1, 2, 3, TimeUnit.SECONDS, new LinkedBlockingQueue(), 
    //我们自定义一个线程工厂和重写线程的setUncaughtExceptionHandler方法
    new ThreadFactory() {
        final AtomicInteger threadNumber = new AtomicInteger(1);
        public Thread newThread(Runnable r) {
            Thread thread = new Thread(Thread.currentThread().getThreadGroup(), r, "thread-"
                    + (threadNumber.getAndIncrement()));
            thread.setUncaughtExceptionHandler((t,e) -> System.out.println(e));
            return thread;
        }
    });