Многопоточность Java — создание и использование пулов потоков и расширение исходного кода

Java

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

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

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

1 Что такое пул потоков

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

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

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

Все члены пула потоков находятся в пакете java.util.concurrent, который является ядром параллельного пакета JDK. Где ThreadPoolExecutor представляет пул потоков. Класс Executors — это роль фабрики потоков, через Executors можно получить пул потоков с определенной функцией, а через Executors — пул потоков с определенной функцией.

2.1 Метод newFixedThreadPool()

Этот метод возвращает пул потоков с фиксированным количеством потоков. Количество потоков в этом пуле потоков всегда одинаково. При отправке новой задачи, если в пуле потоков есть незанятые потоки, она будет выполнена немедленно. Если нет, новая задача будет временно храниться в очереди задач, а когда поток простаивает, будет обрабатываться очередь в очереди задач.

2.2 Метод newSingleThreadExecutor()

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

2.3 Метод newCachedThreadPool()

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

2.4 Метод newSingleThreadScheduledExecutor()

Этот метод возвращает объекты ScheduledExecutorService, размер пула потоков равен 1. Интерфейсы ScheduledExecutorService на интерфейс ExecutorService расширяют функциональные возможности выполнения задачи в заданное время, например, после фиксированной задержки выполнения или периодического выполнения задачи.

2.5 Метод newScheduledThreadPool()

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

Создайте пул потоков фиксированного размера

public class ThreadPoolThread {

    public static class MyTask implements Runnable{

        @Override
        public void run() {
            System.out.println(System.currentTimeMillis() + ":Thread ID: " + Thread.currentThread().getId());
            try{
                Thread.sleep(1000);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }

        public static void main(String[] args) {
            MyTask myTask = new MyTask();
            ExecutorService executorService = Executors.newFixedThreadPool(5);
            for (int i = 0;i< 10 ;i++){
                executorService.submit(myTask);
            }
        }
    }

}
1562554721820:Thread ID: 12
1562554721820:Thread ID: 15
1562554721820:Thread ID: 16
1562554721820:Thread ID: 13
1562554721820:Thread ID: 14
1562554722821:Thread ID: 15
1562554722821:Thread ID: 16
1562554722821:Thread ID: 12
1562554722821:Thread ID: 13
1562554722821:Thread ID: 14

Запланировать задачу

Метод newScheduledThreadPool() возвращает объект ScheduledExecutorService, который может планировать потоки в соответствии с требованиями времени. Основной метод заключается в следующем

public <V> ScheduledFuture<V> schedule(Callable<V> callable,long delay, TimeUnit unit);
public ScheduledFuture<?> scheduleAtFixedRate(Runnable command,long initialDelay,         long period,TimeUnit unit);
public ScheduledFuture<?> scheduleWithFixedDelay(Runnable command,long initialDelay,long delay,TimeUnit unit);

В отличие от других потоков, ScheduledExecutorService не обязательно планирует немедленное выполнение задач. Он фактически играет роль планирования задач, и будет планировать задачи в указанное время.

schedule() запланирует задачу один раз в заданное время. Методы scheduleAtFixedRate() и scheduleWithFixedDelay() будут периодически планировать задачи, но между ними все еще есть различия. Частота планирования задачи метода scheduleAtFixedRate() фиксирована, она начинается с момента начала выполнения предыдущей задачи, а затем планирует следующую задачу в указанное время. Метод scheduleWithFixedDelay() предназначен для планирования задачи через указанное время после окончания предыдущей задачи.

public class ScheduleExecutorServiceDemo {

    public static void main(String[] args) {
        ScheduledExecutorService scheduledExecutorService = Executors.newScheduledThreadPool(10);
        scheduledExecutorService.scheduleAtFixedRate(new Runnable() {
            @Override
            public void run() {
                try {
                    Thread.sleep(1000);
                    System.out.println(System.currentTimeMillis());
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
            }
        },0,2, TimeUnit.SECONDS);
    }
}
1562555518798
1562555520798
1562555522798
1562555524799
1562555526800

Видно, что задание запланировано каждые две секунды.

Если время выполнения задачи превышает время планирования, задача будет вызвана сразу после завершения предыдущей задачи.

Измените код на 8 секунд

Thread.sleep(8000);
1562555680333
1562555688333
1562555696333
1562555704333

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

3 Внутренняя реализация пула потоков

Для нескольких основных пулов потоков, несмотря на то, что они хотят создать пул потоков, они имеют разные функции, но внутри его используется класс ThreadPoolExecutor.

Посмотрите исходный код для создания нескольких пулов потоков:

public static ExecutorService newFixedThreadPool(int nThreads) {
        return new ThreadPoolExecutor(nThreads, nThreads,0L, TimeUnit.MILLISECONDS,new LinkedBlockingQueue<Runnable>());
    }
    
public static ExecutorService newCachedThreadPool() {
        return new ThreadPoolExecutor(0, Integer.MAX_VALUE,60L, TimeUnit.SECONDS,new SynchronousQueue<Runnable>());
    }

Видно, что все они являются инкапсуляциями класса ThreadPoolExecutor Давайте посмотрим на конструкцию класса ThreadPoolExecutor:

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;
    }
  • corePoolSize: указывает количество потоков в пуле потоков.
  • maxPoolSize: указывает максимальное количество потоков в пуле потоков.
  • keepAliveTime: Когда количество потоков в пуле потоков превышает corePoolSize, время выживания избыточных бездействующих потоков, то есть, как долго простаивающие потоки, превышающие corePoolSize, будут уничтожены.
  • unit: единица измерения keepAliveTime.
  • workQueue: очередь задач, очередь хранения для задач, которые были отправлены, но еще не выполнены.
  • threadFacotry: фабрика потоков, используемая для создания потоков.
  • обработчик: запретить политику. Как отклонять задачи, когда их слишком много.

3.1 workQueue — очередь задач

Параметр workQueue относится к очереди задач, которая отправлена, но не выполнена, является объектом интерфейса BlockingQueue и используется только для хранения объектов Runnable. В конструкции ThreadPoolExecutor можно использовать следующие интерфейсы BlockingQueue:

  • Отправить очередь напрямую: Предоставляется объектом SynchronousQueue. SynchronousQueue не имеет пропускной способности, каждая операция вставки должна ожидать соответствующей операции удаления, в противном случае каждая операция удаления должна ждать соответствующей операции вставки. Используя SynchronousQueue, если всегда есть новые задачи, отправленные в поток для выполнения, если нет незанятого процесса, он попытается создать новый поток, и если количество потоков достигло максимального значения, будет выполнена политика отклонения . При использовании SynchronousQueue обычно устанавливается большое значение maxPoolSize, в противном случае легко применить политику отклонения.
  • Ограниченная очередь задач: ограниченная очередь задач реализована с использованием класса ArrayBlockingQueue.Конструктор класса ArrayBlockingQueue должен принимать параметр емкости, указывающий максимальную емкость очереди. При использовании ограниченной очереди задач, если есть новая задача для выполнения, когда фактический поток пула потоков меньше, чем corePoolSize, поток будет создан первым, а если он больше, чем corePoolSize, задача будет добавить в очередь ожидания. Если очередь ожидания заполнена, создайте новый поток для выполнения задачи, если общий поток не превышает maxPoolSize, и выполните политику отклонения, если он больше maxPoolSize.
  • Неограниченная очередь задач. Неограниченная очередь задач реализована с помощью класса LinkedBlockingQueue. По сравнению с ограниченными очередями задач, неограниченные очереди задач всегда ставятся в очередь. При использовании LinkedBlockingQueue, когда новые задачи должны выполняться потоками, если количество потоков меньше, чем corePoolSize, будут созданы новые потоки, но когда количество потоков достигнет corePoolSize, оно не будет продолжать расти. Если в будущем будет добавлена ​​новая задача, если есть какой-либо незанятый поток, он напрямую войдет в очередь и будет ждать.
  • Очередь задач с приоритетом: очередь с приоритетом выполнения при использовании очереди задач с приоритетом. Реализовано через класс PriorityBlockingQueue, вы можете контролировать порядок выполнения задач. Это специальная неограниченная очередь. Класс PriorityBlockingQueue может выполнять задачи в соответствии с их собственным приоритетом.

4 Политика отказа

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

В JDK встроены четыре политики отказа:

  • Политика AbortPolicy: эта политика напрямую вызывает исключение, предотвращая нормальную работу.
  • Политика CallerRunsPolicy: Пока пул потоков не закрыт, эта политика запускает текущую отброшенную задачу непосредственно в текущем потоке вызывающего объекта. На самом деле это не приведет к падению потока, но снизит производительность потока отправки задачи.
  • Политика DiscardOldestPolicy: эта политика отбрасывает самый старый запрос, то есть задачу, которая должна быть выполнена, и пытается снова отправить текущую задачу.
  • Политика DiscardPolicy: эта политика отбрасывает задачи, которые не могут быть обработаны без какой-либо обработки.

Все вышеперечисленные стратегии реализуют интерфейс RejectedExecutionHandler.Если приведенные выше стратегии не соответствуют фактической разработке, вы можете расширить их самостоятельно.

Конструкция интерфейса RejectedExecutionHandler:

public interface RejectedExecutionHandler {
    void rejectedExecution(Runnable r, ThreadPoolExecutor executor);
}

Пользовательская политика отказа:

//拒绝策略demo
public class RejectThreadPoolDemo {

    public static class MyTask implements Runnable{

        @Override
        public void run() {
            System.out.println(System.currentTimeMillis() + ": Thread ID : " + Thread.currentThread().getId());
            try{
                Thread.sleep(100);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
        }
    }

    public static void main(String[] args) throws InterruptedException {
        MyTask myTask = new MyTask();
        ThreadPoolExecutor es = new ThreadPoolExecutor(5, 5, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>(10), Executors.privilegedThreadFactory(), new RejectedExecutionHandler() {
            @Override
            public void rejectedExecution(Runnable r, ThreadPoolExecutor executor) {
                System.out.println(r.toString() + "被拒绝");
            }
        });
        for (int i = 0;i<Integer.MAX_VALUE;i++){
           es.submit(myTask);
           Thread.sleep(10);
        }
    }

}
1562575292467: Thread ID : 14
1562575292478: Thread ID : 15
1562575292489: Thread ID : 16
java.util.concurrent.FutureTask@b4c966a被拒绝
java.util.concurrent.FutureTask@2f4d3709被拒绝
java.util.concurrent.FutureTask@4e50df2e被拒绝

5 Создание пользовательского потока: ThreadFactory

ThreadFactory — это интерфейс, который имеет только один метод для создания потоков.

Thread newThread(Runnable r);

С созданием пользовательских потоков мы можем отслеживать, сколько потоков создается пулом потоков в какое время, имена пользовательских потоков и т. д.

public class ThreadFactoryDemo {

    static volatile int i = 0;

    public static class TestTask implements Runnable{

        @Override
        public void run() {
            System.out.println(Thread.currentThread().getName());
        }
    }

    public static void main(String[] args) {
        TestTask testTask = new TestTask();
        ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(5, 5, 0L, TimeUnit.MILLISECONDS, new SynchronousQueue<>(), new ThreadFactory() {
            @Override
            public Thread newThread(Runnable r) {
                Thread thread = new Thread(r,"test--" + i);
                i++;
                return thread;
            }
        });
        for (int i = 0;i<5;i++){
            threadPoolExecutor.submit(testTask);
        }
    }

}
test--0
test--1
test--4
test--2
test--3

6 Расширение пула потоков

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

ThreadPoolExecutor — это расширяемый пул потоков, который предоставляет три интерфейса beforeExecutor(), afterExecutor() и terminated() для его расширения.

public class ThreadFactoryDemo {

    static volatile int i = 0;

    public static class TestTask implements Runnable{

        @Override
        public void run() {
            System.out.println(Thread.currentThread().getName());
        }
    }

    public static void main(String[] args) {
        TestTask testTask = new TestTask();
        ThreadPoolExecutor threadPoolExecutor = new ThreadPoolExecutor(5, 5, 0L, TimeUnit.MILLISECONDS, new SynchronousQueue<>(), new ThreadFactory() {
            @Override
            public Thread newThread(Runnable r) {
                Thread thread = new Thread(r,"test--" + i);
                i++;
                return thread;
            }
        }){
            @Override
            protected void beforeExecute(Thread t,Runnable r){
                System.out.println("task-----准备执行");
            }
        };
        for (int i = 0;i<5;i++){
            threadPoolExecutor.submit(testTask);
        }
    }

}
task-----准备执行
task-----准备执行
test--2
task-----准备执行
test--1
task-----准备执行
test--4
task-----准备执行
test--3
test--0

7 Разница между отправкой и выполнением

7.1 Метод выполнения()

Метод отправки execute может отправлять только объект Runnable, а возвращаемое значение метода недействительно, то есть если поток запускается после отправки, он будет отделен от основного потока.Конечно, вы можете установить некоторые переменные для получения результат операции потока. И когда во время выполнения потока генерируется исключение, основной поток обычно не может получить информацию об исключении.Только путем активной установки класса обработки исключений потока через ThreadFactory он может воспринять исключение в отправленном потоке.

7.2 метод sumbit()

Метод submit() имеет три формы:

public Future<?> submit(Runnable task) {
        if (task == null) throw new NullPointerException();
        RunnableFuture<Void> ftask = newTaskFor(task, null);
        execute(ftask);
        return ftask;
    }

    /**
     * @throws RejectedExecutionException {@inheritDoc}
     * @throws NullPointerException       {@inheritDoc}
     */
    public <T> Future<T> submit(Runnable task, T result) {
        if (task == null) throw new NullPointerException();
        RunnableFuture<T> ftask = newTaskFor(task, result);
        execute(ftask);
        return ftask;
    }

    /**
     * @throws RejectedExecutionException {@inheritDoc}
     * @throws NullPointerException       {@inheritDoc}
     */
    public <T> Future<T> submit(Callable<T> task) {
        if (task == null) throw new NullPointerException();
        RunnableFuture<T> ftask = newTaskFor(task);
        execute(ftask);
        return ftask;
    }

Метод sumbit возвращает объект Future, который представляет результат выполнения этого потока.Когда основной поток вызывает метод get Future, он получает данные результата, возвращенные потоком. Если во время выполнения потока возникает исключение, get получит информацию об исключении.