Spec-Zone.ru › PyTorch 2.14

Пакет распределённых коммуникаций — torch.distributed

Создано: 12 июля 2017 г. | Последнее обновление: 7 августа 2026 г.

Примечание

Краткое описание всех функций, связанных с распределённым обучением, см. в документе «Обзор распределённых вычислений в PyTorch».

torch.distributed.elastic.utils.api.get_env_variable_or_raise(env_name) [исходный код]

Пытается получить переменную среды. Вызывает ValueError, если переменная среды не найдена.

Параметры:

env_name (str) – Имя переменной среды

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

str

torch.distributed.elastic.utils.distributed.get_free_port() [исходный код]

Возвращает неиспользуемый порт на localhost.

Эта функция находит неиспользуемый порт на localhost, открывая сокет для привязки к порту, а затем закрывая его.

Возвращает:

неиспользуемый порт на localhost

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

int

Пример

>>> get_free_port()
63976

Примечание

Порт, возвращаемый функцией get_free_port(), не резервируется, поэтому другой процесс может занять его после возврата этой функции.

torch.distributed.elastic.utils.log_level.get_log_level() [исходный код]

Возвращает уровень ведения журнала по умолчанию для PyTorch.

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

str

torch.distributed.elastic.utils.logging.get_logger(name=None) [исходный код]

Вспомогательная функция для настройки простого регистратора, выполняющего запись в stderr. Уровень ведения журнала считывается из переменной среды LOGLEVEL; если она не задана, используется значение WARNING. Если имя не указано, функция использует имя модуля вызывающего кода.

Параметры:

name (str | None) – Имя регистратора. Если имя не указано, оно определяется по стеку вызовов.

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

Logger

torch.distributed.rendezvous.register_rendezvous_handler(scheme, handler) [исходный код]

Регистрирует новый обработчик rendezvous.

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

Результатом процесса rendezvous является кортеж из общего хранилища «ключ/значение», ранга процесса и общего числа участвующих процессов.

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

Параметры:
  • scheme (str) – Схема URL для идентификации обработчика rendezvous.
  • handler (function) – Обработчик, вызываемый при вызове функции rendezvous() с URL, использующим соответствующую схему. Это должна быть функция-генератор, возвращающая кортеж.
torch.distributed.algorithms.model_averaging.utils.average_parameters(params, process_group) [исходный код]

Вычисляет среднее значение всех заданных параметров.

Для повышения эффективности allreduce все параметры объединяются в один непрерывный буфер. Поэтому требуется дополнительная память размером, равным размеру заданных параметров.

torch.distributed.algorithms.model_averaging.utils.average_parameters_or_parameter_groups(params, process_group) [исходный код]

Вычисляет среднее значение параметров модели или групп параметров оптимизатора.

torch.distributed.algorithms.model_averaging.utils.get_params_to_average(params) [исходный код]

Возвращает список параметров, для которых необходимо вычислить среднее значение.

Эта функция отфильтровывает параметры, для которых не вычисляются градиенты. :param params: Параметры модели или группы параметров оптимизатора.

Бэкенды

torch.distributed поддерживает четыре встроенных бэкенда с разными возможностями. В таблице ниже показано, какие функции доступны для CPU и GPU в каждом бэкенде. Для NCCL под GPU подразумевается CUDA GPU, а для XCCL — XPU GPU.

MPI поддерживает CUDA только в том случае, если используемая для сборки PyTorch реализация поддерживает её.

Бэкенд

gloo

mpi

nccl

xccl

Устройство

CPU

GPU

CPU

GPU

CPU

GPU

CPU

GPU

send

✓

✘

✓

?

✘

✓

✘

✓

recv

✓

✘

✓

?

✘

✓

✘

✓

broadcast

✓

✓

✓

?

✘

✓

✘

✓

all_reduce

✓

✓

✓

?

✘

✓

✘

✓

reduce

✓

✓

✓

?

✘

✓

✘

✓

all_gather

✓

✓

✓

?

✘

✓

✘

✓

gather

✓

✓

✓

?

✘

✓

✘

✓

scatter

✓

✓

✓

?

✘

✓

✘

✓

reduce_scatter

✓

✓

✘

✘

✘

✓

✘

✓

all_to_all

✘

✘

✓

?

✘

✓

✘

✓

barrier

✓

✘

✓

?

✘

✓

✘

✓

Бэкенды, входящие в состав PyTorch

Пакет распределённых вычислений PyTorch поддерживает Linux (стабильная версия), macOS (стабильная версия) и Windows (прототип). По умолчанию для Linux бэкенды Gloo и NCCL собираются и включаются в пакет распределённых вычислений PyTorch (NCCL включается только при сборке с CUDA). MPI — необязательный бэкенд, который можно включить только при сборке PyTorch из исходного кода (например, при сборке PyTorch на хосте с установленным MPI).

Примечание

Начиная с PyTorch v1.8, Windows поддерживает все бэкенды коллективных коммуникаций, кроме NCCL. Если аргумент init_method функции init_process_group() указывает на файл, он должен соответствовать следующей схеме:

  • Локальная файловая система, init_method="file:///d:/tmp/some_file"
  • Общая файловая система, init_method="file://////{machine_name}/{share_folder_name}/some_file"

Как и на платформе Linux, можно включить TcpStore, задав переменные среды MASTER_ADDR и MASTER_PORT.

Какой бэкенд использовать?

Раньше нас часто спрашивали: «Какой бэкенд следует использовать?»

  • Общее правило

    • Для распределённого обучения на CUDA GPU используйте бэкенд NCCL.
    • Для распределённого обучения на XPU GPU используйте бэкенд XCCL.
    • Для распределённого обучения на CPU используйте бэкенд Gloo.
  • Хосты с GPU и соединением InfiniBand

    • Используйте NCCL, поскольку это единственный бэкенд с поддержкой InfiniBand и GPUDirect.
  • Хосты с GPU и соединением Ethernet

    • Используйте NCCL, поскольку сейчас он обеспечивает наилучшую производительность распределённого обучения на GPU, особенно при распределённом обучении в нескольких процессах на одном узле или на нескольких узлах. Если возникнут проблемы с NCCL, используйте Gloo в качестве запасного варианта. (Обратите внимание, что на GPU Gloo работает медленнее NCCL.)
  • Хосты с CPU и соединением InfiniBand

    • Если в InfiniBand включён IP over IB, используйте Gloo, в противном случае — MPI. В будущих выпусках планируется добавить поддержку InfiniBand для Gloo.
  • Хосты с CPU и соединением Ethernet

    • Используйте Gloo, если у вас нет особых причин использовать MPI.

Общие переменные среды

Выбор сетевого интерфейса

По умолчанию бэкенды NCCL и Gloo пытаются автоматически выбрать подходящий сетевой интерфейс. Если автоматически выбранный интерфейс не подходит, его можно переопределить с помощью следующих переменных среды (для соответствующего бэкенда):

  • NCCL_SOCKET_IFNAME, например export NCCL_SOCKET_IFNAME=eth0
  • GLOO_SOCKET_IFNAME, например export GLOO_SOCKET_IFNAME=eth0

При использовании бэкенда Gloo можно указать несколько интерфейсов, разделив их запятыми, например: export GLOO_SOCKET_IFNAME=eth0,eth1,eth2,eth3. Бэкенд будет распределять операции по этим интерфейсам по кругу. Крайне важно, чтобы все процессы указали в этой переменной одинаковое количество интерфейсов.

Другие переменные среды NCCL

Отладка — в случае сбоя NCCL можно задать NCCL_DEBUG=INFO, чтобы выводить подробное предупреждение, а также основные сведения об инициализации NCCL.

Также можно использовать NCCL_DEBUG_SUBSYS, чтобы получить более подробные сведения об определённом аспекте NCCL. Например, NCCL_DEBUG_SUBSYS=COLL выведет журналы коллективных вызовов, что может быть полезно при отладке зависаний, особенно вызванных несовпадением типа коллективной операции или размера сообщения. В случае сбоя обнаружения топологии полезно задать NCCL_DEBUG_SUBSYS=GRAPH, чтобы изучить подробные результаты обнаружения и сохранить их для справки, если потребуется обратиться за помощью к команде NCCL.

Настройка производительности — NCCL автоматически настраивает параметры на основе обнаруженной топологии, чтобы упростить настройку для пользователей. В некоторых системах на основе сокетов пользователи могут попробовать настроить NCCL_SOCKET_NTHREADS и NCCL_NSOCKS_PERTHREAD для увеличения пропускной способности сети сокетов. Эти две переменные среды предварительно настроены NCCL для некоторых облачных провайдеров, например AWS и GCP.

Полный список переменных среды NCCL см. в официальной документации NVIDIA NCCL

Дополнительно настраивать коммуникаторы NCCL можно с помощью torch.distributed.ProcessGroupNCCL.NCCLConfig и torch.distributed.ProcessGroupNCCL.Options. Подробнее о них можно узнать с помощью help (например, help(torch.distributed.ProcessGroupNCCL.NCCLConfig)) в интерпретаторе.

Основы

Пакет torch.distributed обеспечивает поддержку PyTorch и предоставляет примитивы коммуникации для параллельной обработки в нескольких процессах на нескольких вычислительных узлах, работающих на одной или нескольких машинах. Класс torch.nn.parallel.DistributedDataParallel() использует эту функциональность для синхронного распределённого обучения, оборачивая любую модель PyTorch. Это отличается от типов параллелизма, обеспечиваемых пакетом multiprocessing — torch.multiprocessing и torch.nn.DataParallel(): он поддерживает несколько машин, соединённых по сети, и требует от пользователя явного запуска отдельной копии основного обучающего скрипта для каждого процесса.

Даже в синхронном режиме на одной машине torch.distributed или обёртка torch.nn.parallel.DistributedDataParallel() могут иметь преимущества перед другими подходами к параллелизму данных, в том числе перед torch.nn.DataParallel():

  • Каждый процесс поддерживает собственный оптимизатор и выполняет полный шаг оптимизации на каждой итерации. Хотя это может показаться избыточным, поскольку градиенты уже собраны и усреднены между процессами, а значит, одинаковы для каждого процесса, это позволяет избежать рассылки параметров и сократить время передачи тензоров между узлами.
  • Каждый процесс содержит независимый интерпретатор Python, что устраняет дополнительные накладные расходы интерпретатора и «конкуренцию за GIL», возникающую при управлении несколькими потоками выполнения, репликами модели или GPU из одного процесса Python. Это особенно важно для моделей, активно использующих среду выполнения Python, в том числе моделей с рекуррентными слоями или множеством небольших компонентов.

Инициализация

Перед вызовом любых других методов пакет необходимо инициализировать с помощью функции torch.distributed.init_process_group() или torch.distributed.device_mesh.init_device_mesh(). Обе функции блокируют выполнение, пока не присоединятся все процессы.

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

Инициализация не является потокобезопасной. Создание группы процессов следует выполнять из одного потока, чтобы предотвратить несогласованное назначение «UUID» для разных рангов и гонки во время инициализации, которые могут привести к зависаниям.

torch.distributed.is_available() [исходный код]

Возвращает True, если пакет распределённых вычислений доступен.

В противном случае torch.distributed не предоставляет другие API. В настоящее время torch.distributed доступен в Linux, macOS и Windows. Установите USE_DISTRIBUTED=1, чтобы включить его при сборке PyTorch из исходного кода. В настоящее время значение по умолчанию — USE_DISTRIBUTED=1 для Linux и Windows и USE_DISTRIBUTED=0 для macOS.

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

bool

torch.distributed.init_process_group(backend=None, init_method=None, timeout=None, world_size=-1, rank=-1, store=None, group_name='', pg_options=None, device_id=None, _ranks=None, enable_reconfigure=False) [исходный код]

Инициализирует группу процессов распределённых вычислений по умолчанию.

Также инициализирует пакет распределённых вычислений.

Существует два основных способа инициализировать группу процессов:
  1. Явно указать store, rank и world_size.
  2. Указать init_method (строку URL), которая определяет, где и как обнаруживать узлы. При необходимости укажите rank и world_size либо закодируйте все необходимые параметры в URL и не указывайте их отдельно.

Если не указан ни один из вариантов, предполагается, что init_method имеет значение «env://».

Параметры:
  • backend (str или Backend, необязательный) – Используемый бэкенд. В зависимости от конфигурации сборки допустимы значения mpi, gloo, nccl, ucc, xccl или бэкенд, зарегистрированный сторонним плагином. Начиная с версии 2.6, если backend не задан, c10d использует бэкенд, зарегистрированный для типа устройства, указанного в именованном аргументе device_id (если он задан). На данный момент известны следующие регистрации по умолчанию: nccl для cuda, gloo для cpu, xccl для xpu. Если не задан ни backend, ни device_id, c10d определит ускоритель на компьютере во время выполнения и использует бэкенд, зарегистрированный для обнаруженного ускорителя (или cpu). Это поле можно задать строкой в нижнем регистре (например, "gloo"); к нему также можно обратиться через атрибуты Backend (например, Backend.GLOO). Если используется несколько процессов на одном компьютере с бэкендом nccl, каждый процесс должен иметь исключительный доступ ко всем используемым им графическим процессорам. Совместное использование графических процессоров несколькими процессами может привести к взаимной блокировке или некорректному использованию NCCL. Бэкенд ucc является экспериментальным. Бэкенд устройства по умолчанию можно узнать с помощью get_default_backend_for_device().
  • init_method (str, необязательный) – URL, определяющий способ инициализации группы процессов. Если не заданы init_method или store, используется значение по умолчанию «env://». Взаимоисключающий с store параметр.
  • world_size (int, необязательный) – Число процессов, участвующих в задаче. Обязателен, если задан store.
  • rank (int, необязательный) – Ранг текущего процесса (это должно быть число от 0 до world_size-1). Обязателен, если задан store.
  • store (Store, необязательный) – Хранилище «ключ/значение», доступное всем рабочим процессам и используемое для обмена информацией о подключении и адресах. Взаимоисключающий с init_method параметр.
  • timeout (timedelta, необязательный) – Тайм-аут операций, выполняемых в группе процессов. Значение по умолчанию составляет 10 минут для NCCL и 30 минут для других бэкендов. По истечении этого времени коллективные операции асинхронно прерываются, а процесс аварийно завершается. Это необходимо, поскольку выполнение CUDA асинхронно, и дальнейшее выполнение пользовательского кода становится небезопасным: неудачные асинхронные операции NCCL могут привести к тому, что последующие операции CUDA будут выполняться с повреждёнными данными. Если задан TORCH_NCCL_BLOCKING_WAIT, процесс заблокируется и будет ожидать истечения этого тайм-аута.
  • group_name (str, необязательный, устарел) – Имя группы. Этот аргумент игнорируется.
  • pg_options (ProcessGroupOptions, необязательный) – Параметры группы процессов, определяющие дополнительные настройки, которые нужно передать при создании конкретных групп процессов. На данный момент поддерживается только ProcessGroupNCCL.Options для бэкенда nccl. Можно указать is_high_priority_stream, чтобы бэкенд nccl мог использовать потоки cuda с высоким приоритетом, когда ожидают выполнения вычислительные ядра. О других доступных параметрах конфигурации nccl см. https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/api/types.html#ncclconfig-t
  • device_id (torch.device | int, необязательный) – Конкретное устройство, на котором будет работать этот процесс; это позволяет применять оптимизации, зависящие от бэкенда. Сейчас этот параметр имеет два эффекта только при использовании NCCL: коммуникатор создаётся немедленно (сразу вызывается ncclCommInit* вместо обычного отложенного вызова), а вложенные группы при возможности используют ncclCommSplit, чтобы избежать лишних затрат на создание группы. Это поле также позволяет заранее узнать об ошибках инициализации NCCL. Если задан int, API предполагает, что будет использоваться тип ускорителя, определённый во время компиляции.
  • _ranks (list[int] | None) – Ранги в группе процессов. Если задан, имя группы процессов будет хешем всех рангов в группе.
  • enable_reconfigure (bool, необязательный) – Если задано True, создаёт бэкенд в режиме переконфигурирования (отказоустойчивости). Коммуникатор не инициализируется, пока не будет вызван ProcessGroup.reconfigure(). Бэкенды, не поддерживающие переконфигурирование, игнорируют этот флаг. Значение по умолчанию — False.

Примечание

Для включения backend == Backend.MPI PyTorch необходимо собрать из исходного кода в системе с поддержкой MPI.

Примечание

Поддержка нескольких бэкендов является экспериментальной. В настоящее время, если бэкенд не указан, будут созданы бэкенды gloo и nccl. Бэкенд gloo будет использоваться для коллективных операций с тензорами CPU, а бэкенд nccl — для коллективных операций с тензорами CUDA. Пользовательский бэкенд можно указать строкой в формате «<device_type>:<backend_name>,<device_type>:<backend_name>», например «cpu:gloo,cuda:custom_backend».

torch.distributed.device_mesh.init_device_mesh(device_type, mesh_shape, *, mesh_dim_names=None, backend_override=None) [исходный код]

Инициализирует DeviceMesh на основе параметров device_type, mesh_shape и mesh_dim_names.

Создаёт DeviceMesh с n-мерным расположением в виде массива, где n — длина mesh_shape. Если задано mesh_dim_names, каждому измерению присваивается метка mesh_dim_names[i].

Примечание

init_device_mesh следует модели программирования SPMD, то есть одна и та же программа PyTorch на Python выполняется на всех процессах/рангах кластера. Убедитесь, что mesh_shape (размерности nD-массива, описывающего расположение устройств) одинаково для всех рангов. Несогласованное значение mesh_shape может привести к зависанию.

Примечание

Если группа процессов не найдена, init_device_mesh автоматически инициализирует группу или группы процессов, необходимые для распределённого взаимодействия.

Параметры:
  • device_type (str) – Тип устройства в сетке. В настоящее время поддерживаются: «cpu», «cuda/cuda-like», «xpu». Указывать тип устройства с индексом графического процессора, например «cuda:0», нельзя.
  • mesh_shape (Tuple[int]) – Кортеж, задающий размерности многомерного массива, описывающего расположение устройств.
  • mesh_dim_names (tuple[str, ...], необязательный) – Кортеж имён размерностей сетки, назначаемых каждой размерности многомерного массива, описывающего расположение устройств. Его длина должна совпадать с длиной mesh_shape. Каждая строка в mesh_dim_names должна быть уникальной.
  • backend_override (Dict[int | str, tuple[str, Options] | str | Options], необязательный) – Переопределения для некоторых или всех ProcessGroups, которые будут созданы для каждой размерности сетки. Ключом может быть индекс размерности или её имя (если задан параметр mesh_dim_names). Значением может быть кортеж, содержащий имя бэкенда и его параметры, либо только один из этих компонентов (в этом случае для другого компонента используется значение по умолчанию).
Возвращает:

Объект DeviceMesh, представляющий расположение устройств.

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

DeviceMesh

Пример:

>>> from torch.distributed.device_mesh import init_device_mesh
>>>
>>> mesh_1d = init_device_mesh("cuda", mesh_shape=(8,))
>>> mesh_2d = init_device_mesh("cuda", mesh_shape=(2, 8), mesh_dim_names=("dp", "tp"))
torch.distributed.is_initialized() [исходный код]

Проверяет, инициализирована ли группа процессов по умолчанию.

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

bool

torch.distributed.is_mpi_available() [исходный код]

Проверяет, доступен ли бэкенд MPI.

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

bool

torch.distributed.is_nccl_available() [исходный код]

Проверяет, доступен ли бэкенд NCCL.

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

bool

torch.distributed.is_gloo_available() [исходный код]

Проверяет, доступен ли бэкенд Gloo.

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

bool

torch.distributed.distributed_c10d.is_xccl_available() [исходный код]

Проверяет, доступен ли бэкенд XCCL.

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

bool

torch.distributed.distributed_c10d.batch_isend_irecv(p2p_op_list) [исходный код]

Асинхронно отправляет или принимает пакет тензоров и возвращает список запросов.

Выполняет каждую из операций в p2p_op_list и возвращает соответствующие запросы. В настоящее время поддерживаются бэкенды NCCL, Gloo и UCC.

Параметры:

p2p_op_list (list[P2POp]) –

Список операций точка-точка (тип каждого оператора — torch.distributed.P2POp). Все операции в списке рассматриваются как единый пакет: относительный порядок отправок и приёмов не имеет значения и не приводит к взаимным блокировкам. Например, [recv_from_1, send_to_1] эквивалентно [send_to_1, recv_from_1].

Порядок имеет значение для нескольких операций в одном направлении между одной и той же парой рангов: если ранг 0 вызывает [send_to_1(A), send_to_1(B)], ранг 1 должен вызвать [recv_from_0(A), recv_from_0(B)] в том же порядке, чтобы сопоставить правильные тензоры.

Возвращает:

Список объектов распределённых запросов, возвращённых соответствующими вызовами операций из op_list.

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

list[Work]

Примеры

>>> send_tensor = torch.arange(2, dtype=torch.float32) + 2 * rank
>>> recv_tensor = torch.randn(2, dtype=torch.float32)
>>> send_op = dist.P2POp(dist.isend, send_tensor, (rank + 1) % world_size)
>>> recv_op = dist.P2POp(
...     dist.irecv, recv_tensor, (rank - 1 + world_size) % world_size
... )
>>> reqs = batch_isend_irecv([send_op, recv_op])
>>> for req in reqs:
>>>     req.wait()
>>> recv_tensor
tensor([2, 3])     # Rank 0
tensor([0, 1])     # Rank 1

Примечание

Обратите внимание: при использовании этого API с бэкендом NCCL PG необходимо установить текущее устройство GPU с помощью torch.cuda.set_device, иначе возможны неожиданные зависания.

Кроме того, если этот вызов API является первой коллективной операцией в group, переданном в dist.P2POp, в вызове API должны участвовать все ранги group; в противном случае поведение не определено. Если этот вызов API не является первой коллективной операцией в group, разрешены пакетные операции P2P, в которых участвует только подмножество рангов group.

torch.distributed.distributed_c10d.destroy_process_group(group=None) [исходный код]

Уничтожает указанную группу процессов и деинициализирует пакет распределённых вычислений.

Параметры:

group (ProcessGroup, необязательный) – Группа процессов, которую нужно уничтожить. Если указана group.WORLD, будут уничтожены все группы процессов, включая группу по умолчанию.

torch.distributed.distributed_c10d.is_backend_available(backend) [исходный код]

Проверяет доступность бэкенда.

Проверяет, доступен ли указанный бэкенд; поддерживаются встроенные бэкенды и сторонние бэкенды через функцию Backend.register_backend.

Параметры:

backend (str) – Имя бэкенда.

Возвращает:

Возвращает true, если бэкенд доступен, и false в противном случае.

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

bool

torch.distributed.distributed_c10d.irecv(tensor, src=None, group=None, tag=0, group_src=None) [исходный код]

Асинхронно принимает тензор.

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

tag не поддерживается бэкендом NCCL.

В отличие от recv, который блокирует выполнение, irecv допускает совпадение рангов src и dst, то есть приём от самого себя.

Параметры:
  • tensor (Tensor) – Тензор для заполнения полученными данными.
  • src (int, необязательный) – Ранг источника в глобальной группе процессов (независимо от аргумента group). Если не указан, данные будут получены от любого процесса.
  • group (ProcessGroup, необязательный) – Группа процессов, с которой нужно работать. Если значение равно None, используется группа процессов по умолчанию.
  • tag (int, необязательный) – Тег для сопоставления recv с удалённой отправкой.
  • group_src (int, необязательный) – Ранг назначения в group. Нельзя одновременно указывать src и group_src.
Возвращает:

Объект распределённого запроса. None, если процесс не входит в группу.

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

Work | None

torch.distributed.distributed_c10d.is_gloo_available() [исходный код]

Проверяет, доступен ли бэкенд Gloo.

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

bool

torch.distributed.distributed_c10d.is_initialized() [исходный код]

Проверяет, инициализирована ли группа процессов по умолчанию.

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

bool

torch.distributed.distributed_c10d.is_mpi_available() [исходный код]

Проверяет, доступен ли бэкенд MPI.

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

bool

torch.distributed.distributed_c10d.is_nccl_available() [исходный код]

Проверяет, доступен ли бэкенд NCCL.

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

bool

torch.distributed.distributed_c10d.is_torchelastic_launched() [исходный код]

Проверяет, был ли этот процесс запущен с помощью torch.distributed.elastic (также известного как torchelastic).

Наличие переменной среды TORCHELASTIC_RUN_ID используется как косвенный признак того, что текущий процесс был запущен с помощью torchelastic. Это достаточно надёжный признак, поскольку TORCHELASTIC_RUN_ID соответствует идентификатору rendezvous, который всегда является ненулевым значением и указывает идентификатор задачи для обнаружения узлов.

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

bool

torch.distributed.distributed_c10d.is_ucc_available() [исходный код]

Проверяет, доступен ли бэкенд UCC.

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

bool

torch.distributed.is_torchelastic_launched() [исходный код]

Проверяет, был ли этот процесс запущен с помощью torch.distributed.elastic (также известного как torchelastic).

Наличие переменной среды TORCHELASTIC_RUN_ID используется как косвенный признак того, что текущий процесс был запущен с помощью torchelastic. Это достаточно надёжный признак, поскольку TORCHELASTIC_RUN_ID соответствует идентификатору rendezvous, который всегда является ненулевым значением и указывает идентификатор задачи для обнаружения узлов.

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

bool

torch.distributed.get_default_backend_for_device(device) [исходный код]

Возвращает бэкенд по умолчанию для указанного устройства.

Параметры:

device (Union[str, torch.device]) – Устройство, для которого нужно получить бэкенд по умолчанию.

Возвращает:

Бэкенд по умолчанию для указанного устройства в виде строки в нижнем регистре.

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

str

В настоящее время поддерживаются три метода инициализации:

Инициализация TCP

Существует два способа инициализации с помощью TCP. Оба требуют сетевого адреса, доступного для всех процессов, и желаемого world_size. Первый способ требует указать адрес, принадлежащий процессу с рангом 0. Этот метод инициализации требует, чтобы все процессы вручную указывали свои ранги.

Обратите внимание, что в последней версии пакета distributed многоадресная рассылка больше не поддерживается. group_name также объявлен устаревшим.

import torch.distributed as dist

# Use address of one of the machines
dist.init_process_group(backend, init_method='tcp://10.1.1.20:23456',
                        rank=args.rank, world_size=4)

Инициализация с общей файловой системой

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

Обратите внимание, что в последней версии пакета distributed автоматическое назначение рангов больше не поддерживается, а group_name также объявлен устаревшим.

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

Этот метод предполагает, что файловая система поддерживает блокировки с помощью fcntl — большинство локальных систем и NFS поддерживают их.

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

Этот метод всегда создаёт файл и старается удалить его в конце выполнения программы. Другими словами, для успешной инициализации с помощью метода file init каждый раз нужен новый пустой файл. Повторное использование файла, применённого при предыдущей инициализации (если он не был удалён), приведёт к непредсказуемому поведению и часто может вызвать взаимоблокировки и сбои. Поэтому, даже несмотря на то что этот метод постарается удалить файл, если автоматическое удаление не удастся, вы должны самостоятельно убедиться, что файл удалён в конце обучения, чтобы он не использовался повторно при следующем запуске. Это особенно важно, если вы планируете несколько раз вызывать init_process_group() с одним и тем же именем файла. Другими словами, если файл не удалён и вы снова вызываете init_process_group() для этого файла, следует ожидать сбоев. Общее правило: каждый раз перед вызовом init_process_group() убедитесь, что файл не существует или пуст.

import torch.distributed as dist

# rank should always be specified
dist.init_process_group(backend, init_method='file:///mnt/nfs/sharedfile',
                        world_size=4, rank=args.rank)

Инициализация с помощью переменных среды

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

  • MASTER_PORT — обязательно; должен быть свободным портом на машине с рангом 0
  • MASTER_ADDR — обязательно (кроме процесса с рангом 0); адрес узла с рангом 0
  • WORLD_SIZE — обязательно; можно задать здесь или при вызове функции инициализации
  • RANK — обязательно; можно задать здесь или при вызове функции инициализации

Машина с рангом 0 будет использоваться для настройки всех соединений.

Это метод по умолчанию, поэтому init_method указывать не нужно (или можно задать значение env://).

Сокращение времени инициализации

  • TORCH_GLOO_LAZY_INIT — устанавливает соединения по мере необходимости, а не использует полную сетку, что может значительно сократить время инициализации для операций, не относящихся к all2all.

После инициализации

После выполнения torch.distributed.init_process_group() можно использовать следующие функции. Чтобы проверить, была ли уже инициализирована группа процессов, используйте torch.distributed.is_initialized().

class torch.distributed.Backend(name) [исходный код]

Класс, подобный перечислению, для бэкендов.

Доступные бэкенды: GLOO, NCCL, UCC, MPI, XCCL, FAKE и другие зарегистрированные бэкенды.

Значения этого класса — строки в нижнем регистре, например "gloo". К ним можно обращаться как к атрибутам, например Backend.NCCL.

Этот класс можно вызывать напрямую для разбора строки: например, Backend(backend_str) проверит, является ли backend_str допустимым значением, и в этом случае вернёт разобранную строку в нижнем регистре. Также принимаются строки в верхнем регистре: например, Backend("GLOO") возвращает "gloo".

Примечание

Элемент Backend.UNDEFINED существует, но используется только как начальное значение некоторых полей. Пользователям не следует использовать его напрямую или полагаться на его наличие.

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

str

classmethod register_backend(name, func, extended_api=False, devices=None, *, _backend_type=None) [исходный код]

Регистрирует новый бэкенд с указанным именем и функцией создания экземпляра.

Этот метод класса используется расширениями ProcessGroup сторонних разработчиков для регистрации новых бэкендов.

Параметры:
  • name (str) – Имя бэкенда расширения ProcessGroup. Оно должно совпадать с именем в init_process_group().
  • func (function) – Функция-обработчик, создающая экземпляр бэкенда. Функция должна быть реализована в расширении бэкенда и принимать четыре аргумента, включая store, rank, world_size и timeout.
  • extended_api (bool, optional) – Поддерживает ли бэкенд расширенную структуру аргументов. Значение по умолчанию: False. Если задано значение True, бэкенд получит экземпляр c10d::DistributedBackendOptions и объект параметров группы процессов, определённый реализацией бэкенда.
  • devices (str or list of str, optional) – Тип устройства, поддерживаемый этим бэкендом, например «cpu», «cuda» и т. д. Если задано значение None, предполагается поддержка CPU и текущего ускорителя.

Примечание

Поддержка сторонних бэкендов является экспериментальной и может измениться.

torch.distributed.get_backend(group=None) [исходный код]

Возвращает бэкенд указанной группы процессов.

Параметры:

group (ProcessGroup, optional) – Группа процессов, с которой нужно работать. По умолчанию используется общая основная группа процессов. Если указана другая конкретная группа, вызывающий процесс должен входить в group.

Возвращает:

Бэкенд указанной группы процессов в виде строки в нижнем регистре.

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

str

torch.distributed.get_backend_impl(group=None, device=None) → torch._C._distributed_c10d.Backend [исходный код]

Возвращает базовую реализацию бэкенда указанной группы процессов.

Это позволяет пользовательским бэкендам предоставлять специфичные для бэкенда методы для экспериментальных целей.

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

Этот API обходит трассировку torch.compile и другие перехватчики. Методы бэкенда являются экспериментальными и могут быть изменены без предупреждения.

Параметры:
  • group (str or ProcessGroup, optional) – Группа процессов или имя группы процессов, с которой нужно работать. По умолчанию используется общая основная группа процессов. Если указана другая конкретная группа, вызывающий процесс должен входить в group.
  • device (torch.device, optional) – Устройство, используемое для выбора реализации бэкенда. Если задано значение None, возвращается бэкенд, связанный с устройством, или единственный бэкенд, общий для зарегистрированных устройств. Значение по умолчанию: None.
Возвращает:

Реализацию бэкенда для указанной группы процессов и устройства.

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

torch._C._distributed_c10d.Backend

torch.distributed.get_backend_config(group=None) [исходный код]

Возвращает конфигурацию бэкенда указанной группы процессов.

Параметры:

group (ProcessGroup, optional) – Группа процессов, с которой нужно работать. По умолчанию используется общая основная группа процессов. Если указана другая конкретная группа, вызывающий процесс должен входить в group.

Возвращает:

Конфигурацию бэкенда указанной группы процессов в виде строки в нижнем регистре.

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

str

torch.distributed.get_rank(group=None) [исходный код]

Возвращает ранг текущего процесса в указанной group, в противном случае — значение по умолчанию.

Ранг — это уникальный идентификатор, назначаемый каждому процессу в распределённой группе процессов. Ранги всегда являются последовательными целыми числами в диапазоне от 0 до world_size.

Параметры:

group (ProcessGroup, optional) – Группа процессов, с которой нужно работать. Если значение None, используется группа процессов по умолчанию.

Возвращает:

Ранг процесса в группе; -1, если процесс не входит в группу

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

int

torch.distributed.get_world_size(group=None) [исходный код]

Возвращает количество процессов в текущей группе процессов.

Параметры:

group (ProcessGroup, optional) – Группа процессов, с которой нужно работать. Если значение None, используется группа процессов по умолчанию.

Возвращает:

Размер мира группы процессов; -1, если процесс не входит в группу

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

int

torch.distributed.get_debug_level() → torch._C._distributed_c10d.DebugLevel

Возвращает уровень отладки пакета torch.distributed.

torch.distributed.get_node_local_rank(fallback_rank=None) [исходный код]

Возвращает локальный ранг текущего процесса относительно узла.

Семантически это полезно для сопоставления процессов с устройствами. Например, на узле с 8 ускорителями можно использовать локальный ранг узла, чтобы определить, к какому устройству-ускорителю привязать процесс.

На практике фактическое назначение локальных рангов узла выполняется средством запуска процессов вне PyTorch и передаётся через переменную среды LOCAL_RANK.

Torchrun автоматически задаёт LOCAL_RANK, однако другие средства запуска могут этого не делать. Если LOCAL_RANK не указана, этот API использует переданный аргумент kwarg ‘fallback_rank’, если он задан, иначе выдаёт ошибку. Это позволяет написать приложение, которое работает в контексте одного или нескольких устройств без ошибок.

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

int

torch.distributed.get_pg_count() [исходный код]

Возвращает количество групп процессов.

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

int

torch.distributed.set_timeout(timeout, group=None) [исходный код]

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

Это переопределяет время ожидания, настроенное при создании группы (с помощью init_process_group() или new_group()). Новое время ожидания передаётся каждому бэкенду, зарегистрированному в group, — например, бэкендам Gloo и NCCL группы, охватывающей устройства CPU и CUDA. Бэкенды, не поддерживающие изменение времени ожидания (например, MPI и UCC), выводят предупреждение и оставляют прежнее значение.

Параметры:
  • timeout (timedelta) – Время ожидания для операций, выполняемых в группе процессов. Значение по умолчанию, заданное при инициализации, составляет 10 минут для NCCL и 30 минут для других бэкендов. Для NCCL это период, по истечении которого коллективные операции асинхронно прерываются, а процесс аварийно завершается. Это необходимо, поскольку выполнение CUDA происходит асинхронно и после сбоя асинхронной операции NCCL небезопасно продолжать выполнение пользовательского кода (последующие операции CUDA могут работать с повреждёнными данными). Если задано TORCH_NCCL_BLOCKING_WAIT, процесс вместо этого блокируется и ждёт истечения этого времени ожидания.
  • group (ProcessGroup, optional) – Группа процессов, с которой нужно работать. По умолчанию используется общая основная группа процессов. Если указана другая конкретная группа, вызывающий процесс должен входить в group.
Вызывает исключения:

ValueError – Если вызывающий процесс не входит в group.

Возвращает:

None

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

None

Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Shorten the timeout of the default process group to 30 seconds.
>>> dist.set_timeout(timedelta(seconds=30))

Отказоустойчивая реконфигурация

torch.distributed.distributed_c10d._supports_reconfigure(group=None) [исходный код]

Возвращает, поддерживает ли group API отказоустойчивости на основе реконфигурации.

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

Этот API является экспериментальным и может измениться или быть удалён.

Параметры:

group (ProcessGroup, optional) – Группа процессов, для которой выполняется проверка. Если задано None, используется группа процессов по умолчанию.

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

bool

torch.distributed.distributed_c10d._get_reconfigure_handle(group=None) [исходный код]

Возвращает непрозрачный дескриптор реконфигурации для group.

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

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

Этот API является экспериментальным и может измениться или быть удалён.

Параметры:

group (ProcessGroup, optional) – Группа процессов, для которой выполняется проверка. Если задано None, используется группа процессов по умолчанию.

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

str

torch.distributed.distributed_c10d._reconfigure(uuid, handles, group=None, timeout=None, hints=None) [исходный код]

Реконфигурирует group с новым набором узлов для обеспечения отказоустойчивости.

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

Этот API является экспериментальным и может измениться или быть удалён.

Параметры:
  • uuid (int) – Уникальный идентификатор этого экземпляра средства связи. Передавайте новое значение при каждой (повторной) инициализации.
  • handles (set[str] or list[str]) – По одному дескриптору реконфигурации для каждого узла, возвращаемому функцией _get_reconfigure_handle(). Список назначает ранги в соответствии с позициями элементов; множество позволяет бэкенду выбрать назначение рангов.
  • group (ProcessGroup, optional) – Группа процессов, которую нужно реконфигурировать. Если задано None, используется группа процессов по умолчанию.
  • timeout (timedelta, optional) – Допустимое время реконфигурации до возникновения сбоя. None использует время ожидания по умолчанию для бэкенда.
  • hints (dict[str, str], optional) – Конфигурация, специфичная для бэкенда.
Возвращает:

Дескриптор асинхронной операции реконфигурации.

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

Work

Пример::
>>> dist.init_process_group("gloo", enable_reconfigure=True)
>>> # Every peer receives the same fresh UUID and rank-ordered handles.
>>> uuid, handles = rendezvous_reconfigure(dist._get_reconfigure_handle())
>>> dist._reconfigure(uuid=uuid, handles=handles).wait()

Завершение работы

Важно освободить ресурсы при выходе, вызвав destroy_process_group().

Самый простой вариант — уничтожить каждую группу процессов и бэкенд, вызвав destroy_process_group() со значением None по умолчанию для аргумента group в той части сценария обучения, где обмен данными уже не требуется, обычно ближе к концу main(). Вызов нужно выполнить один раз для каждого процесса обучения, а не на внешнем уровне средства запуска процессов.

Если destroy_process_group() не будет вызвана всеми рангами группы процессов в течение времени ожидания, особенно когда приложение использует несколько групп процессов, например для параллелизма по N измерениям, при завершении возможны зависания. Это происходит потому, что деструктор ProcessGroupNCCL вызывает ncclCommAbort, который должен вызываться коллективно, однако порядок вызова деструкторов ProcessGroupNCCL при сборке мусора Python не детерминирован. Вызов destroy_process_group() помогает обеспечить согласованный порядок вызова ncclCommAbort на всех рангах и избежать вызова ncclCommAbort из деструктора ProcessGroupNCCL.

Повторная инициализация

destroy_process_group также можно использовать для уничтожения отдельных групп процессов. Один из вариантов применения — отказоустойчивое обучение, при котором группа процессов может быть уничтожена, а затем заново инициализирована во время выполнения. В этом случае крайне важно синхронизировать процессы обучения каким-либо способом, отличным от примитивов torch.distributed, после вызова destroy и до последующей инициализации. Сейчас такое поведение не поддерживается и не тестируется из-за сложности обеспечения этой синхронизации и считается известной проблемой. Если это препятствует вашему варианту использования, создайте issue или RFC на GitHub.

Группы

По умолчанию коллективные операции выполняются в группе по умолчанию (также называемой world) и требуют, чтобы все процессы вызвали распределённую функцию. Однако для некоторых рабочих нагрузок может быть полезно более детальное управление обменом данными. Для этого используются распределённые группы. Функцию new_group() можно использовать для создания новых групп с произвольными подмножествами всех процессов. Она возвращает непрозрачный дескриптор группы, который можно передавать в качестве аргумента group всем коллективным операциям (коллективные операции — это распределённые функции для обмена данными в рамках определённых общеизвестных шаблонов программирования).

torch.distributed.new_group(ranks=None, timeout=None, backend=None, pg_options=None, use_local_synchronization=False, group_desc=None, device_id=None, sort_ranks=True) [исходный код]

Создаёт новую распределённую группу.

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

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

Безопасное параллельное использование: при использовании нескольких групп процессов с бэкендом NCCL пользователь должен обеспечить глобально согласованный порядок выполнения коллективных операций на всех рангах.

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

При использовании асинхронных вариантов API обмена данными torch.distributed возвращается объект работы, а ядро обмена данными помещается в очередь отдельного потока CUDA, что позволяет выполнять обмен данными и вычисления параллельно. После вызова одной или нескольких асинхронных операций в одной группе процессов их необходимо синхронизировать с другими потоками CUDA, вызвав work.wait(), прежде чем использовать другую группу процессов.

Подробнее см. Using multiple NCCL communicators concurrently <https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/communicators.html#using-multiple-nccl-communicators-concurrently>.

Параметры:
  • ranks (list[int]) – Список рангов участников группы. Если None, будут использованы все ранги. Значение по умолчанию — None.
  • timeout (timedelta, optional) – сведения и значение по умолчанию см. в init_process_group.
  • backend (str or Backend, optional) – Используемый бэкенд. В зависимости от конфигурации сборки допустимы значения gloo и nccl. По умолчанию используется тот же бэкенд, что и для глобальной группы. Это поле следует задавать строкой в нижнем регистре (например, "gloo"); также к нему можно обратиться через атрибуты Backend (например, Backend.GLOO). Если передано None, будет использоваться бэкенд, соответствующий группе процессов по умолчанию. Значение по умолчанию — None.
  • pg_options (ProcessGroupOptions, optional) – параметры группы процессов, определяющие, какие дополнительные параметры следует передать при создании конкретных групп процессов. Например, для бэкенда nccl можно указать is_high_priority_stream, чтобы группа процессов могла использовать потоки CUDA с высоким приоритетом. Другие доступные параметры конфигурации nccl см. в документации https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/api/types.html#ncclconfig-tuse_local_synchronization (bool, optional): выполнить локальный для группы барьер в конце создания группы процессов. Отличие этого режима в том, что ранги, не входящие в группу, не должны вызывать API и не участвуют в барьере.
  • group_desc (str, optional) – строка с описанием группы процессов.
  • device_id (torch.device, optional) – одно конкретное устройство, к которому будет «привязан» этот процесс. Если задано это поле, вызов new_group попытается немедленно инициализировать для устройства бэкенд обмена данными.
  • sort_ranks (bool, optional) – Если True (значение по умолчанию), перед созданием группы отсортировать список ranks. Если False, сохранить переданный вызывающей стороной порядок рангов, чтобы положение в списке ranks определяло ранг в группе. Все процессы должны передать идентичный список ranks. Значение по умолчанию — True.
Возвращает:

Дескриптор распределённой группы, который можно передавать вызовам коллективных операций, либо GroupMember.NON_GROUP_MEMBER, если ранг не входит в ranks.

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

ProcessGroup | Literal[-100]

Примечание. use_local_synchronization не работает с MPI.

Примечание. Хотя use_local_synchronization=True может значительно повысить производительность в крупных кластерах с небольшими группами процессов, следует соблюдать осторожность: этот параметр меняет поведение кластера, поскольку ранги, не входящие в группу, не участвуют в барьере group().

Примечание. use_local_synchronization=True может привести к взаимоблокировкам, если каждый ранг создаёт несколько перекрывающихся групп процессов. Чтобы этого избежать, убедитесь, что все ранги соблюдают один и тот же глобальный порядок создания.

Примечание. Если включён TorchComms (torch.distributed.config.use_torchcomms / TORCH_DISTRIBUTED_USE_TORCHCOMMS=1), подгруппы создаются напрямую через new_comm TorchComms (обычный путь _new_group_with_tag), а не путём разделения родительского коммуникатора. Передайте backend="nccl-lazy", чтобы создать лениво инициализируемую группу для каждого узла (отдельный comm и поток для каждого узла отправки/получения, как в ProcessGroupNCCL), что позволяет параллельно выполнять операции P2P с разными узлами; передайте бэкенд по умолчанию / "nccl" для создания группы с немедленной инициализацией.

torch.distributed.distributed_c10d.shrink_group(ranks_to_exclude, group=None, shrink_flags=0, pg_options=None) [исходный код]

Уменьшает группу процессов, исключая указанные ранги.

Создаёт и возвращает новую, меньшую группу процессов, включающую только те ранги исходной группы, которые не указаны в списке ranks_to_exclude.

Параметры:
  • ranks_to_exclude (List[int]) – Список рангов исходной группы group, которые нужно исключить из новой группы.
  • group (ProcessGroup, optional) – Группа процессов, которую нужно уменьшить. Если None, используется группа процессов по умолчанию. Значение по умолчанию — None.
  • shrink_flags (int, optional) – Флаги управления уменьшением группы. Может быть SHRINK_DEFAULT (по умолчанию) или SHRINK_ABORT. SHRINK_ABORT попытается завершить текущие операции в родительском коммуникаторе перед уменьшением группы. Значение по умолчанию — SHRINK_DEFAULT.
  • pg_options (ProcessGroupOptions, optional) – Специфичные для бэкенда параметры, применяемые к уменьшенной группе процессов. Если они заданы, бэкенд использует их при создании новой группы. Если параметры не указаны, новая группа наследует значения по умолчанию от родительской.
Возвращает:

новую группу, состоящую из оставшихся рангов. Если была уменьшена группа по умолчанию, возвращённая группа становится новой группой по умолчанию.

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

ProcessGroup

Вызывает исключение:
  • TypeError – если бэкенд группы не поддерживает уменьшение.
  • ValueError – если ranks_to_exclude недопустим (пуст, содержит выходящие за допустимый диапазон значения,
  • дубликаты или исключает все ранги). –
  • RuntimeError – если исключённый ранг вызывает эту функцию или бэкенд
  • не может выполнить операцию. –

Примечания

  • Вызывать эту функцию должны только ранги, не подлежащие исключению; исключённые ранги не должны участвовать в операции уменьшения группы.
  • Уменьшение группы по умолчанию уничтожает все остальные группы процессов, поскольку переназначение рангов приводит к несогласованности.
torch.distributed.get_group_rank(group, global_rank) [исходный код]

Преобразует глобальный ранг в ранг группы.

global_rank должен входить в group, иначе будет вызвано исключение RuntimeError.

Параметры:
  • group (ProcessGroup) – ProcessGroup, для которой нужно найти относительный ранг.
  • global_rank (int) – Глобальный ранг, который нужно получить.
Возвращает:

Ранг группы для global_rank относительно group

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

int

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

torch.distributed.get_global_rank(group, group_rank) [исходный код]

Преобразует ранг группы в глобальный ранг.

group_rank должен входить в group, иначе будет вызвано исключение RuntimeError.

Параметры:
  • group (ProcessGroup) – ProcessGroup, по которой нужно найти глобальный ранг.
  • group_rank (int) – Ранг группы, который нужно получить.
Возвращает:

Глобальный ранг для group_rank относительно group

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

int

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

torch.distributed.get_process_group_ranks(group) [исходный код]

Получает все ранги, связанные с group.

Параметры:

group (Optional[ProcessGroup]) – ProcessGroup, из которой нужно получить все ранги. Если None, используется группа процессов по умолчанию.

Возвращает:

Список глобальных рангов, упорядоченный по рангам группы.

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

list[int]

torch.distributed.split_group(parent_pg=None, split_ranks=None, timeout=None, pg_options=None, group_desc=None, backend=None) [исходный код]

Создаёт новую группу процессов, разделённую из указанной родительской группы процессов.

Предупреждение:: Это экспериментальный API. Поддерживаются только бэкенды NCCL и пользовательские подключаемые бэкенды. Другие бэкенды вызовут ошибку. Пользователи этого API должны гарантировать, что все ранги родительской группы вызовут этот API и что разделение на подгруппы будет одинаковым для всех рангов родительской группы.

Параметры:
  • parent_pg (ProcessGroup, optional) – Родительская группа процессов. Если None, используется группа процессов по умолчанию. Пользователи должны гарантировать, что родительская группа полностью инициализирована (например, инициализированы коммуникаторы).
  • split_ranks (list[list[int]]) – ранги для разделения — список списков рангов. Пользователи должны убедиться, что список рангов для разделения корректен, то есть одно подмножество (представленное внутренним списком целых чисел) не пересекается с другими подмножествами. Обратите внимание, что ранги в каждом подмножестве — это ранги группы (а не глобальные ранги) в родительской группе процессов. Например, если в родительской группе 4 ранга, split_ranks может иметь значение [[0, 1], [2, 3]]. Также допустимо значение [[0,1]], в этом случае ранги 2 и 3 будут возвращены как участники, не входящие в группу. Порядок рангов внутри каждого подмножества сохраняется: положение в списке определяет ранг в новой группе. Все ранги должны передать одинаковый порядок.
  • timeout (timedelta, optional) – сведения и значение по умолчанию см. в init_process_group.
  • pg_options (ProcessGroupOptions, optional) – Дополнительные параметры, которые нужно передать при создании конкретных групп процессов. Например, можно указать ``is_high_priority_stream``, чтобы группа процессов могла использовать потоки CUDA с высоким приоритетом.
  • group_desc (str, optional) – строка с описанием группы процессов.
  • backend (str or Backend, optional) – выбирает подмножество бэкендов родительской группы процессов для отдельных устройств, которые нужно сохранить в дочерней группе. Строка должна иметь тот же формат, что и init_process_group (например, "nccl" или "cpu:gloo,cuda:nccl"). Каждый тип устройства, указанный в запрошенном бэкенде, должен присутствовать в родительской группе, а имя бэкенда для этого типа устройства должно в точности совпадать с именем в родительской группе. Тип устройства бэкенда дочерней группы по умолчанию должен входить в указанный набор. Если задано None (значение по умолчанию), дочерняя группа наследует всю конфигурацию бэкендов родительской группы.
Возвращает:

ProcessGroup, если текущий ранг входит в одно из подмножеств/подгрупп, заданных split_ranks, или GroupMember.NON_GROUP_MEMBER, если текущий ранг не входит ни в одно подмножество в split_ranks.

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

ProcessGroup | Literal[-100]

torch.distributed.distributed_c10d.GroupName = torch.distributed.distributed_c10d.GroupName

NewType создаёт простые уникальные типы практически без накладных расходов во время выполнения. Статические средства проверки типов считают NewType(name, tp) подтипом tp. Во время выполнения NewType(name, tp) возвращает фиктивный вызываемый объект, который просто возвращает переданный ему аргумент. Пример использования:

UserId = NewType('UserId', int)

def name_by_id(user_id: UserId) -> str:
    ...

UserId('user')          # Fails type check

name_by_id(42)          # Fails type check
name_by_id(UserId(42))  # OK

num = UserId(5) + 1     # type: int

DeviceMesh

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

class torch.distributed.device_mesh.DeviceMesh(device_type, mesh=None, *, mesh_dim_names=None, backend_override=None, _init_backend=True, _rank=None, _layout=None, _rank_map=None, _root_mesh=None) [исходный код]

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

DeviceMesh можно использовать для настройки N-мерных соединений устройств в кластере и управления ProcessGroups для N-мерного параллелизма. Обмен данными может выполняться отдельно по каждому измерению DeviceMesh. DeviceMesh учитывает устройство, уже выбранное пользователем (то есть если пользователь вызвал torch.cuda.set_device до инициализации DeviceMesh), и выберет/установит устройство для текущего процесса, если пользователь не установил его заранее. Обратите внимание, что устройство следует выбирать вручную ДО инициализации DeviceMesh.

DeviceMesh также можно использовать как менеджер контекста при работе с API DTensor.

Примечание

DeviceMesh использует модель программирования SPMD, то есть одна и та же программа PyTorch на Python выполняется на всех процессах/рангах кластера. Поэтому пользователи должны убедиться, что массив mesh (описывающий расположение устройств) одинаков на всех рангах. Несогласованные значения mesh приведут к незаметному зависанию.

Параметры:
  • device_type (str) – Тип устройства сетки. В настоящее время поддерживаются: «cpu», «cuda/cuda-like».
  • mesh (ndarray) – Многомерный массив или целочисленный тензор, описывающий расположение устройств, где идентификаторы являются глобальными идентификаторами группы процессов по умолчанию.
  • _rank (int) – (экспериментальный/внутренний) Глобальный ранг текущего процесса. Если не задан, определяется по группе процессов по умолчанию.
Возвращает:

Объект DeviceMesh, представляющий расположение устройств.

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

DeviceMesh

Следующая программа выполняется на каждом процессе/ранге в режиме SPMD. В этом примере у нас есть 2 хоста, на каждом из которых по 4 GPU. Сокращение по первому измерению сетки будет выполняться по столбцам (0, 4), .. и (3, 7), а сокращение по второму измерению сетки — по строкам (0, 1, 2, 3) и (4, 5, 6, 7).

Пример:

>>> from torch.distributed.device_mesh import DeviceMesh
>>>
>>> # Initialize device mesh as (2, 4) to represent the topology
>>> # of cross-host(dim 0), and within-host (dim 1).
>>> mesh = DeviceMesh(device_type="cuda", mesh=[[0, 1, 2, 3],[4, 5, 6, 7]])
property device_type: str

Возвращает тип устройства сетки.

static from_group(group, device_type, mesh=None, *, mesh_dim_names=None) [исходный код]

Создаёт DeviceMesh с device_type на основе существующей ProcessGroup или списка существующих ProcessGroup.

Созданная сетка устройств имеет столько измерений, сколько передано групп. Например, если передана одна группа процессов, результирующий DeviceMesh будет одномерной сеткой. Если передан список из 2 групп процессов, результирующий DeviceMesh будет двумерной сеткой.

Если передано более одной группы, требуются аргументы mesh и mesh_dim_names. Порядок переданных групп процессов определяет топологию сетки. Например, первая группа процессов будет нулевым измерением DeviceMesh. Переданный тензор mesh должен иметь столько же измерений, сколько передано групп процессов, а порядок измерений в тензоре mesh должен соответствовать порядку переданных групп процессов.

Параметры:
  • group (ProcessGroup or list[ProcessGroup]) – существующая ProcessGroup или список существующих ProcessGroup.
  • device_type (str) – Тип устройства сетки. В настоящее время поддерживаются: «cpu», «cuda/cuda-like». Передавать тип устройства с индексом GPU, например «cuda:0», нельзя.
  • mesh (torch.Tensor or ArrayLike, optional) – Многомерный массив или целочисленный тензор, описывающий расположение устройств, где идентификаторы являются глобальными идентификаторами группы процессов по умолчанию. По умолчанию — None.
  • mesh_dim_names (tuple[str, ...], optional) – Кортеж имён измерений сетки, назначаемых каждому измерению многомерного массива, описывающего расположение устройств. Его длина должна совпадать с длиной mesh_shape. Каждая строка в mesh_dim_names должна быть уникальной. По умолчанию — None.
Возвращает:

Объект DeviceMesh, представляющий расположение устройств.

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

DeviceMesh

get_all_groups() [исходный код]

Возвращает список ProcessGroup для всех измерений сетки.

Возвращает:

Список объектов ProcessGroup.

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

list[ProcessGroup]

get_coordinate() [исходный код]

Возвращает относительные индексы этого ранга по всем измерениям сетки. Если этот ранг не входит в сетку, возвращает None.

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

tuple[int, …] | None

get_group(mesh_dim=None) [исходный код]

Возвращает единственную ProcessGroup, указанную параметром mesh_dim; если mesh_dim не указан, а DeviceMesh одномерный, возвращает единственную ProcessGroup в сетке.

Параметры:
  • mesh_dim (str/python:int, optional) – имя измерения сетки или его индекс
  • None. (измерения сетки. По умолчанию —) –
Возвращает:

Объект ProcessGroup.

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

ProcessGroup

get_local_rank(mesh_dim=None) [исходный код]

Возвращает локальный ранг указанного измерения mesh_dim в DeviceMesh.

Параметры:
  • mesh_dim (str/python:int, optional) – имя измерения сетки или его индекс
  • None. (измерения сетки. По умолчанию —) –
Возвращает:

Целое число, обозначающее локальный ранг.

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

int

Следующая программа выполняется на каждом процессе/ранге в режиме SPMD. В этом примере у нас есть 2 хоста, на каждом из которых по 4 GPU. Вызов mesh_2d.get_local_rank(mesh_dim=0) на рангах 0, 1, 2, 3 вернёт 0. Вызов mesh_2d.get_local_rank(mesh_dim=0) на рангах 4, 5, 6, 7 вернёт 1. Вызов mesh_2d.get_local_rank(mesh_dim=1) на рангах 0, 4 вернёт 0. Вызов mesh_2d.get_local_rank(mesh_dim=1) на рангах 1, 5 вернёт 1. Вызов mesh_2d.get_local_rank(mesh_dim=1) на рангах 2, 6 вернёт 2. Вызов mesh_2d.get_local_rank(mesh_dim=1) на рангах 3, 7 вернёт 3.

Пример:

>>> from torch.distributed.device_mesh import DeviceMesh
>>>
>>> # Initialize device mesh as (2, 4) to represent the topology
>>> # of cross-host(dim 0), and within-host (dim 1).
>>> mesh = DeviceMesh(device_type="cuda", mesh=[[0, 1, 2, 3],[4, 5, 6, 7]])
get_rank() [исходный код]

Возвращает текущий глобальный ранг.

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

int

property mesh: Tensor

Возвращает тензор, представляющий расположение устройств.

property mesh_dim_names: tuple[str, ...] | None

Возвращает имена измерений сетки.

Обмен данными между парами процессов

torch.distributed.send(tensor, dst=None, group=None, tag=0, group_dst=None) [исходный код]

Синхронно отправляет тензор.

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

tag не поддерживается бэкендом NCCL.

Параметры:
  • tensor (Tensor) – Тензор для отправки.
  • dst (int) – Ранг получателя в глобальной группе процессов (независимо от аргумента group). Ранг получателя не должен совпадать с рангом текущего процесса.
  • group (ProcessGroup, optional) – Группа процессов для выполнения операции. Если значение равно None, используется группа процессов по умолчанию.
  • tag (int, optional) – Тег для сопоставления отправки с удалённым вызовом recv
  • group_dst (int, optional) – Ранг получателя в group. Нельзя одновременно указывать dst и group_dst.
torch.distributed.recv(tensor, src=None, group=None, tag=0, group_src=None) [исходный код]

Синхронно получает тензор.

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

tag не поддерживается бэкендом NCCL.

Параметры:
  • tensor (Tensor) – Тензор для заполнения полученными данными.
  • src (int, optional) – Ранг отправителя в глобальной группе процессов (независимо от аргумента group). Если не указан, данные будут получены от любого процесса.
  • group (ProcessGroup, optional) – Группа процессов для выполнения операции. Если значение равно None, используется группа процессов по умолчанию.
  • tag (int, optional) – Тег для сопоставления recv с удалённым вызовом send
  • group_src (int, optional) – Ранг получателя в group. Нельзя одновременно указывать src и group_src.
Возвращает:

Ранг отправителя; -1, если процесс не входит в группу

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

int

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

  • is_completed() — возвращает True, если операция завершена
  • wait() — блокирует процесс до завершения операции. Гарантируется, что is_completed() вернёт True после завершения.
torch.distributed.isend(tensor, dst=None, group=None, tag=0, group_dst=None) [исходный код]

Асинхронно отправляет тензор.

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

Изменение tensor до завершения запроса приводит к неопределённому поведению.

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

tag не поддерживается бэкендом NCCL.

В отличие от блокирующего send, isend допускает совпадение рангов src и dst, то есть отправку самому себе.

Параметры:
  • tensor (Tensor) – Тензор для отправки.
  • dst (int) – Ранг получателя в глобальной группе процессов (независимо от аргумента group)
  • group (ProcessGroup, optional) – Группа процессов для выполнения операции. Если значение равно None, используется группа процессов по умолчанию.
  • tag (int, optional) – Тег для сопоставления отправки с удалённым вызовом recv
  • group_dst (int, optional) – Ранг получателя в group. Нельзя одновременно указывать dst и group_dst
Возвращает:

Объект запроса распределённой операции. None, если процесс не входит в группу

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

Work | None

torch.distributed.irecv(tensor, src=None, group=None, tag=0, group_src=None) [исходный код]

Асинхронно получает тензор.

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

tag не поддерживается бэкендом NCCL.

В отличие от блокирующего recv, irecv допускает совпадение рангов src и dst, то есть получение данных от самого себя.

Параметры:
  • tensor (Tensor) – Тензор для заполнения полученными данными.
  • src (int, optional) – Ранг отправителя в глобальной группе процессов (независимо от аргумента group). Если не указан, данные будут получены от любого процесса.
  • group (ProcessGroup, optional) – Группа процессов для выполнения операции. Если значение равно None, используется группа процессов по умолчанию.
  • tag (int, optional) – Тег для сопоставления recv с удалённым вызовом send
  • group_src (int, optional) – Ранг получателя в group. Нельзя одновременно указывать src и group_src.
Возвращает:

Объект запроса распределённой операции. None, если процесс не входит в группу

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

Work | None

torch.distributed.send_object_list(object_list, dst=None, group=None, device=None, group_dst=None, use_batch=False, weights_only=False) [исходный код]

Синхронно отправляет объекты, поддерживающие сериализацию с помощью pickle, в object_list.

Подобно send(), но позволяет передавать объекты Python. Обратите внимание, что все объекты в object_list должны поддерживать сериализацию с помощью pickle.

Параметры:
  • object_list (List[Any]) – Список входных объектов для отправки. Каждый объект должен поддерживать сериализацию с помощью pickle. Получатель должен предоставить список такого же размера.
  • dst (int) – Ранг получателя, которому отправляются object_list. Ранг получателя определяется относительно глобальной группы процессов (независимо от аргумента group)
  • group (ProcessGroup | None) – (ProcessGroup, optional): Группа процессов для выполнения операции. Если значение равно None, используется группа процессов по умолчанию. Значение по умолчанию — None.
  • device (torch.device, optional) – Если значение не равно None, объекты сериализуются, преобразуются в тензоры и перед отправкой перемещаются на device. Значение по умолчанию — None.
  • group_dst (int, optional) – Ранг получателя в group. Необходимо указать dst или group_dst, но не оба параметра
  • use_batch (bool, optional) – Если значение равно True, вместо обычных операций отправки используются пакетные операции p2p. Это позволяет избежать инициализации коммуникаторов для двух рангов и использовать существующие коммуникаторы для всей группы. Описание использования и предположений см. в batch_isend_irecv. Значение по умолчанию — False.
  • weights_only (bool, optional) – Если значение равно True, объекты сериализуются с помощью torch.save, который получатель может использовать для десериализации с помощью torch.load(weights_only=True). Если значение равно False, используется pickle, небезопасный при работе с недоверенными данными. Значение должно совпадать с переданным получателем в recv_object_list(). Значение по умолчанию — False.
Возвращает:

None.

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

None

Примечание

Для групп процессов на базе NCCL внутренние тензорные представления объектов перед обменом данными необходимо переместить на GPU. В этом случае используемое устройство задаётся параметром torch.cuda.current_device(). Пользователь должен убедиться, что этот параметр настроен так, чтобы каждому рангу соответствовал отдельный GPU, с помощью torch.cuda.set_device().

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

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

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

send_object_list() с weights_only=False неявно использует модуль pickle, который считается небезопасным. Можно создать вредоносные данные pickle, которые при десериализации выполнят произвольный код. Вызывайте эту функцию только с данными, которым доверяете.

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

Вызов send_object_list() с тензорами GPU поддерживается недостаточно хорошо и неэффективен, поскольку при сериализации тензоров с помощью pickle выполняется перенос GPU -> CPU. Вместо этого используйте send().

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> # Assumes backend is not NCCL
>>> device = torch.device("cpu")
>>> if dist.get_rank() == 0:
>>>     # Assumes world_size of 2.
>>>     objects = ["foo", 12, {1: 2}] # any picklable object
>>>     dist.send_object_list(objects, dst=1, device=device)
>>> else:
>>>     objects = [None, None, None]
>>>     dist.recv_object_list(objects, src=0, device=device)
>>> objects
['foo', 12, {1: 2}]
torch.distributed.recv_object_list(object_list, src=None, group=None, device=None, group_src=None, use_batch=False, weights_only=False) [исходный код]

Синхронно получает объекты, поддерживающие сериализацию с помощью pickle, из object_list.

Подобно recv(), но позволяет получать объекты Python.

Параметры:
  • object_list (List[Any]) – Список, в который будут помещены полученные объекты. Размер списка должен совпадать с размером отправляемого списка.
  • src (int, optional) – Ранг отправителя, от которого следует получить object_list. Ранг отправителя определяется относительно глобальной группы процессов (независимо от аргумента group). Если значение равно None, данные будут получены от любого ранга. Значение по умолчанию — None.
  • group (ProcessGroup | None) – (ProcessGroup, optional): Группа процессов для выполнения операции. Если значение равно None, используется группа процессов по умолчанию. Значение по умолчанию — None.
  • device (torch.device, optional) – Если значение не равно None, данные получаются на этом устройстве. Значение по умолчанию — None.
  • group_src (int, optional) – Ранг получателя в group. Нельзя одновременно указывать src и group_src.
  • use_batch (bool, optional) – Если значение равно True, вместо обычных операций отправки используются пакетные операции p2p. Это позволяет избежать инициализации коммуникаторов для двух рангов и использовать существующие коммуникаторы для всей группы. Описание использования и предположений см. в batch_isend_irecv. Значение по умолчанию — False.
  • weights_only (bool, optional) – Если значение равно True, объекты десериализуются с помощью torch.load(weights_only=True), который ограничивает десериализацию безопасными типами. Если значение равно False, используется pickle, небезопасный при работе с недоверенными данными. Значение должно совпадать с переданным отправителем в send_object_list(). Значение по умолчанию — False.
Возвращает:

Ранг отправителя. -1, если ранг не входит в группу. Если ранг входит в группу, object_list будет содержать объекты, отправленные рангом src.

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

int

Примечание

Для групп процессов на базе NCCL внутренние тензорные представления объектов перед обменом данными необходимо переместить на GPU. В этом случае используемое устройство задаётся параметром torch.cuda.current_device(). Пользователь должен убедиться, что этот параметр настроен так, чтобы каждому рангу соответствовал отдельный GPU, с помощью torch.cuda.set_device().

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

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

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

recv_object_list() с weights_only=False неявно использует модуль pickle, который считается небезопасным. Можно создать вредоносные данные pickle, которые при десериализации выполнят произвольный код. Вызывайте эту функцию только с данными, которым доверяете.

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

Вызов recv_object_list() с тензорами GPU поддерживается недостаточно хорошо и неэффективен, поскольку при сериализации тензоров с помощью pickle выполняется перенос GPU -> CPU. Вместо этого используйте recv().

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> # Assumes backend is not NCCL
>>> device = torch.device("cpu")
>>> if dist.get_rank() == 0:
>>>     # Assumes world_size of 2.
>>>     objects = ["foo", 12, {1: 2}] # any picklable object
>>>     dist.send_object_list(objects, dst=1, device=device)
>>> else:
>>>     objects = [None, None, None]
>>>     dist.recv_object_list(objects, src=0, device=device)
>>> objects
['foo', 12, {1: 2}]
torch.distributed.batch_isend_irecv(p2p_op_list) [исходный код]

Асинхронно отправляет или получает пакет тензоров и возвращает список запросов.

Выполняет операции из p2p_op_list и возвращает соответствующие запросы. В настоящее время поддерживаются бэкенды NCCL, Gloo и UCC.

Параметры:

p2p_op_list (list[P2POp]) –

Список операций обмена данными между парами процессов (тип каждого оператора — torch.distributed.P2POp). Все операции в списке рассматриваются как единый пакет: относительный порядок операций отправки и получения не имеет значения и не приводит к взаимным блокировкам. Например, [recv_from_1, send_to_1] эквивалентно [send_to_1, recv_from_1].

Порядок имеет значение для нескольких операций в одном направлении между одной и той же парой рангов: если ранг 0 вызывает [send_to_1(A), send_to_1(B)], ранг 1 должен вызвать [recv_from_0(A), recv_from_0(B)] в том же порядке, чтобы обеспечить правильное сопоставление тензоров.

Возвращает:

Список объектов запросов распределённых операций, возвращённых соответствующими вызовами операций из op_list.

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

list[Work]

Примеры

>>> send_tensor = torch.arange(2, dtype=torch.float32) + 2 * rank
>>> recv_tensor = torch.randn(2, dtype=torch.float32)
>>> send_op = dist.P2POp(dist.isend, send_tensor, (rank + 1) % world_size)
>>> recv_op = dist.P2POp(
...     dist.irecv, recv_tensor, (rank - 1 + world_size) % world_size
... )
>>> reqs = batch_isend_irecv([send_op, recv_op])
>>> for req in reqs:
>>>     req.wait()
>>> recv_tensor
tensor([2, 3])     # Rank 0
tensor([0, 1])     # Rank 1

Примечание

Обратите внимание: при использовании этого API с бэкендом PG NCCL необходимо задать текущий GPU с помощью torch.cuda.set_device, иначе возможны непредвиденные зависания.

Кроме того, если этот вызов API является первым коллективным вызовом в group, переданном в dist.P2POp, в нём должны участвовать все ранги group; в противном случае поведение не определено. Если этот вызов API не является первым коллективным вызовом в group, допускаются пакетные операции P2P с участием только части рангов group.

class torch.distributed.P2POp(op, tensor, peer=None, group=None, tag=0, group_peer=None) [исходный код]

Класс для создания операций обмена данными между парами процессов для batch_isend_irecv.

Этот класс задаёт тип операции P2P, буфер данных, ранг узла-партнёра, группу процессов и тег. Экземпляры этого класса передаются в batch_isend_irecv для обмена данными между парами процессов.

Параметры:
  • op (Callable) – Функция для отправки данных процессу-партнёру или получения данных от него. Тип op — это torch.distributed.isend или torch.distributed.irecv.
  • tensor (Tensor) – Тензор для отправки или получения.
  • peer (int, optional) – Ранг получателя или отправителя.
  • group (ProcessGroup, optional) – Группа процессов для выполнения операции. Если значение равно None, используется группа процессов по умолчанию.
  • tag (int, optional) – Тег для сопоставления отправки с получением.
  • group_peer (int, optional) – Ранг получателя или отправителя.
Тип возвращаемого значения:

P2POp

Синхронные и асинхронные коллективные операции

Каждая функция коллективной операции поддерживает следующие два типа операций, определяемые значением флага async_op, переданного в коллективную операцию:

Синхронная операция — режим по умолчанию, когда async_op установлено в False. К моменту возврата функции гарантируется выполнение коллективной операции. Для операций CUDA завершение самой операции CUDA не гарантируется, поскольку операции CUDA являются асинхронными. В случае коллективных операций на CPU последующие вызовы функций, использующие результат коллективного вызова, будут работать ожидаемым образом. Для коллективных операций CUDA ожидаемым образом будут работать вызовы функций, использующие результат в том же потоке CUDA. При выполнении в разных потоках пользователю необходимо самостоятельно обеспечить синхронизацию. Подробнее о семантике CUDA, например о синхронизации потоков, см. в разделе Семантика CUDA. Примеры ниже показывают различия в семантике операций на CPU и CUDA.

Асинхронная операция — режим, в котором async_op установлено в True. Функция коллективной операции возвращает объект запроса распределённой операции. Как правило, создавать такой объект вручную не нужно; гарантируется поддержка следующих методов:

  • is_completed() — для коллективных операций на CPU возвращает True по завершении. Для операций CUDA возвращает True, если операция успешно помещена в очередь потока CUDA и результат можно использовать в потоке по умолчанию без дополнительной синхронизации.
  • wait() — для коллективных операций на CPU блокирует процесс до завершения операции. Для коллективных операций CUDA блокирует активный поток CUDA до завершения операции (но не блокирует CPU).
  • get_future() — возвращает объект torch._C.Future. Поддерживается для NCCL, а также для большинства операций в GLOO и MPI, кроме операций обмена данными между парами процессов. Примечание: по мере внедрения Futures и объединения API вызов get_future() может стать избыточным.

Пример

Следующий код иллюстрирует семантику операций CUDA при использовании распределённых коллективных операций. Он показывает, что при использовании результатов коллективных операций в разных потоках CUDA необходимо явно выполнять синхронизацию:

# Code runs on each rank.
dist.init_process_group("nccl", rank=rank, world_size=2)
output = torch.tensor([rank]).cuda(rank)
s = torch.cuda.Stream()
handle = dist.all_reduce(output, async_op=True)
# Wait ensures the operation is enqueued, but not necessarily complete.
handle.wait()
# Using result on non-default stream.
with torch.cuda.stream(s):
    s.wait_stream(torch.cuda.default_stream())
    output.add_(100)
if rank == 0:
    # if the explicit call to wait_stream was omitted, the output below will be
    # non-deterministically 1 or 101, depending on whether the allreduce overwrote
    # the value after the add completed.
    print(output)

Коллективные функции

torch.distributed.broadcast(tensor, src=None, group=None, async_op=False, group_src=None) [исходный код]

Рассылает тензор всей группе.

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

Параметры:
  • tensor (Tensor) – Данные для отправки, если src является рангом текущего процесса; в противном случае — тензор для сохранения полученных данных.
  • src (int) – Ранг источника в глобальной группе процессов (независимо от аргумента group).
  • group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Следует ли выполнять эту операцию асинхронно.
  • group_src (int) – Ранг источника в group. Нужно указать либо group_src, либо src, но не оба аргумента.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

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

Work | None

torch.distributed.broadcast_object_list(object_list, src=None, group=None, device=None, group_src=None, weights_only=False) [исходный код]

Рассылает сериализуемые объекты в object_list всей группе.

Подобно broadcast(), но позволяет передавать объекты Python. Обратите внимание: все объекты в object_list должны поддерживать сериализацию для рассылки.

Параметры:
  • object_list (List[Any]) – Список входных объектов для рассылки. Каждый объект должен поддерживать сериализацию. Будут рассылаться только объекты на ранге src, но на каждом ранге необходимо передать списки одинаковой длины.
  • src (int) – Ранг источника, с которого выполняется рассылка object_list. Ранг источника определяется относительно глобальной группы процессов (независимо от аргумента group).
  • group (ProcessGroup | None) – (ProcessGroup, optional): Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию. По умолчанию — None.
  • device (torch.device, optional) – Если значение не None, объекты сериализуются и преобразуются в тензоры, которые перед рассылкой перемещаются на device. По умолчанию — None.
  • group_src (int) – Ранг источника в group. Необходимо указать один из аргументов group_src и src, но не оба.
  • weights_only (bool, optional) – Если True, объекты сериализуются с помощью torch.save и десериализуются с помощью torch.load(weights_only=True), что ограничивает десериализацию безопасными типами. Если False, используется pickle, небезопасный при работе с недоверенными данными. Все ранги должны передавать одинаковое значение. По умолчанию — False.
Возвращает:

None. Если ранг входит в группу, object_list будет содержать объекты, разосланные с ранга src.

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

None

Примечание

Для групп процессов на основе NCCL внутренние тензорные представления объектов перед обменом данными необходимо переместить на GPU. В этом случае используемое устройство задаётся параметром torch.cuda.current_device(), и пользователь должен убедиться, что он настроен так, чтобы каждому рангу соответствовал отдельный GPU, с помощью torch.cuda.set_device().

Примечание

Обратите внимание: этот API немного отличается от коллективной операции broadcast(), поскольку не возвращает дескриптор async_op и поэтому выполняется блокирующим вызовом.

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

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

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

broadcast_object_list() с параметром weights_only=False неявно использует модуль pickle, который считается небезопасным. Можно создать вредоносные данные pickle, которые выполнят произвольный код при десериализации. Вызывайте эту функцию только с данными, которым доверяете.

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

Вызов broadcast_object_list() с тензорами на GPU поддерживается недостаточно хорошо и неэффективен, поскольку при сериализации тензоров pickle данные передаются с GPU на CPU. Вместо этого рекомендуется использовать broadcast().

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> if dist.get_rank() == 0:
>>>     # Assumes world_size of 3.
>>>     objects = ["foo", 12, {1: 2}] # any picklable object
>>> else:
>>>     objects = [None, None, None]
>>> # Assumes backend is not NCCL
>>> device = torch.device("cpu")
>>> dist.broadcast_object_list(objects, src=0, device=device)
>>> objects
['foo', 12, {1: 2}]
torch.distributed.all_reduce(tensor: Tensor, op: _ReduceOp = ReduceOp.SUM, group: ProcessGroup | None = None, *, async_op: Literal[True]) → Work [исходный код]
torch.distributed.all_reduce(tensor:Tensor, op:_ReduceOp=ReduceOp.SUM, group:ProcessGroup|None=None, async_op:bool=False) → Work|None

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

После вызова tensor будет побитово идентичным во всех процессах.

Поддерживаются комплексные тензоры.

Параметры:
  • tensor (Tensor) – Входные и выходные данные коллективной операции. Функция изменяет тензор на месте.
  • op (optional) – Одно из значений перечисления torch.distributed.ReduceOp. Задаёт операцию поэлементной редукции.
  • group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Следует ли выполнять эту операцию асинхронно.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

Примеры

>>> # All tensors below are of torch.int64 type.
>>> # We have 2 process groups, 2 ranks.
>>> device = torch.device(f"cuda:{rank}")
>>> tensor = torch.arange(2, dtype=torch.int64, device=device) + 1 + 2 * rank
>>> tensor
tensor([1, 2], device='cuda:0') # Rank 0
tensor([3, 4], device='cuda:1') # Rank 1
>>> dist.all_reduce(tensor, op=ReduceOp.SUM)
>>> tensor
tensor([4, 6], device='cuda:0') # Rank 0
tensor([4, 6], device='cuda:1') # Rank 1
>>> # All tensors below are of torch.cfloat type.
>>> # We have 2 process groups, 2 ranks.
>>> tensor = torch.tensor(
...     [1 + 1j, 2 + 2j], dtype=torch.cfloat, device=device
... ) + 2 * rank * (1 + 1j)
>>> tensor
tensor([1.+1.j, 2.+2.j], device='cuda:0') # Rank 0
tensor([3.+3.j, 4.+4.j], device='cuda:1') # Rank 1
>>> dist.all_reduce(tensor, op=ReduceOp.SUM)
>>> tensor
tensor([4.+4.j, 6.+6.j], device='cuda:0') # Rank 0
tensor([4.+4.j, 6.+6.j], device='cuda:1') # Rank 1
torch.distributed.all_reduce_coalesced(tensors, op=<RedOpType.SUM: 0>, group=None, async_op=False) [исходный код]

ПРЕДУПРЕЖДЕНИЕ: в настоящее время проверка форм отдельных тензоров между узлами не реализована.

Например, если узел с рангом 0 передаст [torch.rand(4), torch.rand(2)], а узел с рангом 1 — [torch.rand(2), torch.rand(2), torch.rand(2)], операция allreduce выполнится без ошибок и вернёт некорректные результаты. Отсутствие проверки форм значительно повышает производительность, но при использовании этой функции необходимо тщательно следить за тем, чтобы формы передаваемых тензоров совпадали на всех узлах.

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

После вызова каждый тензор в tensors будет побитово идентичным во всех процессах.

Поддерживаются комплексные тензоры.

Параметры:
  • tensors (Union[List[Tensor], Tensor]) – Входные и выходные данные коллективной операции. Функция изменяет данные на месте.
  • op (Optional[ReduceOp]) – Одно из значений перечисления torch.distributed.ReduceOp. Задаёт операцию поэлементной редукции.
  • group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (Optional[bool]) – Следует ли выполнять эту операцию асинхронно.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

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

Future | None

torch.distributed.reduce(tensor, dst=None, op=<RedOpType.SUM: 0>, group=None, async_op=False, group_dst=None) [исходный код]

Выполняет редукцию данных тензора на всех узлах.

Итоговый результат получит только процесс с рангом dst.

Параметры:
  • tensor (Tensor) – Входные и выходные данные коллективной операции. Функция изменяет тензор на месте.
  • dst (int) – Ранг назначения в глобальной группе процессов (независимо от аргумента group).
  • op (optional) – Одно из значений перечисления torch.distributed.ReduceOp. Задаёт операцию поэлементной редукции.
  • group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Следует ли выполнять эту операцию асинхронно.
  • group_dst (int) – Ранг назначения в group. Нужно указать либо group_dst, либо dst, но не оба аргумента.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

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

Work | None

torch.distributed.all_gather(tensor_list: list[Tensor], tensor: Tensor, group: ProcessGroup | C10DBackend | None = None, *, async_op: Literal[True]) → Work [исходный код]
torch.distributed.all_gather(tensor_list:list[Tensor], tensor:Tensor, group:ProcessGroup|C10DBackend|None=None, async_op:bool=False) → Work|None

Собирает тензоры всей группы в список.

Поддерживаются комплексные тензоры и тензоры разного размера.

Параметры:
  • tensor_list (list[Tensor]) – Выходной список. Он должен содержать тензоры подходящего размера для записи результатов коллективной операции. Поддерживаются тензоры разного размера.
  • tensor (Tensor) – Тензор, рассылаемый текущим процессом.
  • group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Следует ли выполнять эту операцию асинхронно.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

Примеры

>>> # All tensors below are of torch.int64 dtype.
>>> # We have 2 process groups, 2 ranks.
>>> device = torch.device(f"cuda:{rank}")
>>> tensor_list = [
...     torch.zeros(2, dtype=torch.int64, device=device) for _ in range(2)
... ]
>>> tensor_list
[tensor([0, 0], device='cuda:0'), tensor([0, 0], device='cuda:0')] # Rank 0
[tensor([0, 0], device='cuda:1'), tensor([0, 0], device='cuda:1')] # Rank 1
>>> tensor = torch.arange(2, dtype=torch.int64, device=device) + 1 + 2 * rank
>>> tensor
tensor([1, 2], device='cuda:0') # Rank 0
tensor([3, 4], device='cuda:1') # Rank 1
>>> dist.all_gather(tensor_list, tensor)
>>> tensor_list
[tensor([1, 2], device='cuda:0'), tensor([3, 4], device='cuda:0')] # Rank 0
[tensor([1, 2], device='cuda:1'), tensor([3, 4], device='cuda:1')] # Rank 1
>>> # All tensors below are of torch.cfloat dtype.
>>> # We have 2 process groups, 2 ranks.
>>> tensor_list = [
...     torch.zeros(2, dtype=torch.cfloat, device=device) for _ in range(2)
... ]
>>> tensor_list
[tensor([0.+0.j, 0.+0.j], device='cuda:0'), tensor([0.+0.j, 0.+0.j], device='cuda:0')] # Rank 0
[tensor([0.+0.j, 0.+0.j], device='cuda:1'), tensor([0.+0.j, 0.+0.j], device='cuda:1')] # Rank 1
>>> tensor = torch.tensor(
...     [1 + 1j, 2 + 2j], dtype=torch.cfloat, device=device
... ) + 2 * rank * (1 + 1j)
>>> tensor
tensor([1.+1.j, 2.+2.j], device='cuda:0') # Rank 0
tensor([3.+3.j, 4.+4.j], device='cuda:1') # Rank 1
>>> dist.all_gather(tensor_list, tensor)
>>> tensor_list
[tensor([1.+1.j, 2.+2.j], device='cuda:0'), tensor([3.+3.j, 4.+4.j], device='cuda:0')] # Rank 0
[tensor([1.+1.j, 2.+2.j], device='cuda:1'), tensor([3.+3.j, 4.+4.j], device='cuda:1')] # Rank 1
torch.distributed.all_gather_single(output_tensor, input_tensor, group=None, async_op=False) [исходный код]

Собирает тензоры всех рангов и помещает их в один выходной тензор.

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

Параметры:
  • output_tensor (Tensor) – Выходной тензор для размещения элементов тензоров со всех рангов. Он должен иметь подходящий размер и одну из следующих форм: (i) конкатенация всех входных тензоров вдоль основного измерения; определение «конкатенации» см. в torch.cat(); (ii) стек всех входных тензоров вдоль основного измерения; определение «стека» см. в torch.stack(). Примеры ниже поясняют поддерживаемые формы выходного тензора.
  • input_tensor (Tensor) – Тензор, собираемый с текущего ранга. В отличие от API all_gather, входные тензоры в этом API должны иметь одинаковый размер на всех рангах.
  • group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Следует ли выполнять эту операцию асинхронно.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

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

Work | None

Примеры

>>> # All tensors below are of torch.int64 dtype and on CUDA devices.
>>> # We have two ranks.
>>> device = torch.device(f"cuda:{rank}")
>>> tensor_in = torch.arange(2, dtype=torch.int64, device=device) + 1 + 2 * rank
>>> tensor_in
tensor([1, 2], device='cuda:0') # Rank 0
tensor([3, 4], device='cuda:1') # Rank 1
>>> # Output in concatenation form
>>> tensor_out = torch.zeros(world_size * 2, dtype=torch.int64, device=device)
>>> dist.all_gather_single(tensor_out, tensor_in)
>>> tensor_out
tensor([1, 2, 3, 4], device='cuda:0') # Rank 0
tensor([1, 2, 3, 4], device='cuda:1') # Rank 1
>>> # Output in stack form
>>> tensor_out2 = torch.zeros(world_size, 2, dtype=torch.int64, device=device)
>>> dist.all_gather_single(tensor_out2, tensor_in)
>>> tensor_out2
tensor([[1, 2],
        [3, 4]], device='cuda:0') # Rank 0
tensor([[1, 2],
        [3, 4]], device='cuda:1') # Rank 1
torch.distributed.all_gather_object(object_list, obj, group=None, weights_only=False) [исходный код]

Собирает сериализуемые объекты всей группы в список.

Подобно all_gather(), но позволяет передавать объекты Python. Обратите внимание: объект должен поддерживать сериализацию для сбора.

Параметры:
  • object_list (list[Any]) – Выходной список. Его размер должен соответствовать размеру группы для этой коллективной операции; в нём будут размещены результаты.
  • obj (Any) – Объект Python, поддерживающий сериализацию, который рассылается текущим процессом.
  • group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию. По умолчанию — None.
  • weights_only (bool, optional) – Если True, объекты сериализуются с помощью torch.save и десериализуются с помощью torch.load(weights_only=True), что ограничивает десериализацию безопасными типами. Если False, используется pickle, небезопасный при работе с недоверенными данными. Все ранги должны передавать одинаковое значение. По умолчанию — False.
Возвращает:

None. Если вызывающий ранг входит в эту группу, результат коллективной операции будет записан во входной object_list. Если вызывающий ранг не входит в группу, переданный object_list останется без изменений.

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

None

Примечание

Обратите внимание: этот API немного отличается от коллективной операции all_gather(), поскольку не возвращает дескриптор async_op и поэтому выполняется блокирующим вызовом.

Примечание

Для групп процессов на основе NCCL внутренние тензорные представления объектов перед обменом данными необходимо переместить на GPU. В этом случае используемое устройство задаётся параметром torch.cuda.current_device(), и пользователь должен убедиться, что он настроен так, чтобы каждому рангу соответствовал отдельный GPU, с помощью torch.cuda.set_device().

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

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

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

all_gather_object() с параметром weights_only=False неявно использует модуль pickle, который считается небезопасным. Можно создать вредоносные данные pickle, которые выполнят произвольный код при десериализации. Вызывайте эту функцию только с данными, которым доверяете.

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

Вызов all_gather_object() с тензорами на GPU поддерживается недостаточно хорошо и неэффективен, поскольку при сериализации тензоров pickle данные передаются с GPU на CPU. Вместо этого рекомендуется использовать all_gather().

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> # Assumes world_size of 3.
>>> gather_objects = ["foo", 12, {1: 2}] # any picklable object
>>> output = [None for _ in gather_objects]
>>> dist.all_gather_object(output, gather_objects[dist.get_rank()])
>>> output
['foo', 12, {1: 2}]
torch.distributed.all_gather_coalesced(output_tensor_lists, input_tensor_list, group=None, async_op=False) [исходный код]

Собирает входные тензоры всей группы в список пакетным способом.

Поддерживаются комплексные тензоры.

Параметры:
  • output_tensor_lists (list[list[Tensor]]) – Выходной список. Он должен содержать тензоры подходящего размера для записи результатов коллективной операции.
  • input_tensor_list (list[Tensor]) – Тензоры, рассылаемые текущим процессом. Как минимум один тензор должен быть непустым.
  • group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Следует ли выполнять эту операцию асинхронно.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

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

Future | None

Пример:

# we have 2 process groups, 2 ranks.
# rank 0 passes:
input_tensor_list = [[[1, 1], [1, 1]], [2], [3, 3]]
output_tensor_lists = [
    [[[-1, -1], [-1, -1]], [-1], [-1, -1]],
    [[[-1, -1], [-1, -1]], [-1], [-1, -1]],
]
# rank 1 passes:
input_tensor_list = [[[3, 3], [3, 3]], [5], [1, 1]]
output_tensor_lists = [
    [[[-1, -1], [-1, -1]], [-1], [-1, -1]],
    [[[-1, -1], [-1, -1]], [-1], [-1, -1]],
]
# both rank 0 and 1 get:
output_tensor_lists = [
    [[[1, 1], [1, 1]], [2], [3, 3]],
    [[[3, 3], [3, 3]], [5], [1, 1]],
]

ПРЕДУПРЕЖДЕНИЕ: в настоящее время проверка форм отдельных тензоров между узлами не реализована. Например, если узел с рангом 0 передаст [torch.rand(4), torch.rand(2)], а узел с рангом 1 — [torch.rand(2), torch.rand(2), torch.rand(2)], операция all_gather_coalesced выполнится без ошибок и вернёт некорректные результаты. Отсутствие проверки форм значительно повышает производительность, но при использовании этой функции необходимо тщательно следить за тем, чтобы формы передаваемых тензоров совпадали на всех узлах.

torch.distributed.gather(tensor, gather_list=None, dst=None, group=None, async_op=False, group_dst=None) [исходный код]

Собирает список тензоров в одном процессе.

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

Параметры:
  • tensor (Tensor) – Входной тензор.
  • gather_list (list[Tensor], optional) – Список тензоров подходящего одинакового размера для сохранения собранных данных (по умолчанию None; должен быть указан на целевом ранге)
  • dst (int, optional) – Целевой ранг в глобальной группе процессов (независимо от аргумента group). (Если dst и group_dst равны None, по умолчанию используется глобальный ранг 0)
  • group (ProcessGroup, optional) – Группа процессов, в которой выполняется операция. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Должна ли эта операция выполняться асинхронно
  • group_dst (int, optional) – Целевой ранг в group. Нельзя одновременно указывать dst и group_dst
Возвращает:

Дескриптор асинхронной работы, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу

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

Work | None

Примечание

Обратите внимание, что все тензоры в gather_list должны иметь одинаковый размер.

Пример::
>>> # We have 2 process groups, 2 ranks.
>>> tensor_size = 2
>>> device = torch.device(f'cuda:{rank}')
>>> tensor = torch.ones(tensor_size, device=device) + rank
>>> if dist.get_rank() == 0:
>>>     gather_list = [torch.zeros_like(tensor, device=device) for i in range(2)]
>>> else:
>>>     gather_list = None
>>> dist.gather(tensor, gather_list, dst=0)
>>> # Rank 0 gets gathered data.
>>> gather_list
[tensor([1., 1.], device='cuda:0'), tensor([2., 2.], device='cuda:0')] # Rank 0
None                                                                   # Rank 1
torch.distributed.gather_single(tensor, gather_tensor=None, dst=None, group=None, async_op=False, group_dst=None) [исходный код]

Собирает входной тензор со всех рангов в один выходной тензор на dst.

Это аналог gather() с одним выходным тензором: вместо заполнения списка Python с тензорами для каждого ранга на целевом ранге вклад каждого ранга записывается непосредственно в один gather_tensor правильного размера. Бэкенды, поддерживающие собственный путь сбора непосредственно в тензор (например, NCCL >= 2.28.3 через ncclGather), позволяют избежать дополнительного копирования для каждого ранга, которое выполняет gather().

Эта функция требует, чтобы размер tensor был одинаковым в каждом процессе.

Параметры:
  • tensor (Tensor) – Входной тензор для сбора с текущего ранга.
  • gather_tensor (Tensor, optional) – Выходной тензор для размещения вкладов со всех рангов. Он должен быть непрерывным и иметь размер, достаточный для хранения world_size копий tensor; проверяется только общее число его элементов, поэтому подойдут как плоская конкатенация (world_size * tensor.numel()), так и стек ([world_size, *tensor.shape]), поскольку они имеют одинаковое непрерывное размещение. Для непрерывного выходного тензора (например, представления стека с шагом) возникает ошибка непрерывности, а не ошибка формы. Требуется (и используется) только на целевом ранге.
  • dst (int, optional) – Целевой ранг в глобальной группе процессов (независимо от аргумента group). (Если dst и group_dst равны None, по умолчанию используется глобальный ранг 0)
  • group (ProcessGroup, optional) – Группа процессов, в которой выполняется операция. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Должна ли эта операция выполняться асинхронно
  • group_dst (int, optional) – Целевой ранг в group. Нельзя одновременно указывать dst и group_dst
Возвращает:

Дескриптор асинхронной работы, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу

Пример::
>>> # We have 2 process groups, 2 ranks.
>>> device = torch.device(f"cuda:{rank}")
>>> tensor = torch.arange(2, dtype=torch.int64, device=device) + 1 + 2 * rank
>>> if dist.get_rank() == 0:
>>>     gather_tensor = torch.zeros(2 * 2, dtype=torch.int64, device=device)
>>> else:
>>>     gather_tensor = None
>>> dist.gather_single(tensor, gather_tensor, dst=0)
>>> gather_tensor
tensor([1, 2, 3, 4], device='cuda:0')  # Rank 0
None                                    # Rank 1
torch.distributed.gather_object(obj, object_gather_list=None, dst=None, group=None, group_dst=None, weights_only=False) [исходный код]

Собирает сериализуемые объекты из всей группы в одном процессе.

Аналогично gather(), но можно передавать объекты Python. Обратите внимание, что объект должен поддерживать сериализацию с помощью pickle, чтобы его можно было собрать.

Параметры:
  • obj (Any) – Входной объект. Должен поддерживать сериализацию с помощью pickle.
  • object_gather_list (list[Any]) – Выходной список. На ранге dst его размер должен соответствовать размеру группы для этой коллективной операции; в него будет записан результат. На рангах, отличных от dst, значение должно быть None. (По умолчанию — None)
  • dst (int, optional) – Целевой ранг в глобальной группе процессов (независимо от аргумента group). (Если dst и group_dst равны None, по умолчанию используется глобальный ранг 0)
  • group (ProcessGroup | None) – (ProcessGroup, необязательный): группа процессов, в которой выполняется операция. Если None, будет использоваться группа процессов по умолчанию. По умолчанию — None.
  • group_dst (int, optional) – Целевой ранг в group. Нельзя одновременно указывать dst и group_dst
  • weights_only (bool, optional) – Если задано True, объекты сериализуются с помощью torch.save, а десериализуются с помощью torch.load(weights_only=True), что ограничивает десериализацию безопасными типами. Если задано False, используется pickle, что небезопасно при работе с недоверенными данными. Все ранги должны передавать одинаковое значение. По умолчанию — False.
Возвращает:

None. На ранге dst в object_gather_list будет записан результат коллективной операции.

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

None

Примечание

Обратите внимание, что этот API немного отличается от коллективной операции gather: он не предоставляет дескриптор async_op и поэтому выполняется синхронно.

Примечание

Для групп процессов на основе NCCL внутренние тензорные представления объектов необходимо переместить на устройство GPU до начала обмена данными. В этом случае используемое устройство задаётся параметром torch.cuda.current_device(); пользователь должен обеспечить его настройку таким образом, чтобы у каждого ранга был отдельный GPU, с помощью torch.cuda.set_device().

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

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

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

gather_object() с weights_only=False неявно использует модуль pickle, который считается небезопасным. Можно создать вредоносные данные pickle, которые выполнят произвольный код при десериализации. Вызывайте эту функцию только с данными, которым доверяете.

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

Вызов gather_object() с тензорами GPU поддерживается недостаточно хорошо и неэффективен, поскольку требует передачи данных с GPU на CPU: тензоры сериализуются с помощью pickle. Рассмотрите возможность использования gather().

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> # Assumes world_size of 3.
>>> gather_objects = ["foo", 12, {1: 2}] # any picklable object
>>> output = [None for _ in gather_objects]
>>> dist.gather_object(
...     gather_objects[dist.get_rank()],
...     output if dist.get_rank() == 0 else None,
...     dst=0
... )
>>> # On rank 0
>>> output
['foo', 12, {1: 2}]
torch.distributed.scatter(tensor, scatter_list=None, src=None, group=None, async_op=False, group_src=None) [исходный код]

Распределяет список тензоров между всеми процессами группы.

Каждый процесс получит ровно один тензор и сохранит его данные в аргументе tensor.

Поддерживаются комплексные тензоры.

Параметры:
  • tensor (Tensor) – Выходной тензор.
  • scatter_list (list[Tensor]) – Список тензоров для распределения (по умолчанию None; должен быть указан на исходном ранге)
  • src (int) – Исходный ранг в глобальной группе процессов (независимо от аргумента group). (Если src и group_src равны None, по умолчанию используется глобальный ранг 0)
  • group (ProcessGroup, optional) – Группа процессов, в которой выполняется операция. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Должна ли эта операция выполняться асинхронно
  • group_src (int, optional) – Исходный ранг в group. Нельзя одновременно указывать src и group_src
Возвращает:

Дескриптор асинхронной работы, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу

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

Work | None

Примечание

Обратите внимание, что все тензоры в scatter_list должны иметь одинаковый размер.

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> tensor_size = 2
>>> device = torch.device(f'cuda:{rank}')
>>> output_tensor = torch.zeros(tensor_size, device=device)
>>> if dist.get_rank() == 0:
>>>     # Assumes world_size of 2.
>>>     # Only tensors, all of which must be the same size.
>>>     t_ones = torch.ones(tensor_size, device=device)
>>>     t_fives = torch.ones(tensor_size, device=device) * 5
>>>     scatter_list = [t_ones, t_fives]
>>> else:
>>>     scatter_list = None
>>> dist.scatter(output_tensor, scatter_list, src=0)
>>> # Rank i gets scatter_list[i].
>>> output_tensor
tensor([1., 1.], device='cuda:0') # Rank 0
tensor([5., 5.], device='cuda:1') # Rank 1
torch.distributed.scatter_object_list(scatter_object_output_list, scatter_object_input_list=None, src=None, group=None, group_src=None, weights_only=False) [исходный код]

Распределяет сериализуемые объекты из scatter_object_input_list по всей группе.

Аналогично scatter(), но можно передавать объекты Python. На каждом ранге распределённый объект будет сохранён как первый элемент scatter_object_output_list. Обратите внимание, что все объекты в scatter_object_input_list должны поддерживать сериализацию с помощью pickle, чтобы их можно было распределить.

Параметры:
  • scatter_object_output_list (List[Any]) – Непустой список, в первый элемент которого будет записан объект, распределённый на этот ранг.
  • scatter_object_input_list (List[Any], optional) – Список входных объектов для распределения. Каждый объект должен поддерживать сериализацию с помощью pickle. Будут распределены только объекты на ранге src; для рангов, отличных от src, аргумент может иметь значение None.
  • src (int) – Исходный ранг, с которого распределяется scatter_object_input_list. Исходный ранг определяется относительно глобальной группы процессов (независимо от аргумента group). (Если src и group_src равны None, по умолчанию используется глобальный ранг 0)
  • group (ProcessGroup | None) – (ProcessGroup, необязательный): группа процессов, в которой выполняется операция. Если None, будет использоваться группа процессов по умолчанию. По умолчанию — None.
  • group_src (int, optional) – Исходный ранг в group. Нельзя одновременно указывать src и group_src
  • weights_only (bool, optional) – Если задано True, объекты сериализуются с помощью torch.save, а десериализуются с помощью torch.load(weights_only=True), что ограничивает десериализацию безопасными типами. Если задано False, используется pickle, что небезопасно при работе с недоверенными данными. Все ранги должны передавать одинаковое значение. По умолчанию — False.
Возвращает:

None. Если ранг входит в группу, первый элемент scatter_object_output_list будет содержать объект, распределённый на этот ранг.

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

None

Примечание

Обратите внимание, что этот API немного отличается от коллективной операции scatter: он не предоставляет дескриптор async_op и поэтому выполняется синхронно.

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

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

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

scatter_object_list() с weights_only=False неявно использует модуль pickle, который считается небезопасным. Можно создать вредоносные данные pickle, которые выполнят произвольный код при десериализации. Вызывайте эту функцию только с данными, которым доверяете.

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

Вызов scatter_object_list() с тензорами GPU поддерживается недостаточно хорошо и неэффективен, поскольку требует передачи данных с GPU на CPU: тензоры сериализуются с помощью pickle. Рассмотрите возможность использования scatter().

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> if dist.get_rank() == 0:
>>>     # Assumes world_size of 3.
>>>     objects = ["foo", 12, {1: 2}] # any picklable object
>>> else:
>>>     # Can be any list on non-src ranks, elements are not used.
>>>     objects = [None, None, None]
>>> output_list = [None]
>>> dist.scatter_object_list(output_list, objects, src=0)
>>> # Rank i gets objects[i]. For example, on rank 2:
>>> output_list
[{1: 2}]
torch.distributed.reduce_scatter(output: Tensor, input_list: list[Tensor], op: _ReduceOp = ReduceOp.SUM, group: ProcessGroup | None = None, *, async_op: Literal[True]) → Work [исходный код]
torch.distributed.reduce_scatter(output:Tensor, input_list:list[Tensor], op:_ReduceOp=ReduceOp.SUM, group:ProcessGroup|None=None, async_op:bool=False) → Work|None

Выполняет редукцию, а затем распределяет список тензоров между всеми процессами группы.

Параметры:
  • output (Tensor) – Выходной тензор.
  • input_list (list[Tensor]) – Список тензоров для редукции и распределения.
  • op (optional) – Одно из значений перечисления torch.distributed.ReduceOp. Задаёт операцию для поэлементной редукции.
  • group (ProcessGroup, optional) – Группа процессов, в которой выполняется операция. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Должна ли эта операция выполняться асинхронно.
Возвращает:

Дескриптор асинхронной работы, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

torch.distributed.reduce_scatter_single(output: Tensor, input: Tensor, op: _ReduceOp = ReduceOp.SUM, group: ProcessGroup | None = None, *, async_op: Literal[True]) → Work [исходный код]
torch.distributed.reduce_scatter_single(output:Tensor, input:Tensor, op:_ReduceOp=ReduceOp.SUM, group:ProcessGroup|None=None, async_op:bool=False) → Work|None

Выполняет редукцию, а затем распределяет тензор между всеми рангами группы.

Параметры:
  • output (Tensor) – Выходной тензор. Он должен иметь одинаковый размер на всех рангах.
  • input (Tensor) – Входной тензор для редукции и распределения. Его размер должен равняться размеру выходного тензора, умноженному на размер мира. Входной тензор может иметь одну из следующих форм: (i) конкатенация выходных тензоров вдоль основного измерения или (ii) стек выходных тензоров вдоль основного измерения. Определение «конкатенации» см. в torch.cat(). Определение «стека» см. в torch.stack().
  • group (ProcessGroup, optional) – Группа процессов, в которой выполняется операция. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Должна ли эта операция выполняться асинхронно.
Возвращает:

Дескриптор асинхронной работы, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

Примеры

>>> # All tensors below are of torch.int64 dtype and on CUDA devices.
>>> # We have two ranks.
>>> device = torch.device(f"cuda:{rank}")
>>> tensor_out = torch.zeros(2, dtype=torch.int64, device=device)
>>> # Input in concatenation form
>>> tensor_in = torch.arange(world_size * 2, dtype=torch.int64, device=device)
>>> tensor_in
tensor([0, 1, 2, 3], device='cuda:0') # Rank 0
tensor([0, 1, 2, 3], device='cuda:1') # Rank 1
>>> dist.reduce_scatter_single(tensor_out, tensor_in)
>>> tensor_out
tensor([0, 2], device='cuda:0') # Rank 0
tensor([4, 6], device='cuda:1') # Rank 1
>>> # Input in stack form
>>> tensor_in = torch.reshape(tensor_in, (world_size, 2))
>>> tensor_in
tensor([[0, 1],
        [2, 3]], device='cuda:0') # Rank 0
tensor([[0, 1],
        [2, 3]], device='cuda:1') # Rank 1
>>> dist.reduce_scatter_single(tensor_out, tensor_in)
>>> tensor_out
tensor([0, 2], device='cuda:0') # Rank 0
tensor([4, 6], device='cuda:1') # Rank 1
torch.distributed.all_to_all_single(output, input, output_split_sizes=None, input_split_sizes=None, group=None, async_op=False) [исходный код]

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

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

Поддерживаются комплексные тензоры.

Параметры:
  • output (Tensor) – Собранный выходной тензор-конкатенация.
  • input (Tensor) – Входной тензор для распределения.
  • output_split_sizes (list[int] | None) – (list[Int], необязательный): размеры частей выходного тензора для измерения 0; если указано None или пустой список, размер измерения 0 тензора output должен без остатка делиться на world_size.
  • input_split_sizes (list[int] | None) – (list[Int], необязательный): размеры частей входного тензора для измерения 0; если указано None или пустой список, размер измерения 0 тензора input должен без остатка делиться на world_size.
  • group (ProcessGroup, optional) – Группа процессов, в которой выполняется операция. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Должна ли эта операция выполняться асинхронно.
Возвращает:

Дескриптор асинхронной работы, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

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

Work | None

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

all_to_all_single является экспериментальной функцией и может измениться.

Примеры

>>> input = torch.arange(4) + rank * 4
>>> input
tensor([0, 1, 2, 3])     # Rank 0
tensor([4, 5, 6, 7])     # Rank 1
tensor([8, 9, 10, 11])   # Rank 2
tensor([12, 13, 14, 15]) # Rank 3
>>> output = torch.empty([4], dtype=torch.int64)
>>> dist.all_to_all_single(output, input)
>>> output
tensor([0, 4, 8, 12])    # Rank 0
tensor([1, 5, 9, 13])    # Rank 1
tensor([2, 6, 10, 14])   # Rank 2
tensor([3, 7, 11, 15])   # Rank 3
>>> # Essentially, it is similar to following operation:
>>> scatter_list = list(input.chunk(world_size))
>>> gather_list = list(output.chunk(world_size))
>>> for i in range(world_size):
>>>     dist.scatter(gather_list[i], scatter_list if i == rank else [], src = i)
>>> # Another example with uneven split
>>> input
tensor([0, 1, 2, 3, 4, 5])                                       # Rank 0
tensor([10, 11, 12, 13, 14, 15, 16, 17, 18])                     # Rank 1
tensor([20, 21, 22, 23, 24])                                     # Rank 2
tensor([30, 31, 32, 33, 34, 35, 36])                             # Rank 3
>>> input_splits
[2, 2, 1, 1]                                                     # Rank 0
[3, 2, 2, 2]                                                     # Rank 1
[2, 1, 1, 1]                                                     # Rank 2
[2, 2, 2, 1]                                                     # Rank 3
>>> output_splits
[2, 3, 2, 2]                                                     # Rank 0
[2, 2, 1, 2]                                                     # Rank 1
[1, 2, 1, 2]                                                     # Rank 2
[1, 2, 1, 1]                                                     # Rank 3
>>> output = ...
>>> dist.all_to_all_single(output, input, output_splits, input_splits)
>>> output
tensor([ 0,  1, 10, 11, 12, 20, 21, 30, 31])                     # Rank 0
tensor([ 2,  3, 13, 14, 22, 32, 33])                             # Rank 1
tensor([ 4, 15, 16, 23, 34, 35])                                 # Rank 2
tensor([ 5, 17, 18, 24, 36])                                     # Rank 3
>>> # Another example with tensors of torch.cfloat type.
>>> input = torch.tensor(
...     [1 + 1j, 2 + 2j, 3 + 3j, 4 + 4j], dtype=torch.cfloat
... ) + 4 * rank * (1 + 1j)
>>> input
tensor([1+1j, 2+2j, 3+3j, 4+4j])                                # Rank 0
tensor([5+5j, 6+6j, 7+7j, 8+8j])                                # Rank 1
tensor([9+9j, 10+10j, 11+11j, 12+12j])                          # Rank 2
tensor([13+13j, 14+14j, 15+15j, 16+16j])                        # Rank 3
>>> output = torch.empty([4], dtype=torch.int64)
>>> dist.all_to_all_single(output, input)
>>> output
tensor([1+1j, 5+5j, 9+9j, 13+13j])                              # Rank 0
tensor([2+2j, 6+6j, 10+10j, 14+14j])                            # Rank 1
tensor([3+3j, 7+7j, 11+11j, 15+15j])                            # Rank 2
tensor([4+4j, 8+8j, 12+12j, 16+16j])                            # Rank 3
torch.distributed.all_to_all(output_tensor_list, input_tensor_list, group=None, async_op=False) [исходный код]

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

Поддерживаются комплексные тензоры.

Параметры:
  • output_tensor_list (list[Tensor]) – Список тензоров для сбора — по одному на каждый ранг.
  • input_tensor_list (list[Tensor]) – Список тензоров для распределения — по одному на каждый ранг.
  • group (ProcessGroup, optional) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Должна ли эта операция выполняться асинхронно.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

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

Work | None

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

all_to_all является экспериментальной и может измениться.

Примеры

>>> input = torch.arange(4) + rank * 4
>>> input = list(input.chunk(4))
>>> input
[tensor([0]), tensor([1]), tensor([2]), tensor([3])]     # Rank 0
[tensor([4]), tensor([5]), tensor([6]), tensor([7])]     # Rank 1
[tensor([8]), tensor([9]), tensor([10]), tensor([11])]   # Rank 2
[tensor([12]), tensor([13]), tensor([14]), tensor([15])] # Rank 3
>>> output = list(torch.empty([4], dtype=torch.int64).chunk(4))
>>> dist.all_to_all(output, input)
>>> output
[tensor([0]), tensor([4]), tensor([8]), tensor([12])]    # Rank 0
[tensor([1]), tensor([5]), tensor([9]), tensor([13])]    # Rank 1
[tensor([2]), tensor([6]), tensor([10]), tensor([14])]   # Rank 2
[tensor([3]), tensor([7]), tensor([11]), tensor([15])]   # Rank 3
>>> # Essentially, it is similar to following operation:
>>> scatter_list = input
>>> gather_list = output
>>> for i in range(world_size):
>>>     dist.scatter(gather_list[i], scatter_list if i == rank else [], src=i)
>>> input
tensor([0, 1, 2, 3, 4, 5])                                       # Rank 0
tensor([10, 11, 12, 13, 14, 15, 16, 17, 18])                     # Rank 1
tensor([20, 21, 22, 23, 24])                                     # Rank 2
tensor([30, 31, 32, 33, 34, 35, 36])                             # Rank 3
>>> input_splits
[2, 2, 1, 1]                                                     # Rank 0
[3, 2, 2, 2]                                                     # Rank 1
[2, 1, 1, 1]                                                     # Rank 2
[2, 2, 2, 1]                                                     # Rank 3
>>> output_splits
[2, 3, 2, 2]                                                     # Rank 0
[2, 2, 1, 2]                                                     # Rank 1
[1, 2, 1, 2]                                                     # Rank 2
[1, 2, 1, 1]                                                     # Rank 3
>>> input = list(input.split(input_splits))
>>> input
[tensor([0, 1]), tensor([2, 3]), tensor([4]), tensor([5])]                   # Rank 0
[tensor([10, 11, 12]), tensor([13, 14]), tensor([15, 16]), tensor([17, 18])] # Rank 1
[tensor([20, 21]), tensor([22]), tensor([23]), tensor([24])]                 # Rank 2
[tensor([30, 31]), tensor([32, 33]), tensor([34, 35]), tensor([36])]         # Rank 3
>>> output = ...
>>> dist.all_to_all(output, input)
>>> output
[tensor([0, 1]), tensor([10, 11, 12]), tensor([20, 21]), tensor([30, 31])]   # Rank 0
[tensor([2, 3]), tensor([13, 14]), tensor([22]), tensor([32, 33])]           # Rank 1
[tensor([4]), tensor([15, 16]), tensor([23]), tensor([34, 35])]              # Rank 2
[tensor([5]), tensor([17, 18]), tensor([24]), tensor([36])]                  # Rank 3
>>> # Another example with tensors of torch.cfloat type.
>>> input = torch.tensor(
...     [1 + 1j, 2 + 2j, 3 + 3j, 4 + 4j], dtype=torch.cfloat
... ) + 4 * rank * (1 + 1j)
>>> input = list(input.chunk(4))
>>> input
[tensor([1+1j]), tensor([2+2j]), tensor([3+3j]), tensor([4+4j])]            # Rank 0
[tensor([5+5j]), tensor([6+6j]), tensor([7+7j]), tensor([8+8j])]            # Rank 1
[tensor([9+9j]), tensor([10+10j]), tensor([11+11j]), tensor([12+12j])]      # Rank 2
[tensor([13+13j]), tensor([14+14j]), tensor([15+15j]), tensor([16+16j])]    # Rank 3
>>> output = list(torch.empty([4], dtype=torch.int64).chunk(4))
>>> dist.all_to_all(output, input)
>>> output
[tensor([1+1j]), tensor([5+5j]), tensor([9+9j]), tensor([13+13j])]          # Rank 0
[tensor([2+2j]), tensor([6+6j]), tensor([10+10j]), tensor([14+14j])]        # Rank 1
[tensor([3+3j]), tensor([7+7j]), tensor([11+11j]), tensor([15+15j])]        # Rank 2
[tensor([4+4j]), tensor([8+8j]), tensor([12+12j]), tensor([16+16j])]        # Rank 3
torch.distributed.barrier(group: ProcessGroup | None = GroupMember.WORLD, *, async_op: Literal[True], device_ids: list[int] | None = None, timeout: timedelta | None = None) → Work [исходный код]
torch.distributed.barrier(group:ProcessGroup|None=GroupMember.WORLD, async_op:bool=False, device_ids:list[int]|None=None, timeout:timedelta|None=None) → Work|None

Синхронизирует все процессы.

Эта коллективная операция блокирует процессы, пока вся группа не вызовет эту функцию, если async_op имеет значение False, либо пока для дескриптора асинхронной операции не будет вызван wait().

Параметры:
  • group (ProcessGroup, optional) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, optional) – Должна ли эта операция выполняться асинхронно.
  • device_ids ([int], optional) – Список идентификаторов устройств/GPU. Ожидается только один идентификатор.
  • timeout (datetime.timedelta, optional) – Тайм-аут для барьера. Если None, будет использоваться тайм-аут группы процессов по умолчанию.
Возвращает:

Дескриптор асинхронной операции, если для async_op задано значение True. None, если async_op не задан или процесс не входит в группу.

Примечание

ProcessGroupNCCL теперь блокирует поток CPU до завершения коллективной операции барьера.

Примечание

ProcessGroupNCCL реализует барьер как all_reduce для тензора из одного элемента. Для выделения памяти под этот тензор необходимо выбрать устройство. Устройство выбирается в следующем порядке: (1) первое устройство, переданное в аргумент device_ids функции barrier, если он не равен None; (2) устройство, переданное в init_process_group, если оно не равно None; (3) устройство, которое первым использовалось с этой группой процессов, если выполнялась другая коллективная операция с тензорными входными данными; (4) индекс устройства, соответствующий остатку от деления глобального ранга на количество локальных устройств.

torch.distributed.monitored_barrier(group=None, timeout=None, wait_all_ranks=False) [исходный код]

Синхронизирует процессы аналогично torch.distributed.barrier, но с настраиваемым тайм-аутом.

Функция может сообщить о рангах, которые не достигли этого барьера в течение заданного времени ожидания. В частности, ранги, отличные от нуля, будут блокироваться до обработки отправки/получения от ранга 0. Ранг 0 будет блокироваться до обработки всех операций отправки/получения от других рангов и сообщит об ошибках для рангов, не ответивших вовремя. Обратите внимание: если один ранг не достигнет monitored_barrier (например, из-за зависания), на всех остальных рангах monitored_barrier завершится с ошибкой.

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

Примечание

Обратите внимание, что эта коллективная операция поддерживается только бэкендом GLOO.

Параметры:
  • group (ProcessGroup, optional) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию.
  • timeout (datetime.timedelta, optional) – Тайм-аут для monitored_barrier. Если None, будет использоваться тайм-аут группы процессов по умолчанию.
  • wait_all_ranks (bool, optional) – Следует ли собирать информацию обо всех сбойных рангах. По умолчанию задано False, и monitored_barrier на ранге 0 выдаст исключение при обнаружении первого сбойного ранга, чтобы быстрее сообщить об ошибке. При задании wait_all_ranks=True monitored_barrier соберёт информацию обо всех сбойных рангах и выдаст ошибку с данными о каждом из них.
Возвращает:

None.

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

None

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> if dist.get_rank() != 1:
>>>     dist.monitored_barrier() # Raises exception indicating that
>>> # rank 1 did not call into monitored_barrier.
>>> # Example with wait_all_ranks=True
>>> if dist.get_rank() == 0:
>>>     dist.monitored_barrier(wait_all_ranks=True) # Raises exception
>>> # indicating that ranks 1, 2, ... world_size - 1 did not call into
>>> # monitored_barrier.
class torch.distributed.Work

Объект Work представляет дескриптор ожидающей выполнения асинхронной операции в распределённом пакете PyTorch. Он возвращается неблокирующими коллективными операциями, например dist.all_reduce(tensor, async_op=True).

block_current_stream(self: torch._C._distributed_c10d.Work) → None

Блокирует текущий активный поток GPU до завершения операции. Для коллективных операций на GPU это эквивалентно синхронизации. Для операций, инициированных на CPU, например с использованием Gloo, поток CUDA будет заблокирован до завершения операции.

Во всех случаях функция возвращает управление немедленно.

Чтобы проверить, успешно ли выполнена операция, следует асинхронно проверить результат объекта Work.

boxed(self: torch._C._distributed_c10d.Work) → object
exception(self: torch._C._distributed_c10d.Work) → object
get_future(self: torch._C._distributed_c10d.Work) → torch.Future
Возвращает:

Объект torch.futures.Future, связанный с завершением Work. Например, объект future можно получить с помощью fut = process_group.allreduce(tensors).get_future().

Пример::

Ниже приведён пример простого коммуникационного хука DDP allreduce, который использует API get_future для получения Future, связанного с завершением allreduce.

>>> def allreduce(process_group: dist.ProcessGroup, bucket: dist.GradBucket): -> torch.futures.Future
>>>     group_to_use = process_group if process_group is not None else torch.distributed.group.WORLD
>>>     tensor = bucket.buffer().div_(group_to_use.size())
>>>     return torch.distributed.all_reduce(tensor, group=group_to_use, async_op=True).get_future()
>>> ddp_model.register_comm_hook(state=None, hook=allreduce)

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

API get_future поддерживает NCCL, а также частично бэкенды GLOO и MPI (без поддержки операций между узлами, таких как send/recv) и возвращает torch.futures.Future.

В приведённом выше примере allreduce работа будет выполняться на GPU с использованием бэкенда NCCL. fut.wait() вернёт управление после синхронизации соответствующих потоков NCCL с текущими потоками устройств PyTorch, чтобы обеспечить асинхронное выполнение CUDA; при этом функция не дожидается полного завершения операции на GPU. Обратите внимание, что CUDAFuture не поддерживает флаг TORCH_NCCL_BLOCKING_WAIT или barrier() из NCCL. Кроме того, если с помощью fut.then() добавлена функция обратного вызова, выполнение будет ожидать синхронизации потоков NCCL для WorkNCCL с выделенным потоком обратных вызовов ProcessGroupNCCL, а затем вызовет функцию обратного вызова напрямую после её выполнения в потоке обратных вызовов. fut.then() вернёт другой объект CUDAFuture, содержащий возвращаемое значение функции обратного вызова, и объект CUDAEvent, записавший поток обратных вызовов.

  1. Для работы на CPU fut.done() возвращает true, когда работа завершена и тензоры value() готовы.
  2. Для работы на GPU fut.done() возвращает true только в том случае, если операция была поставлена в очередь.
  3. Для смешанной работы CPU-GPU (например, при отправке тензоров GPU с помощью GLOO) fut.done() возвращает true, когда тензоры поступили на соответствующие узлы, но ещё не обязательно синхронизированы на соответствующих GPU (как и для работы на GPU).
get_future_result(self: torch._C._distributed_c10d.Work) → torch.Future
Возвращает:

Объект torch.futures.Future типа int, соответствующий типу перечисления WorkResult. Например, объект future можно получить с помощью fut = process_group.allreduce(tensor).get_future_result().

Пример::

Пользователи могут использовать fut.wait() для блокирующего ожидания завершения работы и получения WorkResult с помощью fut.value(). Также пользователи могут использовать fut.then(call_back_func) для регистрации функции обратного вызова, которая будет вызвана по завершении работы без блокировки текущего потока.

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

API get_future_result поддерживает NCCL.

is_completed(self: torch._C._distributed_c10d.Work) → bool
is_success(self: torch._C._distributed_c10d.Work) → bool
result(self: torch._C._distributed_c10d.Work) → list[torch.Tensor]
source_rank(self: torch._C._distributed_c10d.Work) → int
synchronize(self: torch._C._distributed_c10d.Work) → None
static unbox(arg0: object) → torch._C._distributed_c10d.Work
wait(self: torch._C._distributed_c10d.Work, timeout: datetime.timedelta = datetime.timedelta(0)) → bool
Возвращает:

true/false.

Пример::
try:

work.wait(timeout)

except:

# some handling

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

В обычных случаях пользователям не нужно задавать тайм-аут. Вызов wait() эквивалентен вызову synchronize(): текущий поток блокируется до завершения работы NCCL. Однако, если задан тайм-аут, поток CPU будет заблокирован до завершения работы NCCL или истечения времени ожидания. При истечении времени ожидания будет выброшено исключение.

class torch.distributed.ReduceOp

Класс, подобный перечислению, для доступных операций редукции: SUM, PRODUCT, MIN, MAX, BAND, BOR, BXOR и PREMUL_SUM.

Редукции BAND, BOR и BXOR недоступны при использовании бэкенда NCCL.

AVG делит значения на размер world перед суммированием по рангам. AVG доступен только с бэкендом NCCL и только для версий NCCL 2.10 или новее.

PREMUL_SUM локально умножает входные данные на заданный скаляр перед редукцией. PREMUL_SUM доступен с бэкендом NCCL (версии NCCL 2.11 или новее) и бэкендом XCCL. Его можно использовать, вызвав ReduceOp.PREMUL_SUM(factor), где factor — число с плавающей точкой или тензор из одного элемента.

Кроме того, MAX, MIN и PRODUCT не поддерживаются для комплексных тензоров.

Значения этого класса доступны как атрибуты, например ReduceOp.SUM. Они используются для указания стратегий коллективных операций редукции, например reduce().

Этот класс не поддерживает свойство __members__.

class RedOpType

Элементы:

SUM

AVG

PRODUCT

MIN

MAX

BAND

BOR

BXOR

PREMUL_SUM

property name
property factor

Множитель операции ReduceOp PREMUL_SUM.

class torch.distributed.reduce_op

Устаревший класс, подобный перечислению, для операций редукции: SUM, PRODUCT, MIN и MAX.

Рекомендуется использовать ReduceOp.

Распределённое хранилище «ключ-значение»

Пакет distributed включает распределённое хранилище «ключ-значение», которое можно использовать для обмена информацией между процессами в группе, а также для инициализации пакета distributed в torch.distributed.init_process_group() (явно создав хранилище вместо указания init_method). Существует 3 варианта хранилищ «ключ-значение»: TCPStore, FileStore и HashStore.

class torch.distributed.Store

Базовый класс для всех реализаций хранилищ, например для 3 хранилищ, предоставляемых PyTorch distributed: (TCPStore, FileStore и HashStore).

__init__(self: torch._C._distributed_c10d.Store) → None
add(self: torch._C._distributed_c10d.Store, arg0: str, arg1: SupportsInt | SupportsIndex) → int

При первом вызове add для заданного key в хранилище создаётся связанный с key счётчик, которому присваивается начальное значение amount. Последующие вызовы add с тем же key увеличивают счётчик на указанное amount. Вызов add() с ключом, который уже был записан в хранилище с помощью set(), приведёт к исключению.

Параметры:
  • key (str) – Ключ хранилища, счётчик которого будет увеличен.
  • amount (int) – Величина, на которую будет увеличен счётчик.
Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Using TCPStore as an example, other store types can also be used
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.add("first_key", 1)
>>> store.add("first_key", 6)
>>> # Should return 7
>>> store.get("first_key")
append(self: torch._C._distributed_c10d.Store, arg0: str, arg1: str) → None

Добавляет в хранилище пару «ключ-значение» на основе переданных key и value. Если key отсутствует в хранилище, он будет создан.

Параметры:
  • key (str) – Ключ, который нужно добавить в хранилище.
  • value (str) – Значение, связанное с key, которое нужно добавить в хранилище.
Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.append("first_key", "po")
>>> store.append("first_key", "tato")
>>> # Should return "potato"
>>> store.get("first_key")
barrier(self: torch._C._distributed_c10d.Store, key: str, world_size: SupportsInt | SupportsIndex, timeout: datetime.timedelta | None = None) → None

Операция барьера, которая блокирует выполнение, пока её не вызовут world_size рабочих с одним и тем же key. Если timeout не задан, используется время ожидания по умолчанию для хранилища.

Параметры:
  • key (str) – Уникальный ключ этого экземпляра барьера.
  • world_size (int) – Число рабочих, которые должны вызвать барьер, прежде чем он разблокируется.
  • timeout (timedelta, необязательный) – Время ожидания перед возникновением исключения. По умолчанию используется время ожидания хранилища.
Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> # This will return immediately since world_size=1
>>> store.barrier("my_barrier", 1)
>>> store.barrier("my_barrier2", 1, timedelta(seconds=10))
check(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str]) → bool

Проверяет, сохранены ли в хранилище значения для заданного списка keys. В обычных случаях вызов возвращает результат сразу, однако возможны некоторые редкие случаи взаимной блокировки, например вызов check после уничтожения TCPStore. Вызов check() со списком ключей проверяет, сохранены ли они в хранилище.

Параметры:

keys (list[str]) – Ключи, наличие которых в хранилище нужно проверить.

Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Using TCPStore as an example, other store types can also be used
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.add("first_key", 1)
>>> # Should return 7
>>> store.check(["first_key"])
clone(self: torch._C._distributed_c10d.Store) → torch._C._distributed_c10d.Store

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

compare_set(self: torch._C._distributed_c10d.Store, arg0: str, arg1: str, arg2: str) → bytes

Добавляет пару «ключ-значение» в хранилище на основе переданного key и перед добавлением сравнивает expected_value и desired_value. desired_value будет задано, только если expected_value для key уже существует в хранилище или если expected_value — пустая строка.

Параметры:
  • key (str) – Ключ, который нужно проверить в хранилище.
  • expected_value (str) – Значение, связанное с key, которое нужно проверить перед добавлением.
  • desired_value (str) – Значение, связанное с key, которое нужно добавить в хранилище.
Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.set("key", "first_value")
>>> store.compare_set("key", "first_value", "second_value")
>>> # Should return "second_value"
>>> store.get("key")
delete_key(self: torch._C._distributed_c10d.Store, arg0: str) → bool

Удаляет из хранилища пару «ключ-значение», связанную с key. Возвращает true, если ключ был успешно удалён, и false, если нет.

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

API delete_key поддерживается только в TCPStore и HashStore. Использование этого API с FileStore приведёт к исключению.

Параметры:

key (str) – Ключ, который нужно удалить из хранилища

Возвращает:

True, если key был удалён, в противном случае — False.

Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Using TCPStore as an example, HashStore can also be used
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.set("first_key")
>>> # This should return true
>>> store.delete_key("first_key")
>>> # This should return false
>>> store.delete_key("bad_key")
get(self: torch._C._distributed_c10d.Store, arg0: str) → bytes

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

Параметры:

key (str) – Функция вернёт значение, связанное с этим ключом.

Возвращает:

Значение, связанное с key, если key присутствует в хранилище.

Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.set("first_key", "first_value")
>>> # Should return "first_value"
>>> store.get("first_key")
has_extended_api(self: torch._C._distributed_c10d.Store) → bool

Возвращает true, если хранилище поддерживает расширенные операции.

list_keys(self: torch._C._distributed_c10d.Store) → list[str]

Возвращает список всех ключей в хранилище.

multi_get(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str]) → list[bytes]

Получает все значения из keys. Если какой-либо ключ из keys отсутствует в хранилище, функция будет ждать timeout

Параметры:

keys (List[str]) – Ключи, значения которых нужно получить из хранилища.

Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.set("first_key", "po")
>>> store.set("second_key", "tato")
>>> # Should return [b"po", b"tato"]
>>> store.multi_get(["first_key", "second_key"])
multi_set(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str], arg1: collections.abc.Sequence[str]) → None

Добавляет в хранилище список пар «ключ-значение» на основе переданных keys и values

Параметры:
  • keys (List[str]) – Ключи для добавления.
  • values (List[str]) – Значения для добавления.
Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.multi_set(["first_key", "second_key"], ["po", "tato"])
>>> # Should return b"po"
>>> store.get("first_key")
num_keys(self: torch._C._distributed_c10d.Store) → int

Возвращает число ключей, заданных в хранилище. Обратите внимание: обычно это число на единицу больше числа ключей, добавленных с помощью set() и add(), поскольку один ключ используется для координации всех рабочих, использующих хранилище.

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

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

Возвращает:

Число ключей, имеющихся в хранилище.

Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Using TCPStore as an example, other store types can also be used
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.set("first_key", "first_value")
>>> # This should return 2
>>> store.num_keys()
queue_len(self: torch._C._distributed_c10d.Store, arg0: str) → int

Возвращает длину указанной очереди.

Если очередь не существует, возвращает 0.

Подробнее см. queue_push.

Параметры:

key (str) – Ключ очереди, длину которой нужно получить.

queue_pop(self: torch._C._distributed_c10d.Store, key: str, block: bool = True) → bytes

Извлекает значение из указанной очереди или ждёт до истечения времени ожидания, если очередь пуста.

Подробнее см. queue_push.

Если block равно False, при пустой очереди будет вызвана ошибка dist.QueueEmptyError.

Параметры:
  • key (str) – Ключ очереди, из которой нужно извлечь элемент.
  • block (bool) – Нужно ли блокировать выполнение в ожидании ключа или сразу вернуть результат.
queue_push(self: torch._C._distributed_c10d.Store, arg0: str, arg1: str) → None

Помещает значение в указанную очередь.

Использование одного ключа для очередей и операций set/get может привести к непредсказуемому поведению.

Для очередей поддерживаются операции wait/check.

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

Параметры:
  • key (str) – Ключ очереди, в которую нужно поместить элемент.
  • value (str) – Значение, которое нужно поместить в очередь.
set(self: torch._C._distributed_c10d.Store, arg0: str, arg1: str) → None

Добавляет в хранилище пару «ключ-значение» на основе переданных key и value. Если key уже существует в хранилище, старое значение будет заменено новым переданным value.

Параметры:
  • key (str) – Ключ, который нужно добавить в хранилище.
  • value (str) – Значение, связанное с key, которое нужно добавить в хранилище.
Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.set("first_key", "first_value")
>>> # Should return "first_value"
>>> store.get("first_key")
set_timeout(self: torch._C._distributed_c10d.Store, arg0: datetime.timedelta) → None

Задаёт время ожидания по умолчанию для хранилища. Это время ожидания используется при инициализации, а также в wait() и get().

Параметры:

timeout (timedelta) – Время ожидания, задаваемое для хранилища.

Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Using TCPStore as an example, other store types can also be used
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> store.set_timeout(timedelta(seconds=10))
>>> # This will throw an exception after 10 seconds
>>> store.wait(["bad_key"])
property timeout

Возвращает время ожидания хранилища.

wait(*args, **kwargs)

Перегруженная функция.

  1. wait(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str]) -> None

Ожидает добавления в хранилище каждого ключа из keys. Если до истечения timeout (заданного при инициализации хранилища) будут добавлены не все ключи, wait вызовет исключение.

Параметры:

keys (list) – Список ключей, добавления которых в хранилище нужно дождаться.

Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Using TCPStore as an example, other store types can also be used
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> # This will throw an exception after 30 seconds
>>> store.wait(["bad_key"])
  1. wait(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str], arg1: datetime.timedelta) -> None

Ожидает добавления в хранилище каждого ключа из keys и вызывает исключение, если ключи не были добавлены до истечения переданного timeout.

Параметры:
  • keys (list) – Список ключей, добавления которых в хранилище нужно дождаться.
  • timeout (timedelta) – Время ожидания добавления ключей перед вызовом исключения.
Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Using TCPStore as an example, other store types can also be used
>>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30))
>>> # This will throw an exception after 10 seconds
>>> store.wait(["bad_key"], timedelta(seconds=10))
class torch.distributed.TCPStore

Реализация распределённого хранилища «ключ-значение» на основе TCP. Серверное хранилище содержит данные, а клиентские хранилища могут подключаться к нему по TCP и выполнять такие действия, как set() для добавления пары «ключ-значение», get() для получения пары «ключ-значение» и т. д. Всегда должно быть инициализировано одно серверное хранилище, поскольку клиентские хранилища будут ждать подключения к серверу.

Параметры:
  • host_name (str) – Имя хоста или IP-адрес, на котором должно работать серверное хранилище.
  • port (int) – Порт, на котором серверное хранилище должно ожидать входящие запросы.
  • world_size (int, необязательный) – Общее число пользователей хранилища (число клиентов + 1 для сервера). Значение по умолчанию — None (None означает, что число пользователей хранилища не фиксировано).
  • is_master (bool, необязательный) – True при инициализации серверного хранилища и False для клиентских хранилищ. Значение по умолчанию — False.
  • timeout (timedelta, необязательный) – Время ожидания, используемое хранилищем при инициализации и для таких методов, как get() и wait(). Значение по умолчанию — timedelta(seconds=300)
  • wait_for_workers (bool, необязательный) – Нужно ли ждать подключения всех рабочих к серверному хранилищу. Применимо только при фиксированном значении world_size. Значение по умолчанию — True.
  • multi_tenant (bool, необязательный) – Если значение True, все экземпляры TCPStore в текущем процессе с одинаковыми host/port будут использовать один и тот же базовый TCPServer. Значение по умолчанию — False.
  • master_listen_fd (int, необязательный) – Если указано, базовый TCPServer будет ожидать подключения на этом файловом дескрипторе, который должен быть сокетом, уже привязанным к port. Чтобы привязать эфемерный порт, рекомендуется задать port равным 0 и считать .port. Значение по умолчанию — None (сервер создаёт новый сокет и пытается привязать его к port).
  • use_libuv (bool, необязательный) – Если значение True, для серверной части TCPServer используется libuv. Значение по умолчанию — True.
Пример::
>>> import torch.distributed as dist
>>> from datetime import timedelta
>>> # Run on process 1 (server)
>>> server_store = dist.TCPStore("127.0.0.1", 1234, 2, True, timedelta(seconds=30))
>>> # Run on process 2 (client)
>>> client_store = dist.TCPStore("127.0.0.1", 1234, 2, False)
>>> # Use any of the store methods from either the client or server after initialization
>>> server_store.set("first_key", "first_value")
>>> client_store.get("first_key")
__init__(self: torch._C._distributed_c10d.TCPStore, host_name: str, port: SupportsInt | SupportsIndex, world_size: SupportsInt | SupportsIndex | None = None, is_master: bool = False, timeout: datetime.timedelta = datetime.timedelta(seconds=300), wait_for_workers: bool = True, multi_tenant: bool = False, master_listen_fd: SupportsInt | SupportsIndex | None = None, use_libuv: bool = True) → None

Создаёт новый TCPStore.

property host

Возвращает имя хоста, на котором хранилище ожидает запросы.

property libuvBackend

Возвращает True, если используется серверная часть libuv.

property port

Возвращает номер порта, на котором хранилище ожидает запросы.

class torch.distributed.HashStore

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

Пример::
>>> import torch.distributed as dist
>>> store = dist.HashStore()
>>> # store can be used from other threads
>>> # Use any of the store methods after initialization
>>> store.set("first_key", "first_value")
__init__(self: torch._C._distributed_c10d.HashStore) → None

Создаёт новое хранилище HashStore.

class torch.distributed.FileStore

Реализация хранилища, в которой пары «ключ-значение» сохраняются в файле.

Параметры:
  • file_name (str) – путь к файлу, в котором будут храниться пары «ключ-значение»
  • world_size (int, необязательный) – Общее число процессов, использующих хранилище. Значение по умолчанию — -1 (отрицательное значение означает, что число пользователей хранилища не фиксировано).
Пример::
>>> import torch.distributed as dist
>>> store1 = dist.FileStore("/tmp/filestore", 2)
>>> store2 = dist.FileStore("/tmp/filestore", 2)
>>> # Use any of the store methods from either the client or server after initialization
>>> store1.set("first_key", "first_value")
>>> store2.get("first_key")
__init__(self: torch._C._distributed_c10d.FileStore, file_name: str, world_size: SupportsInt | SupportsIndex = -1) → None

Создаёт новое хранилище FileStore.

property path

Возвращает путь к файлу, который FileStore использует для хранения пар «ключ-значение».

class torch.distributed.PrefixStore

Обёртка над любым из 3 хранилищ «ключ-значение» (TCPStore, FileStore и HashStore), которая добавляет префикс к каждому ключу, добавляемому в хранилище.

Параметры:
  • prefix (str) – Строка-префикс, добавляемая перед каждым ключом при его помещении в хранилище.
  • store (torch.distributed.store) – Объект хранилища, образующий базовое хранилище «ключ-значение».
__init__(self: torch._C._distributed_c10d.PrefixStore, prefix: str, store: torch._C._distributed_c10d.Store) → None

Создаёт новое хранилище PrefixStore.

property underlying_store

Возвращает базовый объект хранилища, обёрнутый в PrefixStore.

Профилирование коллективных коммуникаций

Обратите внимание: для профилирования коллективных коммуникаций и API точка-точка, описанных здесь, можно использовать torch.profiler (рекомендуется, доступно только начиная с версии 1.8.1) или torch.autograd.profiler. Поддерживаются все встроенные серверные части (gloo, nccl, mpi), а использование коллективных коммуникаций будет корректно отображаться в результатах профилирования и трассировках. Профилировать код можно так же, как любой обычный оператор torch:

import torch
import torch.distributed as dist
with torch.profiler():
    tensor = torch.randn(20, 10)
    dist.all_reduce(tensor)

Полный обзор возможностей профилировщика см. в документации по профилировщику.

torch.distributed.distributed_c10d.record_comm(name) [исходный код]

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

Во время его работы все коллективные операции c10d, вызванные в этом контексте, будут использовать name в качестве заголовка профилирования в базовом классе Work вместо имени по умолчанию, заданного для конкретной серверной части (например, nccl:all_reduce). Работает со всеми серверными частями без изменений для отдельных серверных частей или коллективных операций.

Параметры:

name (str) – Имя профилирования, связываемое с коллективными операциями.

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

Iterator[None]

Пример::
>>> with dist.record_comm("FSDP::all_gather (layer1)"):
...     dist.all_gather_into_tensor(output, input, group=pg)

Оптимизация с симметричной памятью

Симметричные ядра NCCL

В NCCL 2.27 и более поздних версиях доступны ядра устройств, специально разработанные для симметричных буферов, зарегистрированных в окне. В них используются алгоритмы с низкой задержкой, multimem/NVLS и TMA вместо универсального алгоритма на основе кольца или дерева. all_reduce, all_gather_into_tensor и reduce_scatter_tensor автоматически используют их после регистрации буферов — место вызова не меняется. Если буферы не зарегистрированы или для сочетания операции и типа данных нет симметричной реализации, NCCL незаметно переключается на обычный путь.

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

Коллективные операции с использованием движков копирования

Когда коллективные операции NCCL выполняются над тензорами в симметричной памяти с политикой zero-CTA, перемещение данных передаётся движкам копирования графического процессора (движкам DMA), а не потоковым мультипроцессорам CUDA (SM). Это освобождает SM для вычислений и позволяет лучше перекрывать коммуникацию и вычисления.

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

Редукция с повышенной точностью

Когда коллективные операции NCCL, такие как reduce_scatter и all_reduce, работают с тензорами в симметричной памяти, реализация симметричного ядра NCCL автоматически выполняет внутреннюю редукцию с повышенной точностью (например, BF16/FP16 → накопление в FP32 → BF16/FP16). Это повышает численную точность без изменения кода вызова коллективной операции.

Подробную информацию об области применения, поддерживаемых доменах и требованиях к версиям см. в разделе Редукция с повышенной точностью документации по симметричной памяти.

Коллективные функции для нескольких GPU

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

Функции для нескольких GPU (то есть нескольких GPU на один поток CPU) объявлены устаревшими. На данный момент предпочтительная модель программирования PyTorch Distributed — одно устройство на поток, как показано в API этого документа. Если вы разрабатываете бэкенд и хотите поддерживать несколько устройств на поток, свяжитесь с сопровождающими PyTorch Distributed.

Коллективные операции с объектами

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

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

Коллективные операции с объектами — это набор операций, похожих на коллективные, которые работают с произвольными объектами Python, если их можно сериализовать с помощью pickle. Реализованы различные шаблоны коллективных операций (например, broadcast, all_gather, …), но в целом каждый из них выполняет следующие действия:

  1. преобразует входной объект в pickle (необработанные байты), а затем помещает его в байтовый тензор
  2. передаёт размер этого байтового тензора процессам-участникам (первая коллективная операция)
  3. выделяет тензор подходящего размера для выполнения основной коллективной операции
  4. передаёт данные объекта (вторая коллективная операция)
  5. преобразует необработанные данные обратно в объект Python (десериализует pickle)

Коллективные операции с объектами иногда имеют неожиданные характеристики производительности или использования памяти, из-за которых выполнение может занимать много времени или приводить к ошибкам нехватки памяти (OOM), поэтому использовать их следует с осторожностью. Ниже приведены распространённые проблемы.

Асимметричное время сериализации и десериализации pickle — сериализация объектов с помощью pickle может занимать много времени в зависимости от количества, типа и размера объектов. Если коллективная операция выполняется по схеме «многие к одному» (например, gather_object), принимающий процесс или процессы должны десериализовать в N раз больше объектов, чем отправляющие процессы сериализовали. Это может привести к тайм-ауту других процессов при выполнении следующей коллективной операции.

Неэффективная передача тензоров — тензоры следует передавать с помощью обычных API коллективных операций, а не API коллективных операций с объектами. Передача тензоров через API коллективных операций с объектами возможна, но они будут сериализованы и десериализованы (включая синхронизацию с CPU и копирование с устройства на хост для тензоров не на CPU). Почти во всех случаях, кроме отладки или устранения неполадок в коде, стоит приложить усилия и переработать код, чтобы вместо этого использовать коллективные операции без объектов.

Неожиданные устройства тензоров — если вы всё же хотите передавать тензоры через коллективные операции с объектами, следует учитывать ещё одну особенность тензоров CUDA (и, возможно, других ускорителей). Если сериализовать с помощью pickle тензор, который находится на cuda:3, а затем десериализовать его, получится другой тензор на cuda:3 независимо от того, в каком процессе вы находитесь и какое устройство CUDA является для него устройством «по умолчанию». При использовании обычных API коллективных операций с тензорами «выходные тензоры» всегда находятся на том же локальном устройстве, чего обычно и следует ожидать.

Десериализация тензора неявно активирует контекст CUDA, если процесс впервые использует GPU, что может привести к значительному расходу памяти GPU. Этой проблемы можно избежать, переместив тензоры на CPU до передачи их в качестве входных данных коллективной операции с объектами.

Сторонние бэкенды

Помимо встроенных бэкендов GLOO/MPI/NCCL, PyTorch Distributed поддерживает сторонние бэкенды с помощью механизма регистрации во время выполнения. Сведения о разработке стороннего бэкенда с помощью расширения C++ см. в разделах Учебные материалы — пользовательские расширения C++ и CUDA и test/cpp_extensions/cpp_c10d_extension.cpp. Возможности сторонних бэкендов определяются их собственными реализациями.

Новый бэкенд наследуется от c10d::ProcessGroup и при импорте регистрирует имя бэкенда и интерфейс его создания с помощью torch.distributed.Backend.register_backend().

При ручном импорте этого бэкенда и вызове torch.distributed.init_process_group() с соответствующим именем бэкенда пакет torch.distributed работает на новом бэкенде.

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

Поддержка сторонних бэкендов является экспериментальной и может измениться.

Бэкенд TorchComms

TorchComms — дополнительный коммуникационный бэкенд для torch.distributed. При его включении стандартное создание бэкендов в init_process_group() переопределяется: все группы процессов создаются через TorchComms, а не с помощью встроенных реализаций ProcessGroup.

Примечание

TorchComms является экспериментальным и должен устанавливаться отдельно. Чтобы перечисленные ниже флаги вступили в силу, пакет torchcomms должен быть доступен для импорта.

Включение TorchComms

Перед вызовом init_process_group() задайте переменную окружения TORCH_DISTRIBUTED_USE_TORCHCOMMS:

export TORCH_DISTRIBUTED_USE_TORCHCOMMS=1

Или задайте флаг конфигурации программно:

import torch.distributed.config as dist_config

dist_config.use_torchcomms = True

Аргумент backend функции init_process_group() (например, "nccl", "gloo") по-прежнему учитывается — он передаётся в TorchComms, который выбирает соответствующий плагин поставщика. Другие изменения в коде приложения не требуются; все коллективные API torch.distributed продолжают работать как прежде.

Поведение при включении

При включённом TorchComms функция init_process_group() меняет способ создания бэкендов для каждой пары устройство/бэкенд в группе процессов (за исключением бэкенда fake, который всегда обрабатывается встроенными средствами):

  1. Коммуникатор TorchComms создаётся с помощью torchcomms.new_comm() с использованием запрошенной строки бэкенда и устройства.
  2. Коммуникатор оборачивается в _BackendWrapper, реализующий интерфейс C++ c10d::Backend, и становится непосредственной заменой встроенных бэкендов ProcessGroup.
  3. Для коммуникатора автоматически регистрируется FlightRecorderHook. Этот обработчик учитывает переменные окружения TORCH_FR_BUFFER_SIZE и TORCH_NCCL_TRACE_BUFFER_SIZE, задающие размер буфера трассировки.
  4. destroy_process_group() вызывает finalize() для коммуникаторов TorchComms при очистке.
  5. split_group() создаёт подкоммуникаторы с помощью встроенного механизма разделения TorchComms, а не создаёт новую группу процессов с нуля.

Немедленная инициализация

Коммуникаторы TorchComms инициализируются немедленно во время вызова init_process_group() и поддерживают только одно устройство бэкенда на группу. Аргумент device_id необходимо указать во время инициализации:

dist.init_process_group(backend="nccl", device_id=torch.device("cuda", local_rank))

Параллельное выполнение операций «точка-точка»

Каждая группа процессов TorchComms соответствует одному базовому коммуникатору. Для операций «точка-точка» (send/recv), отправленных в одной группе и одном потоке, не гарантируется параллельное выполнение. Код, зависящий от параллельного выполнения операций «точка-точка», должен использовать один из следующих вариантов:

  • Использовать пакетные API P2P (batch_isend_irecv()) или
  • Выполнять операции в отдельных группах или коммуникаторах.

Утилита запуска

Пакет torch.distributed также предоставляет утилиту запуска в torch.distributed.launch. Эта вспомогательная утилита позволяет запускать несколько процессов на узел для распределённого обучения.

Модуль torch.distributed.launch.

torch.distributed.launch — это модуль, который запускает несколько процессов распределённого обучения на каждом обучающем узле.

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

Этот модуль будет объявлен устаревшим в пользу torchrun.

Эту утилиту можно использовать для распределённого обучения на одном узле, запуская один или несколько процессов на узел. Она подходит как для обучения на CPU, так и для обучения на GPU. При обучении на GPU каждый распределённый процесс будет работать на одном GPU. Это позволяет значительно повысить производительность обучения на одном узле. Утилиту также можно использовать для распределённого обучения на нескольких узлах, запуская несколько процессов на каждом узле для повышения производительности обучения на нескольких узлах. Это особенно полезно для систем с несколькими интерфейсами Infiniband, поддерживающими прямой доступ к GPU, поскольку их все можно задействовать для увеличения совокупной пропускной способности коммуникаций.

В обоих случаях — при распределённом обучении на одном узле или на нескольких узлах — эта утилита запускает заданное количество процессов на каждом узле (--nproc-per-node). При обучении на GPU это число должно быть меньше или равно количеству GPU в текущей системе (nproc_per_node), и каждый процесс будет работать на одном GPU из диапазона от GPU 0 до GPU (nproc_per_node - 1).

Использование этого модуля:

  1. Многопроцессное распределённое обучение на одном узле
python -m torch.distributed.launch --nproc-per-node=NUM_GPUS_YOU_HAVE
           YOUR_TRAINING_SCRIPT.py (--arg1 --arg2 --arg3 and all other
           arguments of your training script)
  1. Многопроцессное распределённое обучение на нескольких узлах (например, на двух узлах)

Узел 1: (IP: 192.168.1.1, свободный порт: 1234)

python -m torch.distributed.launch --nproc-per-node=NUM_GPUS_YOU_HAVE
           --nnodes=2 --node-rank=0 --master-addr="192.168.1.1"
           --master-port=1234 YOUR_TRAINING_SCRIPT.py (--arg1 --arg2 --arg3
           and all other arguments of your training script)

Узел 2:

python -m torch.distributed.launch --nproc-per-node=NUM_GPUS_YOU_HAVE
           --nnodes=2 --node-rank=1 --master-addr="192.168.1.1"
           --master-port=1234 YOUR_TRAINING_SCRIPT.py (--arg1 --arg2 --arg3
           and all other arguments of your training script)
  1. Чтобы узнать, какие необязательные аргументы поддерживает этот модуль:
python -m torch.distributed.launch --help

Важные замечания:

1. На данный момент эта утилита и многопроцессное распределённое обучение на GPU (на одном или нескольких узлах) обеспечивают наилучшую производительность только при использовании бэкенда распределённых вычислений NCCL. Поэтому для обучения на GPU рекомендуется использовать бэкенд NCCL.

2. В программе обучения необходимо обрабатывать аргумент командной строки --local-rank=LOCAL_PROCESS_RANK, который будет передан этим модулем. Если программа обучения использует GPU, следует убедиться, что код выполняется только на GPU-устройстве с номером LOCAL_PROCESS_RANK. Это можно сделать следующим образом:

Разобрать аргумент local_rank

>>> import argparse
>>> parser = argparse.ArgumentParser()
>>> parser.add_argument("--local-rank", "--local_rank", type=int)
>>> args = parser.parse_args()

Назначить устройству номер локального ранга одним из следующих способов:

>>> torch.cuda.set_device(args.local_rank)  # before your code runs

или

>>> with torch.cuda.device(args.local_rank):
>>>    # your code to run
>>>    ...

Изменено в версии 2.0.0: Средство запуска передаёт скрипту аргумент --local-rank=<rank>. Начиная с PyTorch 2.0.0, предпочтителен вариант с дефисами --local-rank, а не ранее использовавшийся вариант с подчёркиваниями --local_rank.

Для обеспечения обратной совместимости может потребоваться обрабатывать оба варианта при разборе аргументов. Это означает, что в анализатор аргументов нужно добавить как "--local-rank", так и "--local_rank". Если указан только "--local_rank", средство запуска выдаст ошибку: «error: unrecognized arguments: –local-rank=<rank>». Для кода обучения, поддерживающего только PyTorch 2.0.0 и более поздние версии, достаточно добавить "--local-rank".

3. В начале программы обучения необходимо вызвать следующую функцию, чтобы запустить распределённый бэкенд. Настоятельно рекомендуется использовать init_method=env://. Другие методы инициализации (например, tcp://) могут работать, однако env:// — единственный метод, официально поддерживаемый этим модулем.

>>> torch.distributed.init_process_group(backend='YOUR BACKEND',
>>>                                      init_method='env://')

4. В программе обучения можно использовать обычные распределённые функции или модуль torch.nn.parallel.DistributedDataParallel(). Если программа обучения использует GPU и вы хотите использовать модуль torch.nn.parallel.DistributedDataParallel(), настройте его следующим образом.

>>> model = torch.nn.parallel.DistributedDataParallel(model,
>>>                                                   device_ids=[args.local_rank],
>>>                                                   output_device=args.local_rank)

Убедитесь, что аргумент device_ids задан как единственный идентификатор GPU-устройства, на котором будет работать код. Обычно это локальный ранг процесса. Иными словами, для использования этой утилиты значение device_ids должно быть [args.local_rank], а output_device должно быть args.local_rank.

5. Ещё один способ передать local_rank дочерним процессам — использовать переменную окружения LOCAL_RANK. Это поведение включается при запуске скрипта с --use-env=True. В приведённом выше примере дочернего процесса замените args.local_rank на os.environ['LOCAL_RANK']; при указании этого флага средство запуска не будет передавать --local-rank.

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

Значение local_rank НЕ является глобально уникальным: оно уникально только для процесса на одной машине. Поэтому не используйте его, чтобы, например, решить, следует ли выполнять запись в сетевую файловую систему. Пример того, что может пойти не так при неправильном использовании, см. в pytorch/pytorch#12042.

torch.distributed.launch.launch(args) [исходный код]
torch.distributed.launch.main(args=None) [исходный код]
torch.distributed.launch.parse_args(args) [исходный код]

Утилита создания процессов

Пакет Пакет многопроцессной обработки — torch.multiprocessing также предоставляет функцию spawn в torch.multiprocessing.spawn(). Эта вспомогательная функция позволяет создавать несколько процессов. Для этого ей передаётся функция, которую нужно выполнить, и запускается N процессов для её выполнения. Её также можно использовать для многопроцессного распределённого обучения.

Инструкции по использованию см. в примере PyTorch — реализации ImageNet

Обратите внимание, что для этой функции требуется Python версии 3.4 или выше.

Отладка приложений torch.distributed

Отладка распределённых приложений может быть непростой задачей из-за труднообъяснимых зависаний, сбоев или несогласованного поведения между рангами. torch.distributed предоставляет набор инструментов для самостоятельной отладки обучающих приложений:

Точка останова Python

Использовать отладчик Python в распределённой среде чрезвычайно удобно, но поскольку он не работает «из коробки», многие вообще им не пользуются. PyTorch предлагает специальную обёртку над pdb, которая упрощает этот процесс.

torch.distributed.breakpoint упрощает этот процесс. Внутри он изменяет поведение точки останова pdb двумя способами, но в остальном работает как обычный pdb.

  1. Подключает отладчик только на одном ранге (указанном пользователем).
  2. Останавливает все остальные ранги с помощью torch.distributed.barrier(), который освободит их после того, как отлаживаемый ранг выполнит continue
  3. Перенаправляет stdin дочернего процесса так, чтобы он был подключён к вашему терминалу.

Чтобы воспользоваться этим, просто выполните torch.distributed.breakpoint(rank) на всех рангах, используя одинаковое значение rank в каждом случае.

Контролируемый барьер

Начиная с версии v1.10, torch.distributed.monitored_barrier() доступен в качестве альтернативы torch.distributed.barrier(). При сбое он предоставляет полезную информацию о том, какой ранг мог стать причиной проблемы, например, если не все ранги вызывают torch.distributed.monitored_barrier() в течение заданного времени ожидания. torch.distributed.monitored_barrier() реализует барьер на стороне хоста с использованием примитивов взаимодействия send/recv в процессе, аналогичном подтверждениям, что позволяет рангу 0 сообщить, какие ранги не подтвердили прохождение барьера вовремя. Рассмотрим, например, следующую функцию, в которой ранг 1 не вызывает torch.distributed.monitored_barrier() (на практике это может быть вызвано ошибкой приложения или зависанием в предыдущей коллективной операции):

import os
from datetime import timedelta

import torch
import torch.distributed as dist
import torch.multiprocessing as mp


def worker(rank):
    dist.init_process_group("nccl", rank=rank, world_size=2)
    # monitored barrier requires gloo process group to perform host-side sync.
    group_gloo = dist.new_group(backend="gloo")
    if rank not in [1]:
        dist.monitored_barrier(group=group_gloo, timeout=timedelta(seconds=2))


if __name__ == "__main__":
    os.environ["MASTER_ADDR"] = "localhost"
    os.environ["MASTER_PORT"] = "29501"
    mp.spawn(worker, nprocs=2, args=())

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

RuntimeError: Rank 1 failed to pass monitoredBarrier in 2000 ms
 Original exception:
[gloo/transport/tcp/pair.cc:598] Connection closed by peer [2401:db00:eef0:1100:3560:0:1c05:25d]:8594

TORCH_DISTRIBUTED_DEBUG

В TORCH_CPP_LOG_LEVEL=INFO переменную среды TORCH_DISTRIBUTED_DEBUG можно использовать для включения дополнительных полезных журналов и проверок синхронизации коллективных операций, чтобы убедиться, что все ранги синхронизированы должным образом. Для TORCH_DISTRIBUTED_DEBUG можно задать значение OFF (по умолчанию), INFO или DETAIL в зависимости от требуемого уровня отладки. Обратите внимание, что наиболее подробный режим, DETAIL, может повлиять на производительность приложения, поэтому его следует использовать только при отладке проблем.

Значение TORCH_DISTRIBUTED_DEBUG=INFO включает дополнительные записи в журнал при инициализации моделей, обучаемых с помощью torch.nn.parallel.DistributedDataParallel(), а TORCH_DISTRIBUTED_DEBUG=DETAIL дополнительно записывает статистику производительности во время выполнения на выбранном числе итераций. Эта статистика включает такие данные, как время прямого прохода, время обратного прохода, время передачи градиентов и т. д. Рассмотрим следующее приложение:

import os

import torch
import torch.distributed as dist
import torch.multiprocessing as mp


class TwoLinLayerNet(torch.nn.Module):
    def __init__(self):
        super().__init__()
        self.a = torch.nn.Linear(10, 10, bias=False)
        self.b = torch.nn.Linear(10, 1, bias=False)

    def forward(self, x):
        a = self.a(x)
        b = self.b(x)
        return (a, b)


def worker(rank):
    dist.init_process_group("nccl", rank=rank, world_size=2)
    torch.cuda.set_device(rank)
    print("init model")
    model = TwoLinLayerNet().cuda()
    print("init ddp")
    ddp_model = torch.nn.parallel.DistributedDataParallel(model, device_ids=[rank])

    inp = torch.randn(10, 10).cuda()
    print("train")

    for _ in range(20):
        output = ddp_model(inp)
        loss = output[0] + output[1]
        loss.sum().backward()


if __name__ == "__main__":
    os.environ["MASTER_ADDR"] = "localhost"
    os.environ["MASTER_PORT"] = "29501"
    os.environ["TORCH_CPP_LOG_LEVEL"]="INFO"
    os.environ[
        "TORCH_DISTRIBUTED_DEBUG"
    ] = "DETAIL"  # set to DETAIL for runtime logging.
    mp.spawn(worker, nprocs=2, args=())

Во время инициализации выводятся следующие записи журнала:

I0607 16:10:35.739390 515217 logger.cpp:173] [Rank 0]: DDP Initialized with:
broadcast_buffers: 1
bucket_cap_bytes: 26214400
find_unused_parameters: 0
gradient_as_bucket_view: 0
is_multi_device_module: 0
iteration: 0
num_parameter_tensors: 2
output_device: 0
rank: 0
total_parameter_size_bytes: 440
world_size: 2
backend_name: nccl
bucket_sizes: 440
cuda_visible_devices: N/A
device_ids: 0
dtypes: float
master_addr: localhost
master_port: 29501
module_name: TwoLinLayerNet
nccl_async_error_handling: N/A
nccl_blocking_wait: N/A
nccl_debug: WARN
nccl_ib_timeout: N/A
nccl_nthreads: N/A
nccl_socket_ifname: N/A
torch_distributed_debug: INFO

Во время выполнения выводятся следующие записи журнала (если задано TORCH_DISTRIBUTED_DEBUG=DETAIL):

I0607 16:18:58.085681 544067 logger.cpp:344] [Rank 1 / 2] Training TwoLinLayerNet unused_parameter_size=0
 Avg forward compute time: 40838608
 Avg backward compute time: 5983335
Avg backward comm. time: 4326421
 Avg backward comm/comp overlap time: 4207652
I0607 16:18:58.085693 544066 logger.cpp:344] [Rank 0 / 2] Training TwoLinLayerNet unused_parameter_size=0
 Avg forward compute time: 42850427
 Avg backward compute time: 3885553
Avg backward comm. time: 2357981
 Avg backward comm/comp overlap time: 2234674

Кроме того, TORCH_DISTRIBUTED_DEBUG=INFO расширяет журналирование сбоев в torch.nn.parallel.DistributedDataParallel(), вызванных неиспользуемыми параметрами модели. В настоящее время, если в прямом проходе могут быть неиспользуемые параметры, при инициализации torch.nn.parallel.DistributedDataParallel() необходимо передать find_unused_parameters=True. Кроме того, начиная с версии v1.10, все выходные данные модели должны использоваться при вычислении функции потерь, поскольку torch.nn.parallel.DistributedDataParallel() не поддерживает неиспользуемые параметры в обратном проходе. Эти ограничения особенно сложны для больших моделей, поэтому при сбое с ошибкой torch.nn.parallel.DistributedDataParallel() записывает полные имена всех неиспользованных параметров. Например, если в приведённом выше приложении изменить loss так, чтобы вместо этого оно вычислялось как loss = output[1], то TwoLinLayerNet.a не получит градиент в обратном проходе, что приведёт к сбою DDP. При сбое пользователь получает сведения о неиспользованных параметрах, которые в больших моделях бывает сложно найти вручную:

RuntimeError: Expected to have finished reduction in the prior iteration before starting a new one. This error indicates that your module has parameters that were not used in producing loss. You can enable unused parameter detection by passing
 the keyword argument `find_unused_parameters=True` to `torch.nn.parallel.DistributedDataParallel`, and by
making sure all `forward` function outputs participate in calculating loss.
If you already have done the above, then the distributed data parallel module wasn't able to locate the output tensors in the return value of your module's `forward` function. Please include the loss function and the structure of the return va
lue of `forward` of your module when reporting this issue (e.g. list, dict, iterable).
Parameters which did not receive grad for rank 0: a.weight
Parameter indices which did not receive grad for rank 0: 0

Значение TORCH_DISTRIBUTED_DEBUG=DETAIL включает дополнительные проверки согласованности и синхронизации при каждом вызове коллективной операции, выполняемом пользователем напрямую или косвенно (например, allreduce DDP). Для этого создаётся группа процессов-обёртка, оборачивающая все группы процессов, возвращаемые API torch.distributed.init_process_group() и torch.distributed.new_group(). В результате эти API возвращают группу процессов-обёртку, которую можно использовать точно так же, как обычную группу процессов, но перед отправкой коллективной операции в базовую группу процессов она выполняет проверки согласованности. В настоящее время эти проверки включают torch.distributed.monitored_barrier(), которая гарантирует, что все ранги завершили незавершённые коллективные операции, и сообщает о рангах, которые зависли. Затем коллективная операция проверяется на согласованность: проверяется, что все коллективные функции совпадают и вызываются с согласованными формами тензоров. Если это условие не выполнено, при сбое приложения выводится подробный отчёт об ошибке вместо зависания или неинформативного сообщения. Рассмотрим, например, следующую функцию, в которой в torch.distributed.all_reduce() передаются входные данные с несовпадающими формами:

import torch
import torch.distributed as dist
import torch.multiprocessing as mp


def worker(rank):
    dist.init_process_group("nccl", rank=rank, world_size=2)
    torch.cuda.set_device(rank)
    tensor = torch.randn(10 if rank == 0 else 20).cuda()
    dist.all_reduce(tensor)
    torch.cuda.synchronize(device=rank)


if __name__ == "__main__":
    os.environ["MASTER_ADDR"] = "localhost"
    os.environ["MASTER_PORT"] = "29501"
    os.environ["TORCH_CPP_LOG_LEVEL"]="INFO"
    os.environ["TORCH_DISTRIBUTED_DEBUG"] = "DETAIL"
    mp.spawn(worker, nprocs=2, args=())

При использовании бэкенда NCCL такое приложение, скорее всего, зависнет, и в нетривиальных случаях причину зависания может быть сложно установить. Если пользователь включит TORCH_DISTRIBUTED_DEBUG=DETAIL и повторно запустит приложение, следующее сообщение об ошибке укажет первопричину:

work = default_pg.allreduce([tensor], opts)
RuntimeError: Error when verifying shape tensors for collective ALLREDUCE on rank 0. This likely indicates that input shapes into the collective are mismatched across ranks. Got shapes:  10
[ torch.LongTensor{1} ]

Примечание

Для точной настройки уровня отладки во время выполнения можно также использовать функции torch.distributed.set_debug_level(), torch.distributed.set_debug_level_from_env() и torch.distributed.get_debug_level().

Кроме того, TORCH_DISTRIBUTED_DEBUG=DETAIL можно использовать вместе с TORCH_SHOW_CPP_STACKTRACES=1 для записи полного стека вызовов при обнаружении рассинхронизации коллективных операций. Эти проверки рассинхронизации работают во всех приложениях, использующих коллективные вызовы c10d на основе групп процессов, созданных с помощью API torch.distributed.init_process_group() и torch.distributed.new_group().

HTTP-сервер отладки torch.distributed

Модуль torch.distributed.debug предоставляет HTTP-сервер для отладки распределённых приложений. Сервер можно запустить вызовом torch.distributed.debug.start_debug_server(). Это позволяет пользователям собирать данные со всех рабочих процессов во время выполнения.

torch.distributed.debug.start_debug_server(port=25999, worker_port=0, start_method=None, dump_dir=None, dump_interval=60.0, enabled_dumps=None, handlers=None, fetch_timeout=60.0) [исходный код]

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

Этот сервер предоставляет внешний HTTP-интерфейс для отладки медленных и заблокированных распределённых задач сразу на всех рангах. Он собирает такие данные, как трассировки стека, события FlightRecorder и профили производительности.

Для работы требуются зависимости, которые не устанавливаются по умолчанию.

Зависимости: - Jinja2 - aiohttp

ПРЕДУПРЕЖДЕНИЕ: Сервер предназначен только для использования в доверенных сетевых средах. Сервер отладки не защищён и не должен быть доступен из общедоступного Интернета. Подробную информацию см. в файле SECURITY.md.

ПРЕДУПРЕЖДЕНИЕ: Эта экспериментальная функция может измениться в любой момент.

Параметры:
  • port (int) – Порт для запуска внешнего сервера отладки.
  • worker_port (int) – Порт для запуска сервера рабочего процесса. По умолчанию равен 0, что приводит к привязке сервера рабочего процесса к эфемерному порту.
  • start_method (str | None) – Метод запуска многопроцессного режима для процесса внешнего сервера. Допустимые значения: “fork”, “spawn” или “forkserver”. Если задано None, используется метод запуска по умолчанию. При использовании CUDA или если важна безопасность fork, рекомендуется использовать “spawn”.
  • dump_dir (str | None) – Каталог для периодической записи отладочных дампов. Если задано None, периодическое создание дампов отключено.
  • dump_interval (float) – Интервал в секундах между периодическими дампами. По умолчанию равен 60.
  • enabled_dumps (set[str] | None) – Набор имён файлов дампов обработчиков, которые нужно включить (например, {“stacks”, “fr_trace”, “tcpstore”}). Если задано None, включаются все обработчики, реализующие dump().
  • handlers (list[DebugHandler] | None) – Список используемых обработчиков отладки. Если задано None, используются обработчики по умолчанию. Список обработчиков по умолчанию см. в torch.distributed.debug._handlers.
  • fetch_timeout (float) – Время ожидания в секундах при получении данных от отдельных рабочих процессов. По умолчанию равно 60. Рабочие процессы, не ответившие за это время, будут отмечены как недоступные.
torch.distributed.debug.stop_debug_server() [исходный код]

Завершает работу сервера отладки и останавливает процесс внешнего сервера отладки.

Ведение журналов

Помимо явной поддержки отладки с помощью torch.distributed.monitored_barrier() и TORCH_DISTRIBUTED_DEBUG, базовая библиотека C++ для torch.distributed также выводит сообщения журнала с различными уровнями детализации. Эти сообщения помогают понять состояние выполнения распределённой задачи обучения и устранить такие проблемы, как сбои сетевого подключения. В следующей таблице показано, как изменить уровень журнала с помощью сочетания переменных среды TORCH_CPP_LOG_LEVEL и TORCH_DISTRIBUTED_DEBUG.

TORCH_CPP_LOG_LEVEL

TORCH_DISTRIBUTED_DEBUG

Эффективный уровень журнала

ERROR

игнорируется

Ошибка

WARNING

игнорируется

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

INFO

игнорируется

Информация

INFO

INFO

Отладка

INFO

DETAIL

Трассировка (то же, что и All)

Распределённые компоненты создают собственные типы исключений, производные от RuntimeError:

  • torch.distributed.DistError: Базовый тип всех распределённых исключений.
  • torch.distributed.DistBackendError: Это исключение возникает при ошибке, специфичной для бэкенда. Например, если используется бэкенд NCCL и пользователь пытается использовать графический процессор, недоступный библиотеке NCCL.
  • torch.distributed.DistNetworkError: Это исключение возникает при ошибках сетевых библиотек (например, Connection reset by peer).
  • torch.distributed.DistStoreError: Это исключение возникает при ошибке хранилища (например, истечении времени ожидания TCPStore).
class torch.distributed.DistError

Исключение, возникающее при ошибке в распределённой библиотеке

class torch.distributed.DistBackendError

Исключение, возникающее при ошибке бэкенда в распределённой среде

class torch.distributed.DistNetworkError

Исключение, возникающее при сетевой ошибке в распределённой среде

class torch.distributed.DistStoreError

Исключение, возникающее при ошибке распределённого хранилища

При обучении на одном узле может быть удобно устанавливать точку останова в скрипте в интерактивном режиме. Мы предлагаем способ удобно установить точку останова на одном ранге:

torch.distributed.breakpoint(rank=0, skip=0, timeout_s=3600) [исходный код]

Устанавливает точку останова только на одном ранге. Все остальные ранги будут ждать, пока вы не завершите отладку, и только затем продолжат работу.

Параметры:
  • rank (int) – На каком ранге установить точку останова. По умолчанию: 0
  • skip (int) – Пропустить первые skip вызовы этой точки останова. По умолчанию: 0.
torch.distributed.collective_utils.all_gather_object_enforce_type(pg, object_list, obj, type_checker=<function <lambda>>) [исходный код]

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

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

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

torch.distributed.launcher.api.launch_agent(config, entrypoint, args, health_check_server=None) [исходный код]
Тип возвращаемого значения:

dict[int, Any]

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

Spec-Zone.ru

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