Пул потоков Java для понимания

Java

предисловие

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

Зачем использовать пулы потоков

1、减少线程创建与切换的开销

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

2、控制线程的数量

  • Используя пул потоков, мы можем эффективно контролировать количество потоков, когда в системе большое количество параллельных потоков, производительность системы резко падает.

что делает пул потоков

重复利用有限的线程

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

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

Фактически, общий пул потоков Java по существу состоит изThreadPoolExecutorилиForkJoinPoolГенерируется то, что он создает экземпляр соответствующего пула потоков в соответствии с различными аргументами, переданными конструктором.

Executors

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

  • newFixedThreadPool(): создать многоразовый пул потоков с фиксированным количеством потоков.
  • newSingleThreadExecutor(): создать пул потоков только с одним потоком.
  • newCachedThreadPool(): создать пул потоков с возможностями кэширования.
  • newWorkStealingPool(): создает пул потоков, который содержит достаточно потоков для поддержки заданного уровня параллелизма.
  • newScheduledThreadPool(): создать пул потоков с указанным количеством потоков, который может выполнять потоки задач после указания задержки.

Интерфейс ExecutorService

Студенты, которые разбираются в шаблонах проектирования, знают, что мы делаем все возможное, чтобы запрограммировать интерфейс, который очень дружелюбен к гибкости программы. Пул потоков Java также использует идею интерфейсно-ориентированного программирования, как вы можете видеть.ThreadPoolExecutorиForkJoinPoolвсеExecutorServiceКласс реализации интерфейса. существуетExecutorServiceНекоторые общие методы определяются в интерфейсе, после чего их можно использовать в различных пулах потоков.ExecutorServiceМетоды, определенные в интерфейсе, обычно используемые методы следующие:

  • 向线程池提交线程
    • Future<?> submit(): передать объект Runnable указанному пулу потоков, и пул потоков будет выполнять задачу, представленную объектом Runnable, когда есть незанятые потоки.Этот метод может получать как объекты Runnable, так и объекты Callable, что означает, что метод sumbit() может иметь возвращаемое значение.
    • void execute(Runnable command): может получать только объекты Runnable, что означает, что метод не имеет возвращаемого значения.
  • 关闭线程池
    • void shutdown(): Предотвратить отправку новых задач и не повлияет на уже отправленные задачи. (Ожидание завершения всех потоков перед закрытием)
    • List<Runnable> shutdownNow(): предотвращает отправку новых задач и прерывает текущий поток, а также удаляет задачи из workQueue и добавляет эти задачи в список для возврата. (сразу закрывается)
  • 检查线程池的状态
    • boolean isShutdown(): возвращает true после вызова метода shutdown() или shutdownNow().
    • boolean isTerminated(): при вызове метода shutdown() и выполнении всех отправленных задач он возвращает значение true; при вызове метода shutdownNow() он возвращает значение true после успешной остановки.

Общие примеры использования пула потоков

1. новый фиксированный пул потоков

Количество потоков в пуле потоков фиксировано, независимо от того, сколько у вас задач.

образец кода

public class MyFixThreadPool {

    public static void main(String[] args) throws InterruptedException {
        // 创建一个线程数固定为5的线程池
        ExecutorService service = Executors.newFixedThreadPool(5);

        System.out.println("初始线程池状态:" + service);

        for (int i = 0; i < 6; i++) {
            service.execute(() -> {
                try {
                    TimeUnit.MILLISECONDS.sleep(1000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                System.out.println(Thread.currentThread().getName());
            });
        }
        System.out.println("线程提交完毕之后线程池状态:" + service);

        service.shutdown();//会等待所有的线程执行完毕才关闭,shutdownNow:立马关闭
        System.out.println("是否全部线程已经执行完毕:" + service.isTerminated());//所有的任务执行完了,就会返回true
        System.out.println("是否已经执行shutdown()" + service.isShutdown());
        System.out.println("执行完shutdown()之后线程池的状态:" + service);

        TimeUnit.SECONDS.sleep(5);
        System.out.println("5秒钟过后,是否全部线程已经执行完毕:" + service.isTerminated());
        System.out.println("5秒钟过后,是否已经执行shutdown()" + service.isShutdown());
        System.out.println("5秒钟过后,线程池状态:" + service);
    }

}

результат операции:

Исходное состояние пула потоков: [Выполняется, размер пула = 0, активные потоки = 0, задачи в очереди = 0, завершенные задачи = 0]
Статус пула потоков после отправки потока: [Выполняется, размер пула = 5, активные потоки = 5, задачи в очереди = 1, завершенные задачи = 0]
Были ли выполнены все потоки: false
Был ли выполнен shutdown(): true
Состояние пула потоков после выполнения shutdown(): [Завершение работы, размер пула = 5, активные потоки = 5, задачи в очереди = 1, выполненные задачи = 0]
pool-1-thread-2
pool-1-thread-1
pool-1-thread-4
pool-1-thread-5
pool-1-thread-3
pool-1-thread-2
Через 5 секунд все ли потоки были выполнены: true
Через 5 секунд, было ли выполнено shutdown(): true
Через 5 секунд состояние пула потоков: [Завершено, размер пула = 0, активные потоки = 0, задачи в очереди = 0, выполненные задачи = 6]

анализ программы

  • Когда мы создаем FixedThreadPool, пул потоков находится вRunningстатус, ноpool size(количество потоков пула потоков),active threads(текущая активная тема)queued tasks(текущий поток в очереди),completed tasks(Количество выполненных задач) все 0
  • После того, как мы отправим все 6 задач в пул потоков,
    • pool size = 5: поскольку мы создали пул потоков с фиксированным числом из 5 потоков (примечание: если мы отправляем только 3 задачи в это время, тоpool size = 3, что указывает на то, что пул потоков также создает потоки посредством отложенной загрузки).
    • active threads = 5: Хотя мы отправили 6 задач в пул потоков, фиксированный размер пула потоков равен 5, поэтому активных потоков всего 5.
    • queued tasks = 1: Хотя мы отправили 6 задач в пул потоков, фиксированный размер пула потоков равен 5, и одновременно могут работать только 5 активных потоков, поэтому есть задача, ожидающая
  • Наше первое исполнениеshutdown(), потому что задачи не были полностью выполнены, поэтомуisTerminated()возвращениеfalse,shutdown()Возвращает true, и состояние пула потоков будет определятьсяRunningстатьShutting down
  • Из текущих результатов задачи мы видим, что имяpool-1-thread-2Задача выполняется дважды, что доказывает, что потоки в пуле потоков действительно используются повторно.
  • Через 5 секундisTerminated()возвращениеtrue,shutdown()возвращениеtrue, чтобы доказать, что все задачи выполнены и пул потоков закрыт, мы снова проверяем состояние пула потоков[Terminated, pool size = 0, active threads = 0, queued tasks = 0, completed tasks = 6], состояние уже вTerminated, то выполненные задачи отображаются как 6
2. новый однопоточный исполнитель

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

пример кода

public class SingleThreadPool {

    public static void main(String[] args) {
        ExecutorService service = Executors.newSingleThreadExecutor();

        for (int i = 0; i < 5; i++) {
            final int j = i;
            service.execute(() -> {
                System.out.println(j + " " + Thread.currentThread().getName());
            });
        }
    }

}

результат операции

0 pool-1-thread-1 1 pool-1-thread-1 2 pool-1-thread-1 3 pool-1-thread-1 4 pool-1-thread-1

анализ программымогу видеть толькоpool-1-thread-1Ветка работает.

Три, новыйCachedThreadPool

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

пример кода

public class CachePool {

    public static void main(String[] args) throws InterruptedException {
        ExecutorService service = Executors.newCachedThreadPool();
        System.out.println("初始线程池状态:" + service);

        for (int i = 0; i < 12; i++) {
            service.execute(() -> {
                try {
                    TimeUnit.MILLISECONDS.sleep(500);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                }
                System.out.println(Thread.currentThread().getName());
            });
        }
        System.out.println("线程提交完毕之后线程池状态:" + service);

        TimeUnit.SECONDS.sleep(50);
        System.out.println("50秒后线程池状态:" + service);

        TimeUnit.SECONDS.sleep(30);
        System.out.println("80秒后线程池状态:" + service);
    }

}

результат операции

Исходное состояние пула потоков: [Выполняется, размер пула = 0, активные потоки = 0, задачи в очереди = 0, завершенные задачи = 0]
Статус пула потоков после отправки потока: [Выполняется, размер пула = 12, активные потоки = 12, задачи в очереди = 0, завершенные задачи = 0]
pool-1-thread-3
pool-1-thread-4
pool-1-thread-1
pool-1-thread-2
pool-1-thread-5
pool-1-thread-8
pool-1-thread-9
pool-1-thread-12
pool-1-thread-7
pool-1-thread-6
pool-1-thread-11
pool-1-thread-10
Состояние пула потоков через 50 секунд: [Выполняется, размер пула = 12, активных потоков = 0, задач в очереди = 0, завершенных задач = 12]
Состояние пула потоков через 80 секунд: [Выполняется, размер пула = 0, активные потоки = 0, задачи в очереди = 0, выполненные задачи = 12]

анализ программы

  • Поскольку каждая задача потока требует не менее 500 миллисекунд времени выполнения, когда мы отправляем 12 задач в пул потоков, у нас в основном нет свободных потоков для повторного использования, поэтому пул потоков создаст 12 потоков.
  • По умолчанию потоки в кеше будут уничтожены, если они не будут активны в течение 60 секунд.Вы можете видеть, что через 50 секунд все задачи были выполнены, но количество потоков в пуле потоков по-прежнему равно 12.
  • Через 80 секунд вы увидите, что все потоки в пуле потоков уничтожены.
Четыре, новый запланированный пул потоков

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

ScheduledThreadPoolExecutor

  • newScheduledThreadPool()Метод возвращаетScheduledThreadPoolExecutorобъект,ScheduledThreadPoolExecutorОпределяется следующим образом:
public class ScheduledThreadPoolExecutor
        extends ThreadPoolExecutor
        implements ScheduledExecutorService {
  • Как видите, он по-прежнему наследуетThreadPoolExecutorи понялScheduledExecutorServiceинтерфейс, при этомScheduledExecutorServiceтакже унаследовалExecutorServiceинтерфейс, поэтому мы также можем использовать его, как и предыдущий объект пула потоков, но у объекта будут дополнительные методы для управления задержкой и циклом:
    • 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,ong initialDelay,long delay,TimeUnit unit): создает и выполняет периодическое действие, которое сначала активируется после заданной начальной задержки, а затем следует заданная задержка между завершением каждого выполнения и началом следующего.

образец кода

Код ниже печатает имя текущего потока вместе со случайным числом каждые 500 миллисекунд.

public class MyScheduledPool {

    public static void main(String[] args) {
        ScheduledExecutorService service = Executors.newScheduledThreadPool(4);
        service.scheduleAtFixedRate(() -> {
            System.out.println(Thread.currentThread().getName() + new Random().nextInt(1000));
        }, 0, 500, TimeUnit.MILLISECONDS);
    }
}
Пять, новыйWorkStealingPool

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

образец кода

public class MyWorkStealingPool {

    public static void main(String[] args) throws IOException {
        ExecutorService service = Executors.newWorkStealingPool(4);
        System.out.println("cpu核心:" + Runtime.getRuntime().availableProcessors());

        service.execute(new R(1000));
        service.execute(new R(2000));
        service.execute(new R(2000));
        service.execute(new R(2000));
        service.execute(new R(2000));

        //由于产生的是精灵线程(守护线程、后台线程),主线程不阻塞的话,看不到输出
        System.in.read();
    }

    static class R implements Runnable {

        int time;

        R(int time) {
            this.time = time;
        }

        @Override
        public void run() {
            try {
                TimeUnit.MILLISECONDS.sleep(time);
            } catch (InterruptedException e) {
                e.printStackTrace();
            }
            System.out.println(time + " " + Thread.currentThread().getName());
        }
    }
}

результат операции

ядра процессора: 4 1000 ForkJoinPool-1-рабочий-1 2000 ForkJoinPool-1-рабочий-0 2000 ForkJoinPool-1-рабочий-3 2000 ForkJoinPool-1-рабочий-2 2000 ForkJoinPool-1-рабочий-1

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

6. ФоркДжоинПул
  • ForkJoinPoolЗадачу можно разделить на несколько «маленьких задач» для параллельных вычислений, а результаты нескольких «маленьких задач» можно объединить в общий результат расчета.ForkJoinPoolПредусмотрены следующие способы созданияForkJoinPoolобъект экземпляра:

    • ForkJoinPool(int parallelism): создание параллелизма, содержащего параллельные потоки.ForkJoinPool, значение параллелизма по умолчанию равноRuntime.getRuntime().availableProcessors()возвращаемое значение метода
    • ForkJoinPool commonPool(): этот метод возвращает общий пул, на рабочее состояние общего пула не будут влиятьshutdown()илиshutdownNow()воздействие метода.
  • созданныйForkJoinPoolПосле примера вы можете позвонитьForkJoinPoolизsubmit(ForkJoinTask task)илиinvoke(ForkJoinTask task)способ выполнения указанной задачи. вForkJoinTask(реализует интерфейс Future) Представляет задачу, которую можно распараллелить и объединить.ForkJoinTaskявляется абстрактным классом и имеет два абстрактных подкласса:RecursiveActionиRecursiveTask. вRecursiveTaskпредставляет задачу с возвращаемым значением, аRecursiveActionПредставляет задачу без возвращаемого значения.

образец кода

Следующий код демонстрирует использованиеForkJoinPoolСумма 1000000 случайных целых чисел.

public class MyForkJoinPool {

    static int[] nums = new int[1000000];
    static final int MAX_NUM = 50000;
    static Random random = new Random();

    static {
        for (int i = 0; i < nums.length; i++) {
            nums[i] = random.nextInt(1000);
        }
        System.out.println(Arrays.stream(nums).sum());
    }

//    static class AddTask extends RecursiveAction {
//
//        int start, end;
//
//        AddTask(int start, int end) {
//            this.start = start;
//            this.end = end;
//        }
//
//        @Override
//        protected void compute() {
//            if (end - start <= MAX_NUM) {
//                long sum = 0L;
//                for (int i = 0; i < end; i++) sum += nums[i];
//                System.out.println("from:" + start + " to:" + end + " = " + sum);
//            } else {
//                int middle = start + (end - start) / 2;
//
//                AddTask subTask1 = new AddTask(start, middle);
//                AddTask subTask2 = new AddTask(middle, end);
//                subTask1.fork();
//                subTask2.fork();
//            }
//        }
//    }

    static class AddTask extends RecursiveTask<Long> {

        int start, end;

        AddTask(int start, int end) {
            this.start = start;
            this.end = end;
        }

        @Override
        protected Long compute() {
            // 当end与start之间的差大于MAX_NUM,将大任务分解成两个“小任务”
            if (end - start <= MAX_NUM) {
                long sum = 0L;
                for (int i = start; i < end; i++) sum += nums[i];
                return sum;
            } else {
                int middle = start + (end - start) / 2;

                AddTask subTask1 = new AddTask(start, middle);
                AddTask subTask2 = new AddTask(middle, end);
                // 并行执行两个“小任务”
                subTask1.fork();
                subTask2.fork();
                // 把两个“小任务”累加的结果合并起来
                return subTask1.join() + subTask2.join();
            }
        }
    }

    public static void main(String[] args) throws IOException {
        ForkJoinPool forkJoinPool = new ForkJoinPool();
        AddTask task = new AddTask(0, nums.length);
        forkJoinPool.execute(task);

        long result = task.join();
        System.out.println(result);

        forkJoinPool.shutdown();
    }
}

дополнительная добавка

Мы упоминали выше: на самом деле обычно используемый пул потоков Java состоит изThreadPoolExecutorилиForkJoinPoolДва класса генерируются, но они генерируют соответствующие пулы потоков в соответствии с разными аргументами, переданными конструктором. Затем давайте взглянем на исходный код, относящийся к нескольким статическим методам создания объектов пула потоков в Executors:

Прототип конструктора ThreadPoolExecutor

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

Параметр Описание

  • corePoolSize: poolSize основной операции, то есть когда он превышает этот диапазон, новый Runnable необходимо поместить в очередь ожидания workQueue.
  • maximumPoolSize: максимальное количество потоков, поддерживаемых пулом потоков. Когда значение больше этого значения, задача будет обрабатываться механизмом обработки отбрасывания (конечно, существуют также пулы потоков, которые никогда не отбрасывают задачи, в зависимости от стратегии). .
  • keepAliveTime: время выживания потока, когда он простаивает (этот параметр действителен, только когда количество потоков больше, чем corePoolSize) [java docЗаписывается так: когда количество потоков больше ядра, это максимальное время, в течение которого лишние простаивающие потоки будут ждать новых задач перед завершением.]
  • unit: единица измерения keepAliveTime.
  • workQueue: очередь блокировки, используемая для хранения задач, ожидающих выполнения, и задачи должны реализовывать интерфейс Runable.

процесс выполнения задач

  1. poolSize (фактическое количество потоков, используемых в настоящее время)
  2. Когда количество отправленных задач превышает corePoolSize, текущий Runnable будет отправлен в BlockingQueue.
  3. После заполнения ограниченной очереди, если poolSize
  4. Если третий шаг не может быть обработан, он перейдет к четвертому шагу для выполнения операции отклонения.

newFixedThreadPool

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

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

newSingleThreadExecutor

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

    public static ExecutorService newSingleThreadExecutor() {
        return new FinalizableDelegatedExecutorService
            (new ThreadPoolExecutor(1, 1,
                                    0L, TimeUnit.MILLISECONDS,
                                    new LinkedBlockingQueue<Runnable>()));
    }

newCachedThreadPool

poolSize равен 0, задача выбрасывается прямо в очередь, а SynchronousQueue используется для хранения (очередь без емкости), поэтому для задачи должен быть создан новый поток. как позволяющее создавать неограниченное количество потоков.

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

newScheduledThreadPool

    public ThreadPoolExecutor(int corePoolSize,
                              int maximumPoolSize,
                              long keepAliveTime,
                              TimeUnit unit,
                              BlockingQueue<Runnable> workQueue) {
        this(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue,
             Executors.defaultThreadFactory(), defaultHandler);
    }

newWorkStealingPool

    public static ExecutorService newWorkStealingPool(int parallelism) {
        return new ForkJoinPool
            (parallelism,
             ForkJoinPool.defaultForkJoinWorkerThreadFactory,
             null, true);
    }

агитационная сессия

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