Многопроцессность
Создано: 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будут содержаться:-
{local_rank}/error.json: если процесс завершился с ошибкой, файл с информацией об ошибке -
{local_rank}/stdout.log: еслиredirect & STDOUT == STDOUT -
{local_rank}/stderr.log: еслиredirect & STDERR == STDERR -
filtered_stdout.log: еслиduplicate_stdout_filtersне пуст -
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) — фильтры для дублируемых журналов стандартного потока ошибок
- Тип возвращаемого значения:
-
Контекст процесса
-
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’stee, в файл, содержащий только строки, совпадающие с _любой_ из строк фильтра. Файл журнала объединяет данные всех рангов, выбранных с помощьюtee. -
duplicate_stderr_filters (list[str] | None) – Если список не пуст, дублирует stderr, указанный в
logs_specs’stee, в файл, содержащий только строки, совпадающие с _любой_ из строк фильтра. Файл журнала объединяет данные всех рангов, выбранных с помощьюtee.
-
duplicate_stdout_filters (list[str] | None) – Если список не пуст, дублирует stdout, указанный в
-
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.Обратите внимание на следующее:
- Все поля сопоставлены локальным рангам
-
return_values— заполняется только для функций (не для исполняемых файлов). -
stdouts— путь к stdout.log (пустая строка, если перенаправление не выполнялось) -
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
- Тип возвращаемого значения:
-
-
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().- Тип возвращаемого значения:
© 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