Spec-Zone.ru › OpenJDK 24

Класс ForkJoinPool

java.lang.Object
java.util.concurrent.AbstractExecutorService
java.util.concurrent.ForkJoinPool
Все реализованные интерфейсы:
AutoCloseable, Executor, ExecutorService
public class ForkJoinPool extends AbstractExecutorService
Интерфейс ExecutorService для запуска задач ForkJoinTask. Класс ForkJoinPool предоставляет точку входа для задач, отправленных клиентами, не являющимися частями ForkJoinTask вычислений, а также операции управления и мониторинга.

Класс ForkJoinPool отличается от других видов ExecutorService главным образом тем, что использует кражу работы: все потоки в пуле пытаются найти и выполнить задачи, отправленные в пул и/или созданные другими активными задачами (в конечном итоге блокируются, ожидая работы, если таковой нет). Это обеспечивает эффективную обработку, когда большинство задач порождают другие подзадачи (как и большинство ForkJoinTask), а также когда множество мелких задач отправляются в пул внешними клиентами. Особенно при установке asyncMode в значение true в конструкторах, ForkJoinPool также могут быть подходящими для использования с задачами в стиле событий, которые никогда не присоединяются. Все рабочие потоки инициализируются с Thread.isDaemon(), установленным в значение true.

Доступен статический commonPool(), подходящий для большинства приложений. Общий пул используется любой задачей ForkJoinTask, которая не отправлена явно в указанный пул. Использование общего пула обычно снижает использование ресурсов (его потоки медленно освобождаются в периоды бездействия и возобновляются при последующем использовании).

Для приложений, которым требуются отдельные или пользовательские пулы, может быть создан ForkJoinPool с заданным уровнем параллелизма; по умолчанию он равен количеству доступных процессоров. Пул пытается поддерживать достаточное количество активных (или доступных) потоков, динамически добавляя, приостанавливая или возобновляя внутренние рабочие потоки, даже если некоторые задачи приостановлены в ожидании присоединения к другим. Однако такие корректировки не гарантируются при наличии заблокированного ввода-вывода или другой неконтролируемой синхронизации. Встроенный интерфейс ForkJoinPool.ManagedBlocker позволяет расширить типы синхронизации, которые поддерживаются. По умолчанию политики могут быть переопределены с помощью конструктора с параметрами, соответствующими документации класса ThreadPoolExecutor.

В дополнение к методам управления выполнением и жизненным циклом, этот класс предоставляет методы проверки состояния (например, getStealCount()), предназначенные для помощи в разработке, настройке и мониторинге приложений fork/join. Кроме того, метод toString() возвращает сведения о состоянии пула в удобной форме для неформального мониторинга.

Как и в случае с другими ExecutorService, существует три основных метода выполнения задач, которые подытожены в следующей таблице. Они предназначены для использования в первую очередь клиентами, которые не участвуют в вычислениях fork/join в текущем пуле. Основные формы этих методов принимают экземпляры ForkJoinTask, но перегруженные формы также позволяют смешанное выполнение обычных операций Runnable или Callable. Однако задачи, которые уже выполняются в пуле, обычно вместо этого должны использовать формы внутри вычислений, перечисленные в таблице, если не используется задача в стиле событий async, которая обычно не присоединяется, в этом случае различий между выбором методов мало.

Сводка методов выполнения задач
Вызов от клиентов, не участвующих в fork/join Вызов изнутри вычислений fork/join
Организация асинхронного выполнения execute(ForkJoinTask) ForkJoinTask.fork()
Ожидание и получение результата invoke(ForkJoinTask) ForkJoinTask.invoke()
Организация exec и получение Future submit(ForkJoinTask) ForkJoinTask.fork() (ForkJoinTasks являются Future)

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

  • java.util.concurrent.ForkJoinPool.common.parallelism - уровень параллелизма, целое неотрицательное число
  • java.util.concurrent.ForkJoinPool.common.threadFactory - имя класса ForkJoinPool.ForkJoinWorkerThreadFactory. Для загрузки этого класса используется системный загрузчик классов.
  • java.util.concurrent.ForkJoinPool.common.exceptionHandler - имя класса Thread.UncaughtExceptionHandler. Для загрузки этого класса используется системный загрузчик классов.
  • java.util.concurrent.ForkJoinPool.common.maximumSpares - максимальное количество дополнительных потоков для поддержания целевого уровня параллелизма (по умолчанию 256).
Если фабрика потоков не указана через системное свойство, общий пул использует фабрику, которая использует системный загрузчик классов как
загрузчик класса контекста потока. При любой ошибке при настройке этих параметров используются значения по умолчанию. Можно отключить или ограничить использование потоков в общем пуле, установив свойство параллелизма в ноль и/или используя фабрику, которая может возвращать null. Однако это может привести к тому, что задачи, не объединенные с результатом, никогда не будут выполнены.
Примечание об реализации:
Эта реализация ограничивает максимальное количество запущенных потоков 32767. Попытки создать пулы с количеством потоков, превышающим максимальное значение, приводят к IllegalArgumentException. Кроме того, эта реализация отклоняет отправленные задачи (то есть, выбрасывая RejectedExecutionException) только при завершении пула или исчерпании внутренних ресурсов.
С:
1.7

Краткое описание вложенных классов

Модификатор и тип Класс Описание
static interface  ForkJoinPool.ForkJoinWorkerThreadFactory
Фабрика для создания новых ForkJoinWorkerThread.
static interface  ForkJoinPool.ManagedBlocker
Интерфейс для расширения управляемой параллельности для задач, выполняемых в ForkJoinPool.

Краткое описание полей

Модификатор и тип Поле Описание
static final ForkJoinPool.ForkJoinWorkerThreadFactory defaultForkJoinWorkerThreadFactory
Создает новый ForkJoinWorkerThread.

Краткое описание конструкторов

Конструктор Описание
ForkJoinPool()
Создает ForkJoinPool с параллелизмом, равным Runtime.availableProcessors(), используя значения по умолчанию для всех остальных параметров (см. ForkJoinPool(int, ForkJoinWorkerThreadFactory, UncaughtExceptionHandler, boolean, int, int, int, Predicate, long, TimeUnit)).
ForkJoinPool(int parallelism)
Создает ForkJoinPool с указанным уровнем параллелизма, используя значения по умолчанию для всех остальных параметров (см. ForkJoinPool(int, ForkJoinWorkerThreadFactory, UncaughtExceptionHandler, boolean, int, int, int, Predicate, long, TimeUnit)).
ForkJoinPool(int parallelism, ForkJoinPool.ForkJoinWorkerThreadFactory factory, Thread.UncaughtExceptionHandler handler, boolean asyncMode)
Создает ForkJoinPool с заданными параметрами (используя значения по умолчанию для других — см. ForkJoinPool(int, ForkJoinWorkerThreadFactory, UncaughtExceptionHandler, boolean, int, int, int, Predicate, long, TimeUnit)).
ForkJoinPool(int parallelism, ForkJoinPool.ForkJoinWorkerThreadFactory factory, Thread.UncaughtExceptionHandler handler, boolean asyncMode, int corePoolSize, int maximumPoolSize, int minimumRunnable, Predicate<? super ForkJoinPool> saturate, long keepAliveTime, TimeUnit unit)
Создает ForkJoinPool с заданными параметрами.

Краткое описание методов

Модификатор и тип Метод Описание
boolean awaitQuiescence(long timeout, TimeUnit unit)
Если вызывается задачей ForkJoinTask, работающей в этом пуле, эквивалентно ForkJoinTask.helpQuiesce().
boolean awaitTermination(long timeout, TimeUnit unit)
Ожидает завершения всех задач после запроса остановки, или истечения таймаута, или прерывания текущего потока — в зависимости от того, что произойдёт раньше.
void close()
Запускает упорядоченную остановку (если это не commonPool()), в которой выполняются ранее отправленные задачи, но новые задачи не принимаются, и ожидает завершения всех задач и завершения работы исполнителя.
static ForkJoinPool commonPool()
Возвращает экземпляр общего пула.
protected int drainTasksTo(Collection<? super ForkJoinTask<?>> c)
Удаляет все доступные невыполненные отправленные и разветвлённые задачи из очередей планирования и добавляет их в заданный коллекцию, не изменяя их статус выполнения.
void execute(Runnable task)
Выполняет заданный метод в будущем.
void execute(ForkJoinTask<?> task)
Организует (асинхронное) выполнение заданной задачи.
<T> ForkJoinTask<T> externalSubmit(ForkJoinTask<T> task)
Отправляет заданную задачу так, как будто она была отправлена не-ForkJoinTask клиентом.
int getActiveThreadCount()
Возвращает оценку количества потоков, которые в данный момент крадут или выполняют задачи.
boolean getAsyncMode()
Возвращает true, если этот пул использует локальный режим планирования FIFO для разветвлённых задач, которые никогда не объединяются.
static int getCommonPoolParallelism()
Возвращает целевой уровень параллелизма общего пула.
ForkJoinPool.ForkJoinWorkerThreadFactory getFactory()
Возвращает фабрику, используемую для создания новых рабочих потоков.
int getParallelism()
Возвращает целевой уровень параллелизма этого пула.
int getPoolSize()
Возвращает количество рабочих потоков, которые были запущены, но ещё не завершены.
int getQueuedSubmissionCount()
Возвращает оценку количества задач, отправленных в этот пул, которые ещё не начали выполняться.
long getQueuedTaskCount()
Возвращает оценку общего количества задач, которые в данный момент находятся в очередях рабочих потоков (но не включает задачи, отправленные в пул, которые ещё не начали выполняться).
int getRunningThreadCount()
Возвращает оценку количества рабочих потоков, которые не заблокированы в ожидании объединения задач или других управляемых синхронизаций.
long getStealCount()
Возвращает оценку общего количества завершённых задач, которые были выполнены потоком, отличным от потока, их отправившего.
Thread.UncaughtExceptionHandler getUncaughtExceptionHandler()
Возвращает обработчик для внутренних рабочих потоков, которые завершаются из-за непреодолимых ошибок, возникших при выполнении задач.
boolean hasQueuedSubmissions()
Возвращает true, если в этот пул отправлены какие-либо задачи, которые ещё не начали выполняться.
<T> T invoke(ForkJoinTask<T> task)
Выполняет заданную задачу и возвращает её результат по завершении.
<T> List<Future<T>> invokeAllUninterruptibly(Collection<? extends Callable<T>> tasks)
Непрерывная версия invokeAll.
boolean isQuiescent()
Возвращает true, если все рабочие потоки в данный момент бездействуют.
boolean isShutdown()
Возвращает true, если этот пул был остановлен.
boolean isTerminated()
Возвращает true, если все задачи были завершены после остановки.
boolean isTerminating()
Возвращает true, если процесс завершения начался, но еще не завершен.
<T> ForkJoinTask<T> lazySubmit(ForkJoinTask<T> task)
Отправляет задачу без гарантии, что она будет в конечном итоге выполнена в отсутствие доступных активных потоков.
static void managedBlock(ForkJoinPool.ManagedBlocker blocker)
Выполняет заданную, возможно, блокирующую задачу.
protected ForkJoinTask<?> pollSubmission()
Удаляет и возвращает следующую невыполненную отправку, если она доступна.
int setParallelism(int size)
Изменяет целевую параллельность этого пула, управляя будущим созданием, использованием и завершением потоков-работников.
void shutdown()
Возможно, инициирует упорядоченное завершение, в котором ранее отправленные задачи выполняются, но новые задачи не будут приниматься.
List<Runnable> shutdownNow()
Возможно, пытается отменить и/или остановить все задачи и отклонить все последующие отправленные задачи.
ForkJoinTask<?> submit(Runnable task)
Отправляет задачу Runnable для выполнения и возвращает Future, представляющую эту задачу.
<T> ForkJoinTask<T> submit(Runnable task, T result)
Отправляет задачу Runnable для выполнения и возвращает Future, представляющую эту задачу.
<T> ForkJoinTask<T> submit(Callable<T> task)
Отправляет возвращающую значение задачу для выполнения и возвращает Future, представляющую ожидаемые результаты задачи.
<T> ForkJoinTask<T> submit(ForkJoinTask<T> task)
Отправляет ForkJoinTask для выполнения.
String toString()
Возвращает строку, идентифицирующую этот пул, а также его состояние, включая указания на состояние выполнения, уровень параллельности и количество потоков-работников и задач.

Methods declared in class java.util.concurrent.AbstractExecutorService

invokeAll, invokeAll, invokeAny, invokeAny, newTaskFor, newTaskFor

Methods declared in class java.lang.Object

clone, equals, finalize, getClass, hashCode, notify, notifyAll, wait, wait, wait

Подробное описание полей

defaultForkJoinWorkerThreadFactory

public static final ForkJoinPool.ForkJoinWorkerThreadFactory defaultForkJoinWorkerThreadFactory
Создаёт новый ForkJoinWorkerThread. Этот фабричный метод используется, если не переопределён в конструкторах ForkJoinPool.

Подробное описание конструкторов

ForkJoinPool

public ForkJoinPool()
Создаёт ForkJoinPool с параллелизмом, равным Runtime.availableProcessors(), используя значения по умолчанию для всех остальных параметров (см. ForkJoinPool(int, ForkJoinWorkerThreadFactory, UncaughtExceptionHandler, boolean, int, int, int, Predicate, long, TimeUnit)).

ForkJoinPool

public ForkJoinPool(int parallelism)
Создаёт ForkJoinPool с указанным уровнем параллелизма, используя значения по умолчанию для всех остальных параметров (см. ForkJoinPool(int, ForkJoinWorkerThreadFactory, UncaughtExceptionHandler, boolean, int, int, int, Predicate, long, TimeUnit)).
Параметры:
parallelism - уровень параллелизма
Исключения:
IllegalArgumentException - если уровень параллелизма меньше или равен нулю, или больше предела реализации

ForkJoinPool

public ForkJoinPool(int parallelism, ForkJoinPool.ForkJoinWorkerThreadFactory factory, Thread.UncaughtExceptionHandler handler, boolean asyncMode)
Создаёт ForkJoinPool с заданными параметрами (используя значения по умолчанию для других — см. ForkJoinPool(int, ForkJoinWorkerThreadFactory, UncaughtExceptionHandler, boolean, int, int, int, Predicate, long, TimeUnit)).
Параметры:
parallelism - уровень параллелизма. Для значения по умолчанию используйте Runtime.availableProcessors().
factory - фабрика для создания новых потоков. Для значения по умолчанию используйте defaultForkJoinWorkerThreadFactory.
handler - обработчик для внутренних рабочих потоков, которые завершаются из-за непреодолимых ошибок, возникающих при выполнении задач. Для значения по умолчанию используйте null.
asyncMode - если true, устанавливает локальный режим планирования FIFO для разветвлённых задач, которые никогда не объединяются. Этот режим может быть более подходящим, чем режим по умолчанию на основе локального стека в приложениях, в которых рабочие потоки обрабатывают только асинхронные задачи типа событий. Для значения по умолчанию используйте false.
Исключения:
IllegalArgumentException - если уровень параллелизма меньше или равен нулю, или больше предела реализации
NullPointerException - если фабрика равна null

ForkJoinPool

public ForkJoinPool(int parallelism, ForkJoinPool.ForkJoinWorkerThreadFactory factory, Thread.UncaughtExceptionHandler handler, boolean asyncMode, int corePoolSize, int maximumPoolSize, int minimumRunnable, Predicate<? super ForkJoinPool> saturate, long keepAliveTime, TimeUnit unit)
Создаёт ForkJoinPool с заданными параметрами.
Параметры:
parallelism — уровень параллелизма. Для значения по умолчанию используйте Runtime.availableProcessors().
factory — фабрика для создания новых потоков. Для значения по умолчанию используйте defaultForkJoinWorkerThreadFactory.
handler — обработчик для внутренних рабочих потоков, которые завершаются из-за невосстановимых ошибок, возникших во время выполнения задач. Для значения по умолчанию используйте null.
asyncMode — если true, устанавливает локальный режим планирования «первый пришёл — первый обслужен» для разветвлённых задач, которые никогда не присоединяются. Этот режим может быть более подходящим, чем режим по умолчанию, основанный на локальном стеке, в приложениях, в которых рабочие потоки обрабатывают только асинхронные задачи типа событий. Для значения по умолчанию используйте false.
corePoolSize — количество потоков, которые необходимо поддерживать в пуле (если они не отключаются после истечения времени ожидания). Обычно (и по умолчанию) это значение такое же, как и уровень параллелизма, но может быть установлено большим значением, чтобы уменьшить динамическую нагрузку, если задачи регулярно блокируются. Использование меньшего значения (например, 0) имеет тот же эффект, что и по умолчанию.
maximumPoolSize — максимальное количество разрешённых потоков. Когда максимум достигнут, попытки заменить заблокированные потоки терпят неудачу. (Однако, поскольку создание и завершение различных потоков могут перекрываться и могут управляться заданной фабрикой потоков, это значение может временно превышаться.) Чтобы установить такое же значение, как используется по умолчанию для общего пула, используйте 256 плюс parallelism уровень. (По умолчанию общий пул допускает максимум 256 резервных потоков.) Использование значения (например, Integer.MAX_VALUE), большего, чем общий предел потоков реализации, имеет тот же эффект, что и использование этого предела (который является значением по умолчанию).
minimumRunnable — минимальное разрешённое число базовых потоков, не заблокированных объединением или ForkJoinPool.ManagedBlocker. Для обеспечения прогресса, когда существует слишком мало разблокированных потоков и могут существовать невыполненные задачи, создаются новые потоки до заданного максимального размера пула. Для значения по умолчанию используйте 1, что гарантирует жизнеспособность. Более высокое значение может улучшить пропускную способность при наличии заблокированных действий, но может и не сделать этого из-за увеличенной нагрузки. Значение ноль может быть приемлемым, когда отправленные задачи не могут иметь зависимости, требующие дополнительных потоков.
saturate — если не null, предикат, вызываемый при попытках создать больше, чем максимальное разрешённое количество потоков. По умолчанию, когда поток собирается заблокироваться на присоединении или ForkJoinPool.ManagedBlocker, но не может быть заменён, потому что maximumPoolSize будет превышен, выбрасывается RejectedExecutionException. Но если этот предикат возвращает true, исключение не выбрасывается, поэтому пул продолжает работать с меньшим, чем целевое, числом работающих потоков, что может не гарантировать прогресс.
keepAliveTime — время с момента последнего использования, прежде чем поток завершается (а затем позже заменяется при необходимости). Для значения по умолчанию используйте 60, TimeUnit.SECONDS.
unit — единица измерения времени для аргумента keepAliveTime
Исключения:
IllegalArgumentException — если уровень параллелизма меньше или равен нулю, или превышает предел реализации, или если maximumPoolSize меньше parallelism, или если keepAliveTime меньше или равно нулю.
NullPointerException — если factory равно null
С:
9

Подробное описание методов

commonPool

public static ForkJoinPool commonPool()
Возвращает экземпляр общего пула. Этот пул статически построен; его состояние выполнения не затрагивается попытками shutdown() или shutdownNow(). Однако этот пул и любое текущее выполнение автоматически завершаются при завершении программы System.exit(int). Любая программа, которая полагается на завершение асинхронной обработки задач до завершения программы, должна вызвать commonPool().awaitQuiescence перед выходом.
Возвращает:
экземпляр общего пула
С:
1.8
END_OF_DOCUMENT_MARKER

invoke

public <T> T invoke(ForkJoinTask<T> task)
Выполняет заданную задачу, возвращая ее результат по завершении. Если вычисление столкнется с необработанным исключением или ошибкой, оно перебрасывается как результат этого вызова. Переброшенные исключения ведут себя так же, как обычные исключения, но, по возможности, содержат трассировки стека (например, отображаемые с помощью ex.printStackTrace()) как текущей, так и потока, фактически столкнувшегося с исключением; в минимальном случае только последнего.
Параметры типа:
T - тип результата задачи
Параметры:
task - задача
Возвращает:
результат задачи
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения

execute

public void execute(ForkJoinTask<?> task)
Организует (асинхронное) выполнение заданной задачи.
Параметры:
task - задача
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения

execute

public void execute(Runnable task)
Описание скопировано из интерфейса: Executor
Выполняет данный командный объект в какой-то момент в будущем. Команда может быть выполнена в новом потоке, в пуле потоков или в вызывающем потоке, по усмотрению реализации Executor.
Параметры:
task - задача Runnable
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения

submit

public <T> ForkJoinTask<T> submit(ForkJoinTask<T> task)
Отправляет ForkJoinTask на выполнение.
Требования к реализации:
Этот метод эквивалентен externalSubmit(ForkJoinTask), когда вызывается из потока, который не находится в этом пуле.
Параметры типа:
T - тип результата задачи
Параметры:
task - задача для отправки
Возвращает:
задачу
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения

submit

public <T> ForkJoinTask<T> submit(Callable<T> task)
Описание скопировано из интерфейса: ExecutorService
Отправляет задачу, возвращающую значение, для выполнения и возвращает Future, представляющую ожидаемые результаты задачи. Метод get Future вернет результат задачи при успешном завершении.

Если вы хотите немедленно заблокироваться, ожидая задачу, вы можете использовать конструкции вида result = exec.submit(aCallable).get();

Примечание: Класс Executors включает набор методов, которые могут преобразовывать некоторые другие распространенные объекты, похожие на замыкания, например, PrivilegedAction в форму Callable, чтобы они могли быть отправлены.

Задано:
submit в интерфейсе ExecutorService
Переопределяет:
submit в классе AbstractExecutorService
Параметры типа:
T - тип результата задачи
Параметры:
task - задача для отправки
Возвращает:
Future, представляющая ожидаемое завершение задачи
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения

submit

public <T> ForkJoinTask<T> submit(Runnable task, T result)
Описание скопировано из интерфейса: ExecutorService
Отправляет задачу Runnable для выполнения и возвращает Future, представляющий эту задачу. Метод Future's get вернёт заданный результат при успешном завершении.
Определено в:
submit в интерфейсе ExecutorService
Переопределено в:
submit в классе AbstractExecutorService
Параметры типа:
T - тип результата
Параметры:
task - задача для отправки
result - результат для возврата
Возвращает:
Future, представляющий ожидаемое завершение задачи
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения

submit

public ForkJoinTask<?> submit(Runnable task)
Описание скопировано из интерфейса: ExecutorService
Отправляет задачу Runnable для выполнения и возвращает Future, представляющий эту задачу. Метод Future's get вернёт null при успешном завершении.
Определено в:
submit в интерфейсе ExecutorService
Переопределено в:
submit в классе AbstractExecutorService
Параметры:
task - задача для отправки
Возвращает:
Future, представляющий ожидаемое завершение задачи
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения

externalSubmit

public <T> ForkJoinTask<T> externalSubmit(ForkJoinTask<T> task)
Отправляет задачу, как если бы она была отправлена из клиента, не являющегося ForkJoinTask. Задача добавляется в очередь планирования для отправки в пул, даже если вызывается из потока в этом пуле.
Требования к реализации:
Этот метод эквивалентен submit(ForkJoinTask), если вызывается из потока, который не принадлежит этому пулу.
Параметры типа:
T - тип результата задачи
Параметры:
task - задача для отправки
Возвращает:
задачу
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения
С:
20

lazySubmit

public <T> ForkJoinTask<T> lazySubmit(ForkJoinTask<T> task)
Отправляет задачу без гарантии, что она будет в конечном итоге выполнена в отсутствие активных потоков. В некоторых контекстах этот метод может снизить конкуренцию и накладные расходы, опираясь на контекстно-зависимые знания о том, что существующие потоки (включая вызывающий поток, если он работает в этом пуле) в конечном итоге станут доступны для выполнения задачи.
Параметры типа:
T - тип результата задачи
Параметры:
task - задача
Возвращает:
задачу
Исключения:
NullPointerException - если задача равна null
RejectedExecutionException - если задача не может быть запланирована для выполнения
С:
19

setParallelism

public int setParallelism(int size)
Изменяет целевую параллельность этого пула, контролируя будущее создание, использование и завершение рабочих потоков. Применения включают контексты, в которых количество доступных процессоров со временем меняется.
Примечание реализации:
Эта реализация ограничивает максимальное количество запущенных потоков 32767
Параметры:
size - целевой уровень параллельности
Возвращает:
предыдущий уровень параллельности.
Исключения:
IllegalArgumentException - если размер меньше 1 или больше максимального, поддерживаемого этим пулом.
UnsupportedOperationException - это commonPool() и уровень параллельности был задан системной переменной java.util.concurrent.ForkJoinPool.common.parallelism.
С:
19

invokeAllUninterruptibly

public <T> List<Future<T>> invokeAllUninterruptibly(Collection<? extends Callable<T>> tasks)
Непрерывная версия invokeAll. Выполняет заданные задачи, возвращая список Futures, содержащих их статус и результаты, когда все завершены, игнорируя прерывания. Future.isDone() имеет значение true для каждого элемента возвращаемого списка. Обратите внимание, что завершенная задача могла завершиться как нормально, так и сбросив исключение. Результаты этого метода не определены, если заданный список изменяется во время выполнения этой операции.
Примечание API:
Этот метод поддерживает использование, которое ранее полагалось на несовместимую переопределённую версию ExecutorService.invokeAll(java.util.Collection).
Параметры типа:
T - тип значений, возвращаемых задачами
Параметры:
tasks - коллекция задач
Возвращает:
список Futures, представляющих задачи в том же последовательном порядке, что и итератор для данного списка задач, каждая из которых завершилась
Исключения:
NullPointerException - если задачи или какой-либо из их элементов являются null
RejectedExecutionException - если какая-либо задача не может быть запланирована для выполнения
С:
22

getFactory

public ForkJoinPool.ForkJoinWorkerThreadFactory getFactory()
Возвращает фабрику, используемую для создания новых рабочих потоков.
Возвращает:
фабрика, используемая для создания новых рабочих потоков

getUncaughtExceptionHandler

public Thread.UncaughtExceptionHandler getUncaughtExceptionHandler()
Возвращает обработчик для внутренних потоков рабочих, которые завершаются из-за необработанных ошибок, возникших при выполнении задач.
Возвращает:
обработчик, или null, если он отсутствует

getParallelism

public int getParallelism()
Возвращает целевой уровень параллельности этого пула.
Возвращает:
целевой уровень параллельности этого пула

getCommonPoolParallelism

public static int getCommonPoolParallelism()
Возвращает целевой уровень параллельности общего пула.
Возвращает:
целевой уровень параллельности общего пула
С:
1.8

getPoolSize

public int getPoolSize()
Возвращает количество рабочих потоков, которые были начаты, но ещё не завершены. Результат, возвращаемый этим методом, может отличаться от getParallelism(), когда потоки создаются для поддержания параллельности, когда другие добровольно блокируются.
Возвращает:
количество рабочих потоков

getAsyncMode

public boolean getAsyncMode()
Возвращает true, если этот пул использует локальный режим планирования FIFO для разветвлённых задач, которые никогда не объединяются.
Возвращает:
true, если этот пул использует асинхронный режим

getRunningThreadCount

public int getRunningThreadCount()
Возвращает приблизительное количество рабочих потоков, которые не заблокированы, ожидая объединения задач или других управляемых синхронизаций. Этот метод может переоценивать количество работающих потоков.
Возвращает:
количество рабочих потоков

getActiveThreadCount

public int getActiveThreadCount()
Возвращает приблизительное количество потоков, которые в данный момент крадут или выполняют задачи. Этот метод может переоценивать количество активных потоков.
Возвращает:
количество активных потоков

isQuiescent

public boolean isQuiescent()
Возвращает true, если все рабочие потоки в данный момент бездействуют. Бездействующий рабочий - это тот, который не может получить задачу для выполнения, потому что ни одна не доступна для кражи у других потоков, и нет ожидающих отправлений в пул. Этот метод является консервативным; он может не вернуть true сразу же после бездействия всех потоков, но в конечном итоге станет истинным, если потоки останутся неактивными.
Возвращает:
true, если все потоки в данный момент бездействуют
END_OF_DOCUMENT_MARKER

getStealCount

public long getStealCount()
Возвращает оценку общего количества завершённых задач, выполненных потоком, отличным от потока-задающего. Сообщённое значение недооценивает фактическое общее количество краж, когда пул не находится в состоянии покоя. Это значение может быть полезно для мониторинга и настройки программ fork/join: как правило, счёт краж должен быть достаточно высоким, чтобы держать потоки занятыми, но достаточно низким, чтобы избежать накладных расходов и конкуренции между потоками.
Возвращает:
количество краж

getQueuedTaskCount

public long getQueuedTaskCount()
Возвращает оценку общего количества задач, в настоящее время находящихся в очередях потоками-работниками (но не включая задачи, отправленные в пул, которые ещё не начали выполнение). Это значение является лишь приблизительным, полученным путём итерации по всем потокам в пуле. Данный метод может быть полезен для настройки гранул задач.
Возвращает:
количество задач в очереди
См. также:
  • ForkJoinWorkerThread.getQueuedTaskCount()

getQueuedSubmissionCount

public int getQueuedSubmissionCount()
Возвращает оценку количества задач, отправленных в этот пул, которые ещё не начали выполнение. Данный метод может занимать время, пропорциональное количеству отправленных задач.
Возвращает:
количество задач в очереди

hasQueuedSubmissions

public boolean hasQueuedSubmissions()
Возвращает true, если в этот пул были отправлены какие-либо задачи, которые ещё не начали выполнение.
Возвращает:
true, если есть задачи в очереди

pollSubmission

protected ForkJoinTask<?> pollSubmission()
Удаляет и возвращает следующую неосуществлённую отправку, если она доступна. Данный метод может быть полезен в расширениях этого класса, которые повторно назначают работу в системах с несколькими пулами.
Возвращает:
следующую отправку или null, если нет

drainTasksTo

protected int drainTasksTo(Collection<? super ForkJoinTask<?>> c)
Удаляет все доступные неосуществлённые отправленные и форкнутые задачи из очередей планирования и добавляет их в заданный набор, не изменяя их статус выполнения. Они могут включать искусственно сгенерированные или обернутые задачи. Данный метод предназначен для вызова только в том случае, когда известно, что пул находится в состоянии покоя. Вызовы в другие моменты времени могут не удалить все задачи. Ошибка, возникшая при попытке добавить элементы в набор c, может привести к тому, что элементы окажутся ни в одном, ни в том, ни в другом, или в обоих наборах, когда будет выброшено соответствующее исключение. Поведение этой операции не определено, если заданный набор изменяется во время выполнения операции.
Параметры:
c - набор для переноса элементов
Возвращает:
количество перенесённых элементов

toString

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

shutdown

public void shutdown()
Возможно инициирует упорядоченное завершение, в котором ранее отправленные задачи выполняются, но новые задачи не принимаются. Вызов не оказывает никакого влияния на состояние выполнения, если это commonPool(), и не оказывает никакого дополнительного влияния, если уже завершён. Задачи, которые находятся в процессе отправки одновременно во время выполнения этого метода, могут быть отклонены или нет.

shutdownNow

public List<Runnable> shutdownNow()
Возможно пытается отменить и/или остановить все задачи и отклонить все последующие отправленные задачи. Вызов не оказывает никакого влияния на состояние выполнения, если это commonPool(), и не оказывает никакого дополнительного влияния, если уже завершён. В противном случае задачи, которые находятся в процессе отправки или выполнения одновременно во время выполнения этого метода, могут быть отклонены или нет. Этот метод отменяет как существующие, так и невыполненные задачи, чтобы позволить завершение в присутствии зависимостей задач. Таким образом, метод всегда возвращает пустой список (в отличие от случая с некоторыми другими Executors).
Возвращает:
пустой список

isTerminated

public boolean isTerminated()
Возвращает true, если все задачи завершены после завершения работы.
Возвращает:
true, если все задачи завершены после завершения работы

isTerminating

public boolean isTerminating()
Возвращает true, если процесс завершения начался, но еще не завершен. Этот метод может быть полезен для отладки. Возврат значения true через достаточный промежуток времени после завершения работы может указывать на то, что отправленные задачи игнорируют или подавляют прерывание или ожидают ввода-вывода, что приводит к неправильному завершению данного исполнителя. (См. справочные заметки к классу ForkJoinTask, где указывается, что задачи обычно не должны включать блокирующие операции. Но если они это делают, они должны прерывать их при прерывании.)
Возвращает:
true, если завершается, но еще не завершен

isShutdown

public boolean isShutdown()
Возвращает true, если этот пул был остановлен.
Возвращает:
true, если этот пул был остановлен

awaitTermination

public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException
Блокирует до тех пор, пока все задачи не будут завершены после запроса на завершение работы, или не наступит таймаут, или текущий поток не будет прерван — в зависимости от того, что произойдет раньше. Поскольку commonPool() никогда не завершается до завершения программы, при применении к общему пулу этот метод эквивалентен awaitQuiescence(long, TimeUnit), но всегда возвращает false.
Параметры:
timeout - максимальное время ожидания
unit - единица измерения времени для аргумента таймаута
Возвращает:
true, если этот исполнитель завершился, и false, если таймаут истек до завершения
Исключения:
InterruptedException - если произошел перерыв во время ожидания

awaitQuiescence

public boolean awaitQuiescence(long timeout, TimeUnit unit)
Если вызывается ForkJoinTask, работающим в этом пуле, то по сути эквивалентно ForkJoinTask.helpQuiesce(). В противном случае ожидает и/или пытается помочь выполнить задачи до тех пор, пока этот пул isQuiescent() или не истечёт указанный таймаут.
Параметры:
timeout - максимальное время ожидания
unit - единица измерения времени для аргумента таймаута
Возвращает:
true, если пул в состоянии покоя; false, если таймаут истек.

close

public void close()
Если это не commonPool(), инициирует упорядоченное завершение, в котором ранее отправленные задачи выполняются, но новые задачи не будут приниматься, и ожидает, пока все задачи будут завершены, и исполнитель завершится.

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

С:
19

managedBlock

public static void managedBlock(ForkJoinPool.ManagedBlocker blocker) throws InterruptedException
Выполняет заданную, возможно, блокирующую задачу. При выполнении в ForkJoinPool этот метод, возможно, организует активацию резервного потока, если необходимо, чтобы обеспечить достаточную параллельность, в то время как текущий поток заблокирован в blocker.block().

Этот метод многократно вызывает blocker.isReleasable() и blocker.block() до тех пор, пока один из методов не вернет значение true. Каждый вызов blocker.block() предваряется вызовом blocker.isReleasable(), вернувшим false.

Если не выполняется в ForkJoinPool, этот метод поведенчески эквивалентен

 
 while (!blocker.isReleasable())
   if (blocker.block())
     break;
Если выполняется в ForkJoinPool, пул может быть расширен, чтобы обеспечить достаточную параллельность во время вызова blocker.block().
Параметры:
blocker - задача блокирования
Исключения:
InterruptedException - если blocker.block() сделал это

© 1993, 2025, Oracle and/or its affiliates. All rights reserved.
Documentation extracted from Debian's OpenJDK Development Kit package.
Licensed under the GNU General Public License, version 2, with the Classpath Exception.
Various third party code in OpenJDK is licensed under different licenses (see Debian package).
Java and OpenJDK are trademarks or registered trademarks of Oracle and/or its affiliates.
https://download.java.net/java/early_access/jdk24/docs/api/java.base/java/util/concurrent/ForkJoinPool.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API