Spec-Zone.ru › PyTorch 2.14

Эластичный агент

Создано: 4 мая 2021 г. | Последнее обновление: 14 мая 2026 г.

Сервер

Эластичный агент — это плоскость управления torchelastic.

Это процесс, который запускает и управляет нижележащими рабочими процессами. Агент отвечает за:

  1. Работу с распределённым torch: рабочие процессы запускаются со всей необходимой информацией, чтобы успешно и без лишних усилий вызвать torch.distributed.init_process_group().
  2. Отказоустойчивость: агент отслеживает рабочие процессы и при обнаружении сбоев или нездорового состояния останавливает все рабочие процессы и перезапускает их.
  3. Эластичность: реагирует на изменения состава участников и перезапускает рабочие процессы с новым составом.

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

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

../_images/agent_diagram.jpg

Основные понятия

В этом разделе описаны высокоуровневые классы и понятия, необходимые для понимания роли agent в torchelastic.

class torch.distributed.elastic.agent.server.ElasticAgent [источник]

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

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

Точная топология развёртывания и соотношение агентов и рабочих процессов зависят от конкретной реализации агента и предпочтений пользователя относительно размещения задач. Например, чтобы запустить распределённую задачу обучения на GPU с 8 обучающими процессами (по одному на GPU), можно:

  1. Использовать 8 экземпляров с одним GPU, разместив агент на каждом экземпляре, который будет управлять одним рабочим процессом.
  2. Использовать 4 экземпляра с двумя GPU, разместив агент на каждом экземпляре, который будет управлять двумя рабочими процессами.
  3. Использовать 2 экземпляра с четырьмя GPU, разместив агент на каждом экземпляре, который будет управлять четырьмя рабочими процессами.
  4. Использовать один экземпляр с 8 GPU, разместив на нём агента, который будет управлять 8 рабочими процессами.

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

group_result = agent.run()
 if group_result.is_failed():
   # workers failed
   failure = group_result.failures[0]
   logger.exception("worker 0 failed with exit code : %s", failure.exit_code)
 else:
   return group_result.return_values[0] # return rank 0's results
abstract get_worker_group(role='default') [источник]

Возвращает WorkerGroup для указанного role.

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

Тип возвращаемого значения:

WorkerGroup

abstract run(role='default') [источник]

Запускает агента.

Поддерживает повторные попытки запуска группы рабочих процессов при сбоях — до max_restarts.

Возвращает:

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

Вызывает исключения:

Exception — любые другие сбои, НЕ связанные с рабочим процессом –

Тип возвращаемого значения:

RunResult

class torch.distributed.elastic.agent.server.WorkerSpec(role, local_world_size, rdzv_handler, fn=None, entrypoint=None, args=(), max_restarts=3, monitor_interval=0.1, master_port=None, master_addr=None, local_addr=None, event_log_handler='null', numa_options=None, duplicate_stdout_filters=None, duplicate_stderr_filters=None, virtual_local_rank=False) [источник]

Описание-шаблон для конкретного типа рабочего процесса.

Для каждой роли должна существовать только одна спецификация рабочего процесса. Предполагается, что спецификация рабочего процесса одинакова для всех узлов (машин): каждый узел запускает одинаковое число рабочих процессов для данной спецификации.

Параметры:
  • role (str) – определяемая пользователем роль рабочих процессов с этой спецификацией
  • local_world_size (int) – число запускаемых локальных рабочих процессов
  • fn (Callable | None) – (устарело, используйте вместо этого entrypoint)
  • entrypoint (Callable | str | None) – функция рабочего процесса или команда
  • args (tuple) – аргументы, передаваемые в entrypoint
  • rdzv_handler (RendezvousHandler) – обрабатывает rdzv для этого набора рабочих процессов
  • max_restarts (int) – максимальное число повторных попыток для рабочих процессов
  • monitor_interval (float) – интервал проверки состояния рабочих процессов в секундах: n
  • master_port (int | None) – фиксированный порт для запуска хранилища c10d с рангом 0; если не указан, будет выбран случайный свободный порт
  • master_addr (str | None) – фиксированный master_addr для запуска хранилища c10d с рангом 0; если не указан, будет выбрано имя узла агента с рангом 0
  • redirects – перенаправляет стандартные потоки в файл; перенаправление можно выборочно настроить для определённого локального ранга, передав карту
  • tee – одновременно направляет указанные стандартные потоки в консоль и файл; для выборочной настройки определённого локального ранга передайте карту; имеет приоритет над настройками redirects.
  • event_log_handler (str) – имя обработчика журналирования событий, зарегистрированного в elastic/events/handlers.py.
  • duplicate_stdout_filters (list[str] | None) – если список не пуст, копирует stdout в файл, содержащий только строки, соответствующие _любому_ из строк фильтра.
  • duplicate_stderr_filters (list[str] | None) – если список не пуст, копирует stderr в файл, содержащий только строки, соответствующие _любому_ из строк фильтра.
  • virtual_local_rank (bool) – включает режим виртуального локального ранга для рабочих процессов (по умолчанию False). Если этот режим включён, для всех рабочих процессов LOCAL_RANK устанавливается в 0, а CUDA_VISIBLE_DEVICES настраивается так, чтобы каждый рабочий процесс обращался к назначенному ему GPU с индексом устройства 0.
get_entrypoint_name() [источник]

Возвращает имя точки входа.

Если entrypoint — функция (например, Callable), возвращает её __qualname__; если entrypoint — исполняемый файл (например, str), возвращает имя исполняемого файла.

class torch.distributed.elastic.agent.server.WorkerState(value) [источник]

Состояние WorkerGroup.

Рабочие процессы в группе меняют состояние как единое целое. Если один рабочий процесс в группе завершается с ошибкой, вся группа считается завершившейся с ошибкой:

UNKNOWN - agent lost track of worker group state, unrecoverable
INIT - worker group object created not yet started
HEALTHY - workers running and healthy
UNHEALTHY - workers running and unhealthy
STOPPED - workers stopped (interrupted) by the agent
SUCCEEDED - workers finished running (exit 0)
FAILED - workers failed to successfully finish (exit !0)

Группа рабочих процессов начинает с исходного состояния INIT, затем переходит в состояние HEALTHY или UNHEALTHY и в конечном итоге достигает терминального состояния SUCCEEDED или FAILED.

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

  1. Обнаружен сбой или нездоровое состояние группы рабочих процессов
  2. Обнаружено изменение состава участников

Если действия с группой рабочих процессов (запуск, остановка, rdzv, повторная попытка и т. д.) завершаются с ошибкой, в результате которой действие было применено к группе лишь частично, её состояние будет UNKNOWN. Обычно это происходит из-за неперехваченных или необработанных исключений при изменении состояния агента. Агент не должен восстанавливать группы рабочих процессов в состоянии UNKNOWN; лучше завершить работу самого агента, чтобы диспетчер задач повторил запуск узла.

static is_running(state) [источник]

Возвращает состояние рабочего процесса.

Возвращает:

True, если состояние указывает на то, что рабочие процессы ещё выполняются (например, процесс существует, но не обязательно находится в исправном состоянии).

Тип возвращаемого значения:

bool

class torch.distributed.elastic.agent.server.Worker(local_rank, global_rank=-1, role_rank=-1, world_size=-1, role_world_size=-1) [источник]

Экземпляр рабочего процесса.

В отличие от WorkerSpec, который содержит спецификацию рабочего процесса. Worker создаётся на основе WorkerSpec. Worker относится к WorkerSpec так же, как объект к классу.

Значение id рабочего процесса интерпретируется конкретной реализацией ElasticAgent. Для локального агента это может быть pid (int) рабочего процесса, а для удалённого агента оно может быть закодировано как host:port (string).

Параметры:
  • id (Any) – уникально идентифицирует рабочий процесс (интерпретируется агентом)
  • local_rank (int) – локальный ранг рабочего процесса
  • global_rank (int) – глобальный ранг рабочего процесса
  • role_rank (int) – ранг рабочего процесса среди всех рабочих процессов с той же ролью
  • world_size (int) – общее число рабочих процессов
  • role_world_size (int) – число рабочих процессов с той же ролью
class torch.distributed.elastic.agent.server.WorkerGroup(spec) [источник]

Набор экземпляров Worker.

Класс определяет набор экземпляров Worker для указанного WorkerSpec, управляемого ElasticAgent. Наличие в группе рабочих процессов экземпляров с разных узлов зависит от реализации агента.

Реализации

Ниже представлены реализации агента, предоставляемые torchelastic.

class torch.distributed.elastic.agent.server.local_elastic_agent.LocalElasticAgent(spec, logs_specs, start_method='spawn', exit_barrier_timeout=300, log_line_prefix_template=None, shutdown_timeout=30, health_check_server=None, uninterruptible_state_timeout=None) [источник]

Реализация torchelastic.agent.server.ElasticAgent для рабочих процессов на локальном узле.

Этот агент развёртывается на каждом узле и настраивается для запуска n рабочих процессов. При использовании GPU значение n соответствует числу GPU, доступных на узле.

Локальный агент не взаимодействует с другими локальными агентами, развёрнутыми на других узлах, даже если рабочие процессы могут взаимодействовать между узлами. Идентификатор рабочего процесса трактуется как локальный процесс. Агент запускает и останавливает все рабочие процессы как единое целое.

Функция рабочего процесса и передаваемые ей аргументы должны быть совместимы с модулем multiprocessing Python. Чтобы передавать рабочим процессам структуры данных multiprocessing, можно создать структуру данных в том же контексте multiprocessing, что и указанный start_method, и передать её в качестве аргумента функции.

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

Сторожевой процесс на основе именованного канала можно включить в `LocalElasticAgent`, если в процессе `LocalElasticAgent` определена переменная окружения TORCHELASTIC_ENABLE_FILE_TIMER со значением 1. При необходимости можно задать другую переменную окружения `TORCHELASTIC_TIMER_FILE` с уникальным именем файла для именованного канала. Если переменная окружения `TORCHELASTIC_TIMER_FILE` не задана, `LocalElasticAgent` самостоятельно создаст уникальное имя файла и запишет его в переменную окружения `TORCHELASTIC_TIMER_FILE`. Эта переменная окружения будет передана рабочим процессам, чтобы они могли подключиться к тому же именованному каналу, который использует `LocalElasticAgent`.

Журналы записываются в указанный каталог журналов. По умолчанию перед каждой строкой журнала добавляется префикс [${role_name}${local_rank}]: (например, [trainer0]: foobar). Префиксы журналов можно настроить, передав строку шаблона в качестве аргумента log_line_prefix_template. Во время выполнения подставляются следующие макросы (идентификаторы): ${role_name}, ${local_rank}, ${rank}, ${hostname}. Например, чтобы добавлять перед каждой строкой журнала глобальный ранг вместо локального, задайте log_line_prefix_template = "[${rank}]:. ${hostname} раскрывается в имя узла, на котором работает агент, что позволяет определить узел, вызвавший проблему, в задаче с несколькими узлами; например, log_line_prefix_template = "${hostname}:${rank}: " отображается как r12i0n8:3: foobar.

Пример запуска функции

def trainer(args) -> str:
    return "do train"

def main():
    start_method="spawn"
    shared_queue= multiprocessing.get_context(start_method).Queue()
    spec = WorkerSpec(
                role="trainer",
                local_world_size=nproc_per_process,
                entrypoint=trainer,
                args=("foobar",),
                ...<OTHER_PARAMS...>)
    agent = LocalElasticAgent(spec, start_method)
    results = agent.run()

    if results.is_failed():
        print("trainer failed")
    else:
        print(f"rank 0 return value: {results.return_values[0]}")
        # prints -> rank 0 return value: do train

Пример запуска исполняемого файла

def main():
    spec = WorkerSpec(
                role="trainer",
                local_world_size=nproc_per_process,
                entrypoint="/usr/local/bin/trainer",
                args=("--trainer-args", "foobar"),
                ...<OTHER_PARAMS...>)
    agent = LocalElasticAgent(spec)
    results = agent.run()

    if not results.is_failed():
        print("binary launches do not have return values")

Расширение агента

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

class torch.distributed.elastic.agent.server.SimpleElasticAgent(spec, exit_barrier_timeout=300, shutdown_timeout=30) [источник]

ElasticAgent, управляющий конкретным типом роли рабочего процесса.

ElasticAgent, управляющий рабочими процессами (WorkerGroup) для одного WorkerSpec, например для конкретного типа роли рабочего процесса.

_assign_worker_ranks(store, group_rank, group_world_size, spec) [источник]

Определяет правильные ранги рабочих процессов.

Быстрый путь: когда у всех рабочих процессов одинаковые роль и размер мира. Глобальный ранг вычисляется как group_rank * group_world_size + local_rank. При этом role_world_size совпадает с global_world_size. В этом случае TCP-хранилище не используется. Этот режим доступен, только если пользователь установил переменную окружения TORCH_ELASTIC_WORKER_IDENTICAL в значение 1.

Временная сложность: O(1) для каждого рабочего процесса, O(1) в целом

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

  1. Каждый агент записывает в общее хранилище свою конфигурацию (group_rank, group_world_size, num_workers).
  2. Агент с рангом 0 считывает из хранилища все значения role_info и определяет ранги рабочих процессов каждого агента.
  3. Определяется глобальный ранг: он вычисляется как накопленная сумма local_world_size всех предшествующих рабочих процессов. Для повышения эффективности каждому рабочему процессу назначается базовый глобальный ранг, так что его рабочие процессы находятся в диапазоне [base_global_rank, base_global_rank + local_world_size).
  4. Определяется ранг роли: он рассчитывается по алгоритму из пункта 3, но ранги вычисляются относительно имени роли.
  5. Агент с рангом 0 записывает назначенные ранги в хранилище.
  6. Каждый агент считывает назначенные ранги из хранилища.

Временная сложность: O(1) для каждого рабочего процесса, O(n) для rank0, O(n) в целом

Тип возвращаемого значения:

list[Worker]

_exit_barrier() [источник]

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

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

_initialize_workers(worker_group) [источник]

Запускает новый набор рабочих процессов для worker_group.

По сути, выполняет rendezvous, а затем start_workers. Перед вызовом этого метода вызывающий код должен сначала вызвать _stop_workers(), чтобы остановить работающие процессы.

Состояние только что запущенной группы рабочих процессов оптимистично устанавливается в HEALTHY, а фактический мониторинг состояния передаётся методу _monitor_workers()

abstract _monitor_workers(worker_group) [источник]

Проверяет worker_group рабочих процессов.

Эта функция также возвращает новое состояние группы рабочих процессов.

Тип возвращаемого значения:

RunResult

_rendezvous(worker_group) [источник]

Выполняет rendezvous для рабочих процессов, указанных в спецификации.

Назначает рабочим процессам новые глобальные ранги и размер мира. Обновляет хранилище rendezvous для группы рабочих процессов.

_restart_workers(worker_group) [источник]

Перезапускает (останавливает, выполняет rendezvous и запускает) все локальные рабочие процессы в группе.

abstract _shutdown(death_sig=Signals.SIGTERM, timeout=30) [источник]

Освобождает все ресурсы, выделенные во время работы агента.

Параметры:
  • death_sig (Signals) – сигнал для отправки дочернему процессу; по умолчанию SIGTERM
  • timeout (int) – время ожидания корректного завершения работы перед отправкой SIGKILL
abstract _start_workers(worker_group) [источник]

Запускает worker_group.spec.local_world_size рабочих процессов.

Количество определяется спецификацией группы рабочих процессов. Возвращает карту соответствий между local_rank и id рабочего процесса.

Тип возвращаемого значения:

dict[int, Any]

abstract _stop_workers(worker_group) [источник]

Останавливает все рабочие процессы в указанной группе.

Разработчики должны учитывать все состояния рабочих процессов, определённые в WorkerState. То есть необходимо корректно обрабатывать остановку несуществующих рабочих процессов, нездоровых (зависших) рабочих процессов и т. д.

class torch.distributed.elastic.agent.server.api.RunResult(state, return_values=<factory>, failures=<factory>) [источник]

Результаты выполнения рабочих процессов.

Для результатов запуска действует принцип «всё или ничего»: запуск считается успешным тогда и только тогда, когда ВСЕ локальные рабочие процессы, управляемые этим агентом, завершаются успешно.

Если запуск успешен (например, is_failed() = False), поле return_values содержит выходные данные (возвращаемые значения) рабочих процессов, управляемых ЭТИМ агентом, сопоставленные с их ГЛОБАЛЬНЫМИ рангами. Иными словами, result.return_values[0] — это возвращаемое значение рабочего процесса с глобальным рангом 0.

Примечание

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

Если is_failed() возвращает True, поле failures содержит сведения о сбоях, также сопоставленные с ГЛОБАЛЬНЫМ рангом рабочего процесса, в котором произошёл сбой.

Ключи в return_values и failures взаимоисключающие: итоговое состояние рабочего процесса может быть только одним из следующих: успешно завершён или завершился с ошибкой. Рабочие процессы, намеренно остановленные агентом в соответствии с политикой его перезапуска, не представлены ни в return_values, ни в failures.

Сторожевой процесс агента

Сторожевой процесс на основе именованного канала можно включить в LocalElasticAgent, если в процессе LocalElasticAgent определена переменная окружения TORCHELASTIC_ENABLE_FILE_TIMER со значением 1. При необходимости можно задать другую переменную окружения TORCHELASTIC_TIMER_FILE с уникальным именем файла для именованного канала. Если переменная окружения TORCHELASTIC_TIMER_FILE не задана, LocalElasticAgent самостоятельно создаст уникальное имя файла и запишет его в переменную окружения TORCHELASTIC_TIMER_FILE. Эта переменная окружения будет передана рабочим процессам, чтобы они могли подключиться к тому же именованному каналу, который использует LocalElasticAgent.

Сервер проверки работоспособности

Сервер мониторинга работоспособности можно включить в LocalElasticAgent, если в процессе LocalElasticAgent задана переменная среды TORCHELASTIC_HEALTH_CHECK_PORT. Добавляется интерфейс сервера проверки работоспособности, который можно расширить, запустив TCP/HTTP-сервер на указанном номере порта. Кроме того, сервер проверки работоспособности будет иметь обратный вызов для проверки того, что сторожевой таймер активен.

class torch.distributed.elastic.agent.server.health_check_server.HealthCheckServer(alive_callback, port, timeout) [источник]

Интерфейс сервера мониторинга работоспособности, который можно расширить, запустив TCP/HTTP-сервер на указанном порту.

Параметры:
  • alive_callback (Callable[[], int]) – Callable[[], int], обратный вызов, возвращающий время последнего продвижения агента
  • port (int) – int, номер порта для запуска TCP/HTTP-сервера
  • timeout (int) – int, время ожидания в секундах для определения, активен или неактивен агент
start() [источник]

Функциональность не поддерживается в PyTorch; сервер проверки работоспособности не запускается

stop() [источник]

Функция остановки сервера проверки работоспособности

torch.distributed.elastic.agent.server.health_check_server.create_healthcheck_server(alive_callback, port, timeout) [источник]

создаёт объект сервера проверки работоспособности

Тип возвращаемого значения:

HealthCheckServer

© 2026, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://docs.pytorch.org/docs/2.14/elastic/agent.html

Spec-Zone.ru

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