В сценарии с поэтапной обработкой, где число участников меняется, почему Phaser подходит лучше CountDownLatch?
Phaser подходит лучше, когда синхронизация состоит из повторяющихся фаз, а состав участников может изменяться. Он позволяет регистрировать новых участников, завершать участие отдельных потоков и автоматически переводить всех ожидающих к следующей фазе. CountDownLatch имеет фиксированный счётчик и после срабатывания не переиспользуется.
Для простой одноразовой координации достаточно CountDownLatch: потоки уменьшают счётчик, а ожидающий поток продолжает работу после его обнуления. На практике часто требовалась более сложная схема — несколько последовательных этапов обработки с повторным ожиданием и динамическим числом рабочих потоков.
Phaser решает эту проблему как переиспользуемый барьер. Он объединяет регистрацию участников, ожидание завершения текущей фазы и переход к следующей фазе в одном примитиве.
Предположим, обработка выполняется партиями. На каждом этапе часть потоков может завершить работу, а новые задачи — подключиться. Использование нового CountDownLatch для каждой фазы требует заранее знать количество участников и отдельно организовывать создание следующего latch.
Ошибки в такой схеме приводят к зависанию: лишний countDown нарушает логику, пропущенный — оставляет потоки в ожидании навсегда. Повторное использование уже сработавшего CountDownLatch невозможно.
У Phaser есть зарегистрированные участники, называемые parties, и номер текущей фазы. Участник может:
register() или bulkRegister();arrive();arriveAndDeregister();arriveAndAwaitAdvance().Фаза завершается, когда все зарегистрированные участники сообщили о прибытии. После этого номер фазы изменяется, а ожидающие участники получают возможность продолжить работу. Следующая фаза использует тот же объект Phaser.
Важно регистрировать участника до того, как он должен учитываться в конкретной фазе. Регистрация, происходящая одновременно с продвижением фазы, может уже не повлиять на текущий барьер. Также каждый зарегистрированный участник обязан либо прибыть, либо явно сняться с регистрации; иначе остальные потоки могут ждать бесконечно.
Минимальный пример повторяющихся фаз:
Phaser не делает работу потоков параллельной сам по себе: он только координирует точки встречи. Для динамического добавления участников регистрацию обычно выполняют управляющий поток или уже зарегистрированный участник, а жизненный цикл регистрации проектируют явно.
По сравнению с ручной комбинацией CountDownLatch, очередей и дополнительных флагов Phaser сокращает служебную логику. Однако он сложнее для одноразовой задачи, а ошибка в управлении регистрацией может быть менее очевидной, чем ошибка в простом счётчике latch.
Сервис обрабатывает пакет файлов в несколько этапов. После первого этапа часть рабочих потоков больше не нужна, а для следующего этапа могут подключаться дополнительные обработчики. Вариант с новым CountDownLatch на каждый этап потребовал бы передавать новый объект всем участникам и отдельно согласовывать точное число задач.
Ручная схема на wait/notify дала бы гибкость, но добавила бы риски ложных пробуждений, ошибок с условием ожидания и проблем с публикацией состояния. Набор latch-объектов проще, но плохо подходит для динамического числа участников и требует дополнительного протокола между этапами.
Был выбран Phaser: управляющий поток регистрирует обработчики, завершившие работу обработчики вызывают arriveAndDeregister(), а оставшиеся используют arriveAndAwaitAdvance(). Это позволило переиспользовать один барьер для всех этапов и избежать зависания из-за ожидания завершивших работу участников. При этом в проекте отдельно добавили контроль тайм-аутов и обработку аварийного завершения, поскольку сам Phaser не исправляет пропущенную регистрацию или прибытие.
Что произойдёт, если участник завершил работу, но вызвал arrive(), а не arriveAndDeregister()?
Он сообщит о завершении текущей фазы, но останется зарегистрированным в Phaser. В следующей фазе барьер продолжит считать его участником, поэтому при отсутствии нового прибытия фаза может навсегда заблокировать остальных. Для окончательного выхода нужен arriveAndDeregister().
Чем отличается arrive() от arriveAndAwaitAdvance()?
arrive() только фиксирует прибытие участника и сразу возвращает управление. Поток не обязан ждать, пока остальные участники завершат фазу. arriveAndAwaitAdvance() сначала фиксирует прибытие, а затем ждёт перехода фазы, поэтому подходит для настоящей точки синхронизации.
Что происходит, если один участник никогда не прибывает к барьеру?
Текущая фаза не завершится, потому что Phaser ждёт прибытия всех зарегистрированных участников. Это не ошибка самого примитива: он не может определить, действительно ли поток задержался или должен быть исключён из участия.
Практическое решение — гарантировать прибытие в finally, использовать arriveAndDeregister() при отмене работы или предусмотреть контролируемое завершение через наследование и переопределение onAdvance. Принудительное завершение (forceTermination) освобождает ожидающих, но переводит phaser в терминальное состояние и должно рассматриваться как аварийный сценарий, а не как обычный способ синхронизации.