Spec-Zone.ru › PyTorch 2.14

Многопроцессность

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

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

Для функций она использует torch.multiprocessing (и, следовательно, python multiprocessing) для порождения/ветвления рабочих процессов. Для исполняемых файлов она использует python subprocessing.Popen для создания рабочих процессов.

Пример 1: запуск двух обучающих процессов в виде функции

from torch.distributed.elastic.multiprocessing import Std, start_processes


def trainer(a, b, c):
    pass  # train


# runs two trainers
# LOCAL_RANK=0 trainer(1,2,3)
# LOCAL_RANK=1 trainer(4,5,6)
ctx = start_processes(
    name="trainer",
    entrypoint=trainer,
    args={0: (1, 2, 3), 1: (4, 5, 6)},
    envs={0: {"LOCAL_RANK": 0}, 1: {"LOCAL_RANK": 1}},
    log_dir="/tmp/foobar",
    redirects=Std.ALL,  # write all worker stdout/stderr to a log file
    tee={0: Std.ERR},  # tee only local rank 0's stderr to console
)

# waits for all copies of trainer to finish
ctx.wait()

Пример 2: запуск 2 рабочих процессов echo в виде исполняемого файла

# same as invoking
# echo hello
# echo world > stdout.log
ctx = start_processes(
        name="echo"
        entrypoint="echo",
        log_dir="/tmp/foobar",
        args={0: "hello", 1: "world"},
        redirects={1: Std.OUT},
       )

Как и в torch.multiprocessing, функция start_processes() возвращает контекст процесса (api.PContext). Если была запущена функция, возвращается api.MultiprocessContext, а если был запущен исполняемый файл — api.SubprocessContext. Оба класса являются конкретными реализациями родительского класса api.PContext.

Запуск нескольких рабочих процессов

torch.distributed.elastic.multiprocessing.start_processes(name, entrypoint, args, envs, logs_specs, log_line_prefixes=None, start_method='spawn', numa_options=None, duplicate_stdout_filters=None, duplicate_stderr_filters=None) [исходный код]

Запускает n копий процессов entrypoint с указанными параметрами.

entrypoint — это либо Callable (функция), либо str (исполняемый файл). Количество копий определяется числом записей в аргументах args и envs; наборы ключей в них должны совпадать.

Параметры args и env — это аргументы и переменные среды, передаваемые точке входа и сопоставленные с индексом реплики (локальным рангом). Необходимо указать все локальные ранги. То есть набор ключей должен быть {0,1,...,(nprocs-1)}.

Примечание

Если entrypoint — это исполняемый файл (str), args может содержать только строки. Если задано значение другого типа, оно преобразуется в строковое представление (например, str(arg1)). Кроме того, при сбое исполняемого файла файл ошибки error.json будет записан только в том случае, если основная функция аннотирована с помощью torch.distributed.elastic.multiprocessing.errors.record. При запуске функций это выполняется по умолчанию, и вручную добавлять аннотацию @record не требуется.

Внутри logs_specs параметры redirects и tee представляют собой битовые маски, указывающие, какие стандартные потоки перенаправлять в файл журнала в log_dir. Допустимые значения масок определены в Std. Чтобы перенаправлять/дублировать в консоль потоки только для определённых локальных рангов, передайте redirects в виде карты, где ключ — локальный ранг, для которого задаётся поведение перенаправления. Для всех отсутствующих локальных рангов по умолчанию используется Std.NONE.

Параметры duplicate_stdout_filters и duplicate_stderr_filters, если они не пусты, копируют стандартные потоки вывода и ошибок соответственно, заданные в logs_specs параметром tee, в файл, содержащий только строки, соответствующие _любому_ из фильтров. Файл журнала объединяет данные всех рангов, выбранных параметром tee.

tee работает подобно команде Unix «tee»: она перенаправляет вывод и одновременно печатает его в консоль. Чтобы вывод и ошибки рабочих процессов не печатались в консоль, используйте параметр redirects.

Для каждого процесса в log_dir будут содержаться:

  1. {local_rank}/error.json: если процесс завершился с ошибкой, файл с информацией об ошибке
  2. {local_rank}/stdout.log: если redirect & STDOUT == STDOUT
  3. {local_rank}/stderr.log: если redirect & STDERR == STDERR
  4. filtered_stdout.log: если duplicate_stdout_filters не пуст
  5. filtered_stderr.log: если duplicate_stderr_filters не пуст

Примечание

Предполагается, что log_dir существует, пуст и является каталогом.

Пример:

log_dir = "/tmp/test"

# ok; two copies of foo: foo("bar0"), foo("bar1")
start_processes(
   name="trainer",
   entrypoint=foo,
   args:{0:("bar0",), 1:("bar1",),
   envs:{0:{}, 1:{}},
   log_dir=log_dir
)

# invalid; envs missing for local rank 1
start_processes(
   name="trainer",
   entrypoint=foo,
   args:{0:("bar0",), 1:("bar1",),
   envs:{0:{}},
   log_dir=log_dir
)

# ok; two copies of /usr/bin/touch: touch file1, touch file2
start_processes(
   name="trainer",
   entrypoint="/usr/bin/touch",
   args:{0:("file1",), 1:("file2",),
   envs:{0:{}, 1:{}},
   log_dir=log_dir
 )

# caution; arguments casted to string, runs:
# echo "1" "2" "3" and echo "[1, 2, 3]"
start_processes(
   name="trainer",
   entrypoint="/usr/bin/echo",
   args:{0:(1,2,3), 1:([1,2,3],),
   envs:{0:{}, 1:{}},
   log_dir=log_dir
 )
Параметры:
  • name (str) — короткое понятное название, описывающее процессы (используется в качестве заголовка при дублировании вывода и ошибок в консоль)
  • entrypoint (Callable | str) — либо Callable (функция), либо cmd (исполняемый файл)
  • args (dict[int, tuple]) — аргументы для каждой реплики
  • envs (dict[int, dict[str, str]]) — переменные среды для каждой реплики
  • log_dir — каталог для записи файлов журнала
  • start_method (str) — метод запуска многопроцессной обработки (spawn, fork, forkserver); для исполняемых файлов не используется
  • logs_specs (LogsSpecs) — определяет log_dir, redirects и tee. внутри logs_specs: - redirects: какие стандартные потоки перенаправлять в файл журнала - tee: какие стандартные потоки перенаправлять и одновременно выводить в консоль
  • local_ranks_filter — журналы каких рангов выводить в консоль
  • duplicate_stdout_filters (list[str] | None) — фильтры для дублируемых журналов стандартного вывода
  • duplicate_stderr_filters (list[str] | None) — фильтры для дублируемых журналов стандартного потока ошибок
Тип возвращаемого значения:

PContext

Контекст процесса

class torch.distributed.elastic.multiprocessing.api.PContext(name, entrypoint, args, envs, logs_specs, log_line_prefixes=None, duplicate_stdout_filters=None, duplicate_stderr_filters=None) [исходный код]

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

Имя PContext выбрано намеренно, чтобы избежать неоднозначности с torch.multiprocessing.ProcessContext.

Предупреждение

stdout и stderr ВСЕГДА должны быть надмножеством tee_stdout и tee_stderr (соответственно), поскольку tee реализован как перенаправление + tail -f <stdout/stderr.log>

Параметры:
  • duplicate_stdout_filters (list[str] | None) – Если список не пуст, дублирует stdout, указанный в logs_specs’s tee, в файл, содержащий только строки, совпадающие с _любой_ из строк фильтра. Файл журнала объединяет данные всех рангов, выбранных с помощью tee.
  • duplicate_stderr_filters (list[str] | None) – Если список не пуст, дублирует stderr, указанный в logs_specs’s tee, в файл, содержащий только строки, совпадающие с _любой_ из строк фильтра. Файл журнала объединяет данные всех рангов, выбранных с помощью tee.
class torch.distributed.elastic.multiprocessing.api.MultiprocessContext(name, entrypoint, args, envs, start_method, logs_specs, log_line_prefixes=None, numa_options=None, duplicate_stdout_filters=None, duplicate_stderr_filters=None) [исходный код]

PContext, содержащий рабочие процессы, вызванные как функция.

class torch.distributed.elastic.multiprocessing.api.SubprocessContext(name, entrypoint, args, envs, logs_specs, log_line_prefixes=None, numa_options=None, duplicate_stdout_filters=None, duplicate_stderr_filters=None) [исходный код]

PContext, содержащий рабочие процессы, вызванные как исполняемый файл.

class torch.distributed.elastic.multiprocessing.api.RunProcsResult(return_values=<factory>, failures=<factory>, stdouts=<factory>, stderrs=<factory>) [исходный код]

Результаты завершённого запуска процессов, начатых с помощью start_processes(). Возвращается методом PContext.

Обратите внимание на следующее:

  1. Все поля сопоставлены локальным рангам
  2. return_values — заполняется только для функций (не для исполняемых файлов).
  3. stdouts — путь к stdout.log (пустая строка, если перенаправление не выполнялось)
  4. stderrs — путь к stderr.log (пустая строка, если перенаправление не выполнялось)
class torch.distributed.elastic.multiprocessing.api.DefaultLogsSpecs(log_dir=None, redirects=Std.NONE, tee=Std.NONE, local_ranks_filter=None) [исходный код]

Реализация LogsSpecs по умолчанию:

  • log_dir будет создан, если он не существует
  • Создаются вложенные папки для каждой попытки и каждого ранга.
reify(envs) [исходный код]

Для формирования путей назначения журналов используется следующая схема:

  • <log_dir>/<rdzv_run_id>/attempt_<attempt>/<rank>/stdout.log
  • <log_dir>/<rdzv_run_id>/attempt_<attempt>/<rank>/stderr.log
  • <log_dir>/<rdzv_run_id>/attempt_<attempt>/<rank>/error.json
  • <log_dir>/<rdzv_run_id>/attempt_<attempt>/filtered_stdout.log
  • <log_dir>/<rdzv_run_id>/attempt_<attempt>/filtered_stderr.log
Тип возвращаемого значения:

LogsDest

class torch.distributed.elastic.multiprocessing.api.LogsDest(stdouts=<factory>, stderrs=<factory>, tee_stdouts=<factory>, tee_stderrs=<factory>, error_files=<factory>, filtered_stdout=<factory>, filtered_stderr=<factory>) [исходный код]

Для каждого типа журнала содержит сопоставление идентификаторов локальных рангов с путями к файлам.

class torch.distributed.elastic.multiprocessing.api.LogsSpecs(log_dir=None, redirects=Std.NONE, tee=Std.NONE, local_ranks_filter=None) [исходный код]

Определяет обработку журналов и перенаправление для каждого рабочего процесса.

Параметры:
  • log_dir (str | None) – Базовый каталог, в который будут записываться журналы.
  • redirects (Std | dict[int, Std]) – Потоки для перенаправления в файлы. Передайте одно значение перечисления Std, чтобы перенаправить потоки для всех рабочих процессов, или словарь с ключами local_rank для выборочного перенаправления.
  • tee (Std | dict[int, Std]) – Потоки для дублирования в stdout/stderr. Передайте одно значение перечисления Std, чтобы дублировать потоки для всех рабочих процессов, или словарь с ключами local_rank для выборочного дублирования.
abstract reify(envs) [исходный код]

На основе переменных среды формирует пути назначения файлов журналов для каждого локального ранга.

Параметр Envs содержит словарь переменных среды для каждого локального ранга; записи определены в: _start_workers().

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

LogsDest

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

Spec-Zone.ru

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