Класс Phaser
public class Phaser extends Object
CyclicBarrier и CountDownLatch, но с более гибким использованием. Регистрация. В отличие от других барьеров, количество участников, зарегистрированных для синхронизации на фазе, может изменяться со временем. Задачи могут быть зарегистрированы в любое время (с помощью методов register(), bulkRegister(int) или форм конструкторов, устанавливающих начальное количество участников), и при желании могут быть дезактивированы по достижении точки прибытия (с помощью arriveAndDeregister()). Как и в случае с большинством основных конструкций синхронизации, регистрация и дерегистрация влияют только на внутренние подсчеты; они не устанавливают дальнейшую внутреннюю бухгалтерию, поэтому задачи не могут запросить, зарегистрированы ли они. (Однако вы можете ввести такую бухгалтерию, создав подкласс этого класса.)
Синхронизация. Как и CyclicBarrier, фазу можно многократно ожидать. Метод arriveAndAwaitAdvance() имеет эффект, аналогичный CyclicBarrier.await. Каждая генерация фазы имеет связанное с ней значение фазы. Значение фазы начинается с нуля и увеличивается, когда все участники прибывают на фазу, возвращаясь к нулю после достижения
Integer.MAX_VALUE. Использование номеров фаз позволяет независимо контролировать действия по прибытию на фазу и ожидания других, с помощью двух видов методов, которые могут вызываться любым зарегистрированным участником:
-
Прибытие. Методы
arrive()иarriveAndDeregister()регистрируют прибытие. Эти методы не блокируют, но возвращают связанное с ними значение фазы прибытия; то есть значение фазы фазера, к которой применяется прибытие. Когда последний участник для данной фазы прибывает, выполняется необязательное действие, и фаза переходит к следующей. Эти действия выполняет участник, запускающий переход на следующую фазу, и они организуются путем переопределения методаonAdvance(int, int), который также контролирует завершение. Переопределение этого метода аналогично, но более гибко, чем предоставление действия барьераCyclicBarrier. -
Ожидание. Метод
awaitAdvance(int)требует аргумента, указывающего значение фазы прибытия, и возвращает значение, когда фаза переходит к (или уже находится на) другой фазе. В отличие от аналогичных конструкций с использованиемCyclicBarrier, методawaitAdvanceпродолжает ждать, даже если поток ожидания прерывается. Также доступны прерываемые и ограниченные по времени версии, но исключения, возникающие во время ожидания задач с прерыванием или с ограничением по времени, не изменяют состояние фазы. При необходимости вы можете выполнить любое связанное восстановление в обработчиках этих исключений, часто после вызоваforceTermination. Фазеры также могут использоваться задачами, выполняющимися вForkJoinPool. Прогресс гарантируется, если уровень параллелизма пула может вместить максимальное количество одновременно заблокированных участников.
Завершение. Фазер может перейти в состояние завершения, которое можно проверить с помощью метода isTerminated(). При завершении все методы синхронизации немедленно возвращаются без ожидания перехода, как указано отрицательным возвращаемым значением. Аналогично, попытки зарегистрироваться при завершении не имеют эффекта. Завершение происходит, когда вызов onAdvance возвращает true. В реализации по умолчанию возвращается
true при дерегистрации, приведшей к тому, что количество зарегистрированных участников стало равным нулю. Как показано ниже, когда фазеры управляют действиями с фиксированным числом итераций, часто удобно переопределить этот метод, чтобы вызвать завершение, когда текущее значение фазы достигает порогового значения. Также доступен метод forceTermination() для внезапного освобождения ожидающих потоков и разрешения им завершиться.
Уровни. Фазеры могут быть многоуровневыми (т. е. созданными в структурах дерева) для уменьшения конкуренции. Фазеры с большим количеством участников, которые в противном случае испытывали бы значительные затраты на синхронизацию из-за конкуренции, могут быть настроены таким образом, чтобы группы под-фазеров делили общего родителя. Это может значительно увеличить пропускную способность, хотя это влечет за собой большие накладные расходы на операцию.
В дереве многоуровневых фазеров регистрация и дерегистрация дочерних фазеров с их родительским элементом управляется автоматически. Всякий раз, когда количество зарегистрированных участников дочернего фазера становится отличным от нуля (как определено в конструкторе Phaser(Phaser,int), register() или bulkRegister(int)), дочерний фазер регистрируется у родителя. Всякий раз, когда количество зарегистрированных участников становится равным нулю в результате вызова arriveAndDeregister(), дочерний фазер деактивируется от родителя.
Мониторинг. Хотя методы синхронизации могут вызываться только зарегистрированными участниками, текущее состояние фазера может быть отслеживаться любым вызывающим элементом. В любой момент времени существует getRegisteredParties() участников в общем, из которых getArrivedParties() прибыли на текущую фазу (getPhase()). Когда оставшиеся (getUnarrivedParties()) участники прибывают, фаза переходит на следующую. Возвращаемые значения этих методов могут отражать промежуточные состояния и, следовательно, в общем случае не подходят для управления синхронизацией. Метод toString() возвращает снимки этих запросов состояния в удобной для неформального мониторинга форме.
Эффекты согласованности памяти: Действия до любого метода arrive происходят раньше соответствующего перехода на следующую фазу и действий onAdvance (если они есть), которые в свою очередь происходят раньше действий после перехода на следующую фазу.
Примеры использования:
Фазу можно использовать вместо CountDownLatch для управления одноразовым действием, обслуживающим переменное количество участников. Типичный шаблон заключается в том, что метод, настраивающий это, сначала регистрирует, затем запускает все действия, затем деактивирует, как показано ниже:
void runTasks(List<Runnable> tasks) {
Phaser startingGate = new Phaser(1); // "1" to register self
// create and start threads
for (Runnable task : tasks) {
startingGate.register();
new Thread(() -> {
startingGate.arriveAndAwaitAdvance();
task.run();
}).start();
}
// deregister self to allow threads to proceed
startingGate.arriveAndDeregister();
} Один из способов заставить набор потоков многократно выполнять действия для заданного количества итераций заключается в переопределении onAdvance:
void startTasks(List<Runnable> tasks, int iterations) {
Phaser phaser = new Phaser() {
protected boolean onAdvance(int phase, int registeredParties) {
return phase >= iterations - 1 || registeredParties == 0;
}
};
phaser.register();
for (Runnable task : tasks) {
phaser.register();
new Thread(() -> {
do {
task.run();
phaser.arriveAndAwaitAdvance();
} while (!phaser.isTerminated());
}).start();
}
// allow threads to proceed; don't wait for them
phaser.arriveAndDeregister();
} Если основной задаче необходимо позже дождаться завершения, она может повторно зарегистрироваться и затем выполнить аналогичный цикл:
// ...
phaser.register();
while (!phaser.isTerminated())
phaser.arriveAndAwaitAdvance(); Связанные конструкции могут использоваться для ожидания конкретных номеров фаз в контекстах, где вы уверены, что фаза никогда не перейдет в состояние циклической замены Integer.MAX_VALUE. Например:
void awaitPhase(Phaser phaser, int phase) {
int p = phaser.register(); // assumes caller not already registered
while (p < phase) {
if (phaser.isTerminated())
// ... deal with unexpected termination
else
p = phaser.arriveAndAwaitAdvance();
}
phaser.arriveAndDeregister();
} Для создания набора n задач с помощью дерева фазеров вы могли бы использовать код следующей формы, предполагая класс задачи с конструктором, принимающим Phaser, с которым он регистрируется при создании. После вызова build(new Task[n], 0, n,
new Phaser()), эти задачи можно было бы запустить, например, путем передачи в пул:
void build(Task[] tasks, int lo, int hi, Phaser ph) {
if (hi - lo > TASKS_PER_PHASER) {
for (int i = lo; i < hi; i += TASKS_PER_PHASER) {
int j = Math.min(i + TASKS_PER_PHASER, hi);
build(tasks, i, j, new Phaser(ph));
}
} else {
for (int i = lo; i < hi; ++i)
tasks[i] = new Task(ph);
// assumes new Task(ph) performs ph.register()
}
} Лучшее значение TASKS_PER_PHASER в основном зависит от ожидаемых скоростей синхронизации. Значение, равное четырем, может быть подходящим для очень небольших задач тела фазы (следовательно, высокая скорость) или до сотен для очень больших. Примечания по реализации: Эта реализация ограничивает максимальное количество участников 65535. Попытки зарегистрировать дополнительные участники приводят к IllegalStateException. Однако вы можете и должны создавать многоуровневые фазеры для вмещения произвольно больших наборов участников.
- Since:
- 1.7
Краткое описание конструкторов
| Конструктор | Описание |
|---|---|
Phaser() |
Создаёт новый объект phaser без зарегистрированных партий, без родителя и начальным номером фазы 0. |
Phaser |
Создаёт новый объект phaser с заданным количеством зарегистрированных незарегистрированных партий, без родителя и начальным номером фазы 0. |
Phaser |
Эквивалентно Phaser(parent, 0). |
Phaser |
Создаёт новый объект phaser с заданным родителем и количеством зарегистрированных незарегистрированных партий. |
Краткое описание методов
| Модификатор и тип | Метод | Описание |
|---|---|---|
int |
arrive() |
Присоединяется к phaser без ожидания других. |
int |
arriveAndAwaitAdvance() |
Присоединяется к phaser и ждёт других. |
int |
arriveAndDeregister() |
Присоединяется к phaser и отменяет регистрацию без ожидания других. |
int |
awaitAdvance |
Ожидает продвижения фазы phaser до заданного значения фазы, возвращаясь немедленно, если текущая фаза не равна заданному значению фазы или этот phaser завершён. |
int |
awaitAdvanceInterruptibly |
Ожидает продвижения фазы phaser до заданного значения фазы, выбрасывая InterruptedException при прерывании ожидания, или возвращаясь немедленно, если текущая фаза не равна заданному значению фазы или этот phaser завершён. |
int |
awaitAdvanceInterruptibly |
Ожидает продвижения фазы phaser до заданного значения фазы или истечения заданного таймаута, выбрасывая
InterruptedException при прерывании ожидания, или возвращаясь немедленно, если текущая фаза не равна заданному значению фазы или этот phaser завершён. |
int |
bulkRegister |
Добавляет заданное количество новых незарегистрированных партий в этот phaser. |
void |
forceTermination() |
Принудительно переводит этот phaser в состояние завершения. |
int |
getArrivedParties() |
Возвращает количество зарегистрированных партий, которые прибыли на текущую фазу этого phaser. |
Phaser |
getParent() |
Возвращает родителя этого phaser, или null если его нет. |
final int |
getPhase() |
Возвращает текущий номер фазы. |
int |
getRegisteredParties() |
Возвращает количество партий, зарегистрированных в этом phaser. |
Phaser |
getRoot() |
Возвращает корневого предка этого phaser, который идентичен этому phaser, если у него нет родителя. |
int |
getUnarrivedParties() |
Возвращает количество зарегистрированных партий, которые ещё не прибыли на текущую фазу этого phaser. |
boolean |
isTerminated() |
Возвращает true если этот phaser завершён. |
protected boolean |
onAdvance |
Переопределяемый метод для выполнения действия при приближении продвижения фазы и управления завершением. |
int |
register() |
Добавляет новую незарегистрированную партию в этот phaser. |
String |
toString() |
Возвращает строку, идентифицирующую этот phaser, а также его состояние. |
Подробное описание конструкторов
Phaser
public Phaser()
Phaser
public Phaser(int parties)
- Параметры:
-
parties- количество участников, необходимое для перехода к следующей фазе - Исключения:
-
IllegalArgumentException- если количество участников меньше нуля или больше максимального поддерживаемого значения
Phaser
public Phaser(Phaser parent)
Phaser(parent, 0).- Параметры:
-
parent- родительский объект Phaser
Phaser
public Phaser(Phaser parent, int parties)
- Параметры:
-
parent- родительский объект Phaser -
parties- количество участников, необходимое для перехода к следующей фазе - Исключения:
-
IllegalArgumentException- если количество участников меньше нуля или больше максимального поддерживаемого значения
Подробное описание методов
register
public int register()
onAdvance(int, int), этот метод может ожидать его завершения перед возвратом. Если у этого объекта Phaser есть родитель и ранее в нём не было зарегистрированных участников, этот дочерний объект Phaser также регистрируется у своего родителя. Если этот объект Phaser завершен, попытка регистрации не имеет эффекта и возвращается отрицательное значение.- Возвращает:
- номер фазы прибытия, к которой относится эта регистрация. Если это значение отрицательно, объект Phaser завершён, и регистрация не имеет эффекта.
- Исключения:
-
IllegalStateException- если попытка зарегистрировать больше, чем максимальное поддерживаемое количество участников
bulkRegister
public int bulkRegister(int parties)
onAdvance(int, int), этот метод может ожидать его завершения перед возвратом. Если у этого объекта Phaser есть родитель, и заданное количество участников больше нуля, и в нём ранее не было зарегистрированных участников, этот дочерний объект Phaser также регистрируется у своего родителя. Если этот объект Phaser завершён, попытка регистрации не имеет эффекта, и возвращается отрицательное значение.- Параметры:
-
parties- количество дополнительных участников, необходимых для перехода к следующей фазе - Возвращает:
- номер фазы прибытия, к которой относится эта регистрация. Если это значение отрицательно, объект Phaser завершён, и регистрация не имеет эффекта.
- Исключения:
-
IllegalStateException- если попытка зарегистрировать больше, чем максимальное поддерживаемое количество участников -
IllegalArgumentException- еслиparties < 0
arrive
public int arrive()
Ошибка использования – попытка вызвать этот метод не зарегистрированным участником. Однако эта ошибка может привести к
IllegalStateException только при последующем выполнении операции с этим объектом Phaser, если таковая произойдёт.
- Возвращает:
- номер фазы прибытия или отрицательное значение, если объект завершен
- Исключения:
-
IllegalStateException- если не завершен и количество не прибывших участников станет отрицательным
arriveAndDeregister
public int arriveAndDeregister()
Ошибка использования – попытка вызвать этот метод не зарегистрированным участником. Однако эта ошибка может привести к
IllegalStateException только при последующем выполнении операции с этим объектом Phaser, если таковая произойдёт.
- Возвращает:
- номер фазы прибытия или отрицательное значение, если объект завершен
- Исключения:
-
IllegalStateException- если не завершен и количество зарегистрированных или не прибывших участников станет отрицательным
arriveAndAwaitAdvance
public int arriveAndAwaitAdvance()
awaitAdvance(arrive()). Если вам нужно ожидание с прерыванием или таймаутом, вы можете организовать это аналогичным образом, используя один из других вариантов метода
awaitAdvance. Если вам нужно отписаться при прибытии, используйте awaitAdvance(arriveAndDeregister()). Ошибка использования – попытка вызвать этот метод не зарегистрированным участником. Однако эта ошибка может привести к
IllegalStateException только при последующем выполнении операции с этим объектом Phaser, если таковая произойдёт.
- Возвращает:
- номер фазы прибытия или (отрицательное) значение текущей фазы, если объект завершен
- Исключения:
-
IllegalStateException- если не завершен и количество не прибывших участников станет отрицательным
awaitAdvance
public int awaitAdvance(int phase)
- Параметры:
-
phase- номер фазы прибытия или отрицательное значение, если объект завершен; этот аргумент обычно является значением, возвращённым предыдущим вызовомarriveилиarriveAndDeregister. - Возвращает:
- следующий номер фазы прибытия, или аргумент, если он отрицательный, или (отрицательное) значение текущей фазы, если объект завершён
awaitAdvanceInterruptibly
public int awaitAdvanceInterruptibly(int phase) throws InterruptedException
InterruptedException если поток прерван во время ожидания, или возвращаясь немедленно, если текущая фаза не равна заданному значению фазы или этот объект Phaser завершён.- Параметры:
-
phase- номер фазы прибытия или отрицательное значение, если объект завершен; этот аргумент обычно является значением, возвращённым предыдущим вызовомarriveилиarriveAndDeregister. - Возвращает:
- следующий номер фазы прибытия, или аргумент, если он отрицательный, или (отрицательное) значение текущей фазы, если объект завершён
- Исключения:
-
InterruptedException- если поток прерван во время ожидания
awaitAdvanceInterruptibly
public int awaitAdvanceInterruptibly(int phase, long timeout, TimeUnit unit) throws InterruptedException, TimeoutException
InterruptedException если поток прерван во время ожидания, или возвращаясь немедленно, если текущая фаза не равна заданному значению фазы или этот объект Phaser завершен.- Параметры:
-
phase- номер фазы прибытия или отрицательное значение, если объект завершен; этот аргумент обычно является значением, возвращённым предыдущим вызовомarriveилиarriveAndDeregister. -
timeout- время ожидания, прежде чем отказаться, в единицахunit -
unit-TimeUnitопределяющий, как интерпретировать параметрtimeout - Возвращает:
- следующий номер фазы прибытия, или аргумент, если он отрицательный, или (отрицательное) значение текущей фазы, если объект завершён
- Исключения:
-
InterruptedException- если поток прерван во время ожидания -
TimeoutException- если истекло время ожидания
forceTermination
public void forceTermination()
getPhase
public final int getPhase()
Integer.MAX_VALUE, после чего он сбрасывается до нуля. При завершении номер фазы отрицательный, в этом случае текущую фазу до завершения можно получить через getPhase() + Integer.MIN_VALUE.- Возвращает:
- номер фазы или отрицательное значение, если объект завершён
getRegisteredParties
public int getRegisteredParties()
- Возвращает:
- количество участников
getArrivedParties
public int getArrivedParties()
- Возвращает:
- количество прибывших участников
getUnarrivedParties
public int getUnarrivedParties()
- Возвращает:
- количество не прибывших участников
getParent
public Phaser getParent()
null , если родителя нет.- Возвращает:
- родителя данного фазера, или
null, если родителя нет
getRoot
public Phaser getRoot()
- Возвращает:
- корневого предка данного фазера
isTerminated
public boolean isTerminated()
true , если данный фазер был завершён.- Возвращает:
-
true, если данный фазер был завершён
onAdvance
protected boolean onAdvance(int phase, int registeredParties)
true, этот фазер будет переведён в конечное состояние завершения при продвижении, и последующие вызовы isTerminated() вернут true. Любое (необработанное) исключение или ошибка, сгенерированные при вызове этого метода, передаются стороне, пытающейся продвинуть этот фазер, в результате чего продвижение не происходит. Аргументы этого метода предоставляют состояние фазера, преобладающее для текущего перехода. Эффекты вызова методов прибытия, регистрации и ожидания для данного фазера изнутри onAdvance не определены и на них не следует полагаться.
Если этот фазер является членом иерархической группы фазеров, то onAdvance вызывается только для корневого фазера при каждом продвижении.
Для поддержки наиболее распространённых случаев использования, по умолчанию этот метод возвращает true , когда количество зарегистрированных сторон стало нулевым в результате вызова arriveAndDeregister стороной. Вы можете отключить это поведение, тем самым разрешив продолжение при последующих регистрациях, переопределив этот метод так, чтобы он всегда возвращал false:
Phaser phaser = new Phaser() {
protected boolean onAdvance(int phase, int parties) { return false; }
};
- Параметры:
-
phase- текущий номер фазы при входе в этот метод, перед продвижением этого фазера -
registeredParties- текущее количество зарегистрированных сторон - Возвращает:
-
true, если этот фазер должен быть завершён
toString
public String toString()
"phase = " , за которой следует номер фазы, "parties = " , за которой следует количество зарегистрированных сторон, и
"arrived = " , за которой следует количество прибывших сторон.- Переопределяет:
-
toStringв классеObject - Возвращает:
- строку, идентифицирующую этот фазер, а также его состояние
© 1993, 2023, 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://docs.oracle.com/en/java/javase/21/docs/api/java.base/java/util/concurrent/Phaser.html