Spec-Zone.ru › PyTorch 2

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

Примечание

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

Бэкэнды

torch.distributed поддерживает три встроенных бэкэнда, каждый из которых обладает различными возможностями. В таблице ниже показано, какие функции доступны для использования с тензорами CPU/CUDA. MPI поддерживает CUDA только если реализация, используемая для построения PyTorch, её поддерживает.

Бэкэнд

gloo

mpi

nccl

Устройство

CPU

GPU

CPU

GPU

CPU

GPU

send

✓

✘

✓

?

✘

✓

recv

✓

✘

✓

?

✘

✓

broadcast

✓

✓

✓

?

✘

✓

all_reduce

✓

✓

✓

?

✘

✓

reduce

✓

✘

✓

?

✘

✓

all_gather

✓

✘

✓

?

✘

✓

gather

✓

✘

✓

?

✘

✓

scatter

✓

✘

✓

?

✘

✓

reduce_scatter

✘

✘

✘

✘

✘

✓

all_to_all

✘

✘

✓

?

✘

✓

barrier

✓

✘

✓

?

✘

✓

Бэкэнды, поставляемые с PyTorch

Пакет PyTorch distributed поддерживает Linux (стабильная версия), MacOS (стабильная версия) и Windows (прототип). По умолчанию для Linux строятся и включаются бэкэнды Gloo и NCCL (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.

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

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

  • Правило большого пальца

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

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

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

    • Если ваш InfiniBand имеет включённый IP через 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

Основы

Пакет torch.distributed предоставляет поддержку PyTorch и примитивы коммуникации для многопроцессорной параллельности на нескольких вычислительных узлах, работающих на одном или нескольких компьютерах. Класс torch.nn.parallel.DistributedDataParallel() использует эту функциональность для обеспечения синхронного распределённого обучения в качестве обертки вокруг любой модели PyTorch. Это отличается от видов параллельности, предоставляемых Пакет многопроцессорности - 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.is_available() [source]

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

Возвращаемый тип

bool

END_OF_DOCUMENT_MARKER
torch.distributed.init_process_group(backend=None, init_method=None, timeout=datetime.timedelta(seconds=1800), world_size=-1, rank=-1, store=None, group_name='', pg_options=None) [source]

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

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

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

Параметры
  • backend (str или Backend, необязательно) – Используемый бэкенд. В зависимости от конфигурации во время сборки, допустимые значения включают mpi, gloo, nccl, и ucc. Если бэкенд не указан, то будут созданы бэкенды gloo и nccl, см. примечания ниже о том, как обрабатываются несколько бэкендов. Это поле может быть задано как строка в нижнем регистре (например, "gloo"), к которой также можно получить доступ через атрибуты Backend (например, Backend.GLOO). При использовании нескольких процессов на одной машине с бэкендом nccl каждый процесс должен иметь эксклюзивный доступ ко всем используемым им GPU, так как совместное использование GPU между процессами может привести к тупиковым ситуациям. Бэкенд ucc находится в стадии разработки.
  • init_method (str, необязательно) – URL, определяющий, как инициализировать группу процессов. По умолчанию “env://”, если не указан init_method или store . Взаимоисключающее с store.
  • world_size (int, необязательно) – Количество процессов, участвующих в работе. Требуется, если указан store.
  • rank (int, необязательно) – Ранг текущего процесса (он должен быть числом от 0 до world_size-1). Требуется, если указан store.
  • store (Store, необязательно) – Хранилище ключей/значений, доступное всем рабочим узлам, используемое для обмена информацией о соединении/адресе. Взаимоисключающее с init_method.
  • timeout (timedelta, необязательно) – Таймаут для операций, выполненных с группой процессов. Значение по умолчанию равно 30 минутам. Это применимо для бэкенда gloo. Для nccl, это применимо только в том случае, если переменная среды NCCL_BLOCKING_WAIT или NCCL_ASYNC_ERROR_HANDLING установлена в 1. Когда NCCL_BLOCKING_WAIT установлена, это время, в течение которого процесс будет заблокирован и ждать завершения коллективов, прежде чем выбросить исключение. Когда NCCL_ASYNC_ERROR_HANDLING установлена, это время, после которого коллективы будут прерваны асинхронно, и процесс завершится аварийно. NCCL_BLOCKING_WAIT предоставит пользователю ошибки, которые могут быть пойманы и обработаны, но из-за своей блокирующей природы, это имеет издержки производительности. С другой стороны, NCCL_ASYNC_ERROR_HANDLING имеет очень низкие издержки производительности, но аварийно завершает процесс при ошибках. Это делается потому, что выполнение CUDA асинхронно, и больше нельзя безопасно продолжать выполнение пользовательского кода, поскольку необработанные асинхронные операции NCCL могут привести к тому, что последующие операции CUDA будут выполняться с поврежденными данными. Должна быть установлена только одна из этих двух переменных среды. Для ucc, блокирующее ожидание поддерживается аналогично NCCL. Однако обработка асинхронных ошибок выполняется по-другому, так как у UCC есть поток выполнения, а не поток-сторож.
  • group_name (str, необязательно, устарело) – Имя группы. Этот аргумент игнорируется.
  • pg_options (ProcessGroupOptions, необязательно) – опции группы процессов, определяющие какие дополнительные опции должны быть переданы во время построения конкретных групп процессов. На данный момент мы поддерживаем только ProcessGroupNCCL.Options для бэкенда nccl, is_high_priority_stream может быть указано, чтобы бэкенд nccl мог выбрать потоки CUDA с высоким приоритетом, когда есть ожидающие вычисления.

Примечание

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

Примечание

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

torch.distributed.is_initialized() [source]

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

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

bool

torch.distributed.is_mpi_available() [source]

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

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

bool

torch.distributed.is_nccl_available() [source]

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

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

bool

torch.distributed.is_gloo_available() [source]

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

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

bool

torch.distributed.is_torchelastic_launched() [source]

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

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

bool

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

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

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

Обратите внимание, что адрес мультикаста больше не поддерживается в последней версии пакета распределенных вычислений. 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() с тем же путём/именем файла.

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

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

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

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

Этот метод всегда создает файл и делает все возможное, чтобы очистить и удалить файл в конце программы. Другими словами, каждое инициализирование с помощью метода 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.distributed.init_process_group() можно использовать следующие функции. Чтобы проверить, была ли группа процессов уже инициализирована, используйте torch.distributed.is_initialized().

class torch.distributed.Backend(name) [source]

Класс-перечисление доступных бэкэндов: GLOO, NCCL, UCC, MPI и другие зарегистрированные бэкэнды.

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

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

Примечание

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

classmethod register_backend(name, func, extended_api=False, devices=None) [source]

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

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

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

Примечание

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

torch.distributed.get_backend(group=None) [source]

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

Parameters

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

Returns

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

Return type

str

torch.distributed.get_rank(group=None) [source]

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

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

Parameters

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

Returns

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

Return type

int

torch.distributed.get_world_size(group=None) [source]

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

Parameters

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

Returns

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

Return type

int

Распределённый хранилище ключей-значений

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

class torch.distributed.Store

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

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_worker (bool, опционально) – Ожидать ли подключения всех рабочих узлов к серверному хранилищу. Применимо только в случае, если world_size имеет фиксированное значение. По умолчанию True.
  • multi_tenant (bool, опционально) – Если True, все экземпляры TCPStore в текущем процессе с одинаковым host/port будут использовать одно и то же базовое хранилище TCPServer. По умолчанию False.
  • master_listen_fd (int, опционально) – Если указано, базовое хранилище TCPServer будет слушать на этом дескрипторе файла, который должен быть уже связанным сокетом к port. Полезно для избежания гонок при назначении портов в некоторых сценариях. По умолчанию None (что означает, что сервер создаёт новый сокет и пытается связать его с port).
Пример::
>>> 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")
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")
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")
class torch.distributed.PrefixStore

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

Параметры
  • prefix (str) – Строка-префикс, добавляемая к каждому ключу перед его вставкой в хранилище.
  • store (torch.distributed.store) – Объект хранилища, который является базовым хранилищем ключей-значений.
torch.distributed.Store.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")
torch.distributed.Store.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")
torch.distributed.Store.add(self: torch._C._distributed_c10d.Store, arg0: str, arg1: int) → 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")
torch.distributed.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")
END_OF_DOCUMENT_MARKER
torch.distributed.Store.wait(*args, **kwargs)

Перегруженный метод.

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

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

Параметры

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))
>>> # This will throw an exception after 30 seconds
>>> store.wait(["bad_key"])
  1. wait(self: torch._C._distributed_c10d.Store, arg0: List[str], arg1: datetime.timedelta) -> None

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

Параметры
  • keys (список) – Список ключей, по которым ожидается установка в хранилище.
  • 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))
torch.distributed.Store.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()
torch.distributed.Store.delete_key(self: torch._C._distributed_c10d.Store, arg0: str) → bool

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

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

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

Параметры

key (строка) – Ключ, подлежащий удалению из хранилища

Возвращает

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")
torch.distributed.Store.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"])

Группы

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

torch.distributed.new_group(ranks=None, timeout=datetime.timedelta(seconds=1800), backend=None, pg_options=None, use_local_synchronization=False) [source]

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

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

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

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

Параметры
  • ranks (список[целое]) – Список рангов членов группы. Если None, будет установлено всем рангам. Значение по умолчанию — None.
  • timeout (timedelta, необязательно) – Таймаут для операций, выполняемых над группой процессов. Значение по умолчанию равно 30 минутам. Это применимо для бекенда gloo. Для nccl, это применимо только если переменная окружения NCCL_BLOCKING_WAIT или NCCL_ASYNC_ERROR_HANDLING установлена в 1. Когда NCCL_BLOCKING_WAIT установлена, это время, в течение которого процесс будет блокироваться и ждать завершения коллективов перед выбросом исключения. Когда NCCL_ASYNC_ERROR_HANDLING установлена, это время, после которого коллективы будут асинхронно прерваны, и процесс аварийно завершит работу. NCCL_BLOCKING_WAIT предоставит пользователю ошибки, которые можно перехватить и обработать, но из-за своей блокирующей природы, она имеет накладные расходы на производительность. С другой стороны, NCCL_ASYNC_ERROR_HANDLING имеет очень низкие накладные расходы на производительность, но аварийно завершает процесс при ошибках. Это делается, так как выполнение CUDA асинхронно, и теперь небезопасно продолжать выполнение пользовательского кода, так как невыполненные асинхронные операции NCCL могут привести к запуску последующих операций CUDA на поврежденных данных. Только одна из этих двух переменных окружения должна быть установлена.
  • backend (строка или Backend, необязательно) – Бэкенд для использования. В зависимости от конфигурации во время сборки, допустимые значения — gloo и nccl. По умолчанию использует тот же бэкенд, что и глобальная группа. Этот параметр должен быть передан в виде строчной строки (например, "gloo"), к которому также можно получить доступ через атрибуты Backend (например, Backend.GLOO). Если None передаётся, будет использован бэкенд, соответствующий группе процессов по умолчанию. Значение по умолчанию — None.
  • pg_options (ProcessGroupOptions, необязательно) – параметры группы процессов, указывающие какие дополнительные параметры необходимо передать при построении конкретных групп процессов. т.е. для бэкенда nccl, is_high_priority_stream может быть указан, чтобы группа процессов могла выбрать потоки CUDA с высоким приоритетом.
  • use_local_synchronization (булево, необязательно) – выполнить барьер группы на локальном уровне в конце создания группы процессов. Это отличается тем, что процессы, не являющиеся членами, не вызывают API и не присоединяются к барьеру.
Возвращает

Хэндл распределённой группы, который можно передать в вызовы коллективов, или None, если ранг не входит в состав ranks.

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

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

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

torch.distributed.get_group_rank(group, global_rank) [source]

Перевод глобального ранга в ранг группы.

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

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

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

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

int

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

torch.distributed.get_global_rank(group, group_rank) [source]

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

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

Параметры
  • group (ProcessGroup) – ProcessGroup для поиска глобального ранга.
  • group_rank (int) – Ранг группы для запроса.
Возвращает

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

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

int

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

torch.distributed.get_process_group_ranks(group) [source]

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

Параметры

group (ProcessGroup) – ProcessGroup для получения всех рангов.

Возвращает

Список глобальных рангов, отсортированных по рангу группы.

Точечная коммуникация

torch.distributed.send(tensor, dst, group=None, tag=0) [source]

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

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

Синхронный приём тензора.

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

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

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

int

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

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

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

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

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

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

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

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

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

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

Work

torch.distributed.irecv(tensor, src=None, group=None, tag=0) [source]

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

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

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

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

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

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

Work

torch.distributed.batch_isend_irecv(p2p_op_list) [source]

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

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

Параметры

p2p_op_list – Список операций точка-точка (тип каждого оператора — torch.distributed.P2POp). Порядок isend/irecv в списке имеет значение и должен соответствовать соответствующим isend/irecv на удалённом конце.

Возвращаемое значение

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

Примеры

>>> send_tensor = torch.arange(2) + 2 * rank
>>> recv_tensor = torch.randn(2)
>>> 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 пользователи должны установить текущий графический процессор с помощью torch.cuda.set_device, в противном случае это приведёт к неожиданным проблемам зависания.

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

class torch.distributed.P2POp(op, tensor, peer, group=None, tag=0) [source]

Класс для построения операций точка-точка для batch_isend_irecv.

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

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

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

Каждая функция коллективной операции поддерживает два вида операций в зависимости от значения флага 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-потоковой очереди и результат может быть использован в потоковой очереди по умолчанию без дополнительной синхронизации.
  • 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, group=None, async_op=False) [source]

Транслирует тензор во всю группу.

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

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

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

torch.distributed.broadcast_object_list(object_list, src=0, group=None, device=None) [source]

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

Параметры
  • object_list (List[Any]) – Список входных объектов для трансляции. Каждый объект должен быть сериализуемым. Только объекты на ранге src будут транслироваться, но каждый ранг должен предоставлять списки одинаковых размеров.
  • src (int) – Источник-ранг для трансляции object_list.
  • group – (ProcessGroup, optional): Группа процессов для работы. Если None, используется группа по умолчанию. Значение по умолчанию — None.
  • device (torch.device, optional) – Если не None, объекты сериализуются и преобразуются в тензоры, которые перемещаются на device перед трансляцией. Значение по умолчанию — None.
Возвращаемое значение

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

Примечание

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

Примечание

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

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

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

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

Вызов broadcast_object_list() с тензорами GPU не хорошо поддерживается и неэффективен, так как происходит передача данных между 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, op=<RedOpType.SUM: 0>, group=None, async_op=False) [source]

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

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

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

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

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

Примеры

>>> # All tensors below are of torch.int64 type.
>>> # We have 2 process groups, 2 ranks.
>>> tensor = torch.arange(2, dtype=torch.int64) + 1 + 2 * rank
>>> tensor
tensor([1, 2]) # Rank 0
tensor([3, 4]) # Rank 1
>>> dist.all_reduce(tensor, op=ReduceOp.SUM)
>>> tensor
tensor([4, 6]) # Rank 0
tensor([4, 6]) # 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) + 2 * rank * (1+1j)
>>> tensor
tensor([1.+1.j, 2.+2.j]) # Rank 0
tensor([3.+3.j, 4.+4.j]) # Rank 1
>>> dist.all_reduce(tensor, op=ReduceOp.SUM)
>>> tensor
tensor([4.+4.j, 6.+6.j]) # Rank 0
tensor([4.+4.j, 6.+6.j]) # Rank 1
torch.distributed.reduce(tensor, dst, op=<RedOpType.SUM: 0>, group=None, async_op=False) [source]

Редуцирует данные тензора по всем машинам.

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

Параметры
  • tensor (Тензор) – Вход и выход коллектива. Функция работает на месте.
  • dst (int) – Целевой ранг
  • op (необязательно) – Одно из значений из torch.distributed.ReduceOp перечисления. Указывает операцию, используемую для поэлементного сокращения.
  • group (ProcessGroup, необязательно) – Процесс-группа для работы. Если None, используется группа по умолчанию.
  • async_op (bool, необязательно) – Будет ли эта операция асинхронной
Возвращает

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

torch.distributed.all_gather(tensor_list, tensor, group=None, async_op=False) [source]

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

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

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

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

Примеры

>>> # All tensors below are of torch.int64 dtype.
>>> # We have 2 process groups, 2 ranks.
>>> tensor_list = [torch.zeros(2, dtype=torch.int64) for _ in range(2)]
>>> tensor_list
[tensor([0, 0]), tensor([0, 0])] # Rank 0 and 1
>>> tensor = torch.arange(2, dtype=torch.int64) + 1 + 2 * rank
>>> tensor
tensor([1, 2]) # Rank 0
tensor([3, 4]) # Rank 1
>>> dist.all_gather(tensor_list, tensor)
>>> tensor_list
[tensor([1, 2]), tensor([3, 4])] # Rank 0
[tensor([1, 2]), tensor([3, 4])] # 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) for _ in range(2)]
>>> tensor_list
[tensor([0.+0.j, 0.+0.j]), tensor([0.+0.j, 0.+0.j])] # Rank 0 and 1
>>> tensor = torch.tensor([1+1j, 2+2j], dtype=torch.cfloat) + 2 * rank * (1+1j)
>>> tensor
tensor([1.+1.j, 2.+2.j]) # Rank 0
tensor([3.+3.j, 4.+4.j]) # Rank 1
>>> dist.all_gather(tensor_list, tensor)
>>> tensor_list
[tensor([1.+1.j, 2.+2.j]), tensor([3.+3.j, 4.+4.j])] # Rank 0
[tensor([1.+1.j, 2.+2.j]), tensor([3.+3.j, 4.+4.j])] # Rank 1
torch.distributed.all_gather_into_tensor(output_tensor, input_tensor, group=None, async_op=False) [source]

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

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

Дескриптор асинхронной работы, если 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_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_into_tensor(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_into_tensor(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

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

Бэкенд Gloo не поддерживает этот API.

torch.distributed.all_gather_object(object_list, obj, group=None) [source]

Собирает пиклируемые объекты из всей группы в список. Аналогично all_gather(), но передаются Python-объекты. Обратите внимание, что объект должен быть пиклируемым для сбора.

Параметры
  • object_list (список[любой]) – Список вывода. Он должен быть правильно размечен как размер группы для этого коллектива и будет содержать вывод.
  • obj (любой) – Пиклируемый Python-объект, передаваемый из текущего процесса.
  • group (ProcessGroup, необязательно) – Процесс-группа для работы. Если None, используется группа по умолчанию. По умолчанию None.
Возвращает

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

Примечание

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

Примечание

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

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

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

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

Вызов all_gather_object() с тензорами GPU не хорошо поддерживается и неэффективен, поскольку он влечёт за собой передачу 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.gather(tensor, gather_list=None, dst=0, group=None, async_op=False) [source]

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

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

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

torch.distributed.gather_object(obj, object_gather_list=None, dst=0, group=None) [source]

Сборка пикелируемых объектов со всей группы в одном процессе. Аналогично gather(), но передаются объекты Python. Обратите внимание, что объект должен быть пикелируемым, чтобы быть собранным.

Параметры
  • obj (Any) – Входной объект. Должен быть пикелируемым.
  • object_gather_list (список[Any]) – Выходной список. На ранге dst, он должен быть правильно размечен как размер группы для этого коллектива и будет содержать вывод. Должен быть None на рангах, не являющихся конечными. (по умолчанию None)
  • dst (int, необязательно) – Конечный ранг. (по умолчанию 0)
  • group – (ProcessGroup, необязательно): Группа процессов для работы. Если None, будет использована группа процессов по умолчанию. Значение по умолчанию None.
Возвращает

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

Примечание

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

Примечание

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

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

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

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

Вызов gather_object() с GPU-тензорами не очень хорошо поддерживается и является неэффективным, так как это приводит к передаче данных с GPU на CPU, так как тензоры будут сериализованы. Пожалуйста, рассмотрите использование 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=0, group=None, async_op=False) [source]

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

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

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

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

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

Примечание

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

Пример::
>>> # Note: Process group initialization omitted on each rank.
>>> import torch.distributed as dist
>>> tensor_size = 2
>>> t_ones = torch.ones(tensor_size)
>>> t_fives = torch.ones(tensor_size) * 5
>>> output_tensor = torch.zeros(tensor_size)
>>> if dist.get_rank() == 0:
>>>     # Assumes world_size of 2.
>>>     # Only tensors, all of which must be the same size.
>>>     scatter_list = [t_ones, t_fives]
>>> else:
>>>     scatter_list = None
>>> dist.scatter(output_tensor, scatter_list, src=0)
>>> # Rank i gets scatter_list[i]. For example, on rank 1:
>>> output_tensor
tensor([5., 5.])
torch.distributed.scatter_object_list(scatter_object_output_list, scatter_object_input_list, src=0, group=None) [source]

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

Параметры
  • scatter_object_output_list (Список[Any]) – Непустой список, первый элемент которого будет хранить объект, распространённый на этот ранг.
  • scatter_object_input_list (Список[Any]) – Список входных объектов для распределения. Каждый объект должен быть пикелируемым. Только объекты на ранге src будут распределены, и аргумент может быть None для рангов, не являющихся исходными.
  • src (int) – Исходный ранг для распределения scatter_object_input_list.
  • group – (ProcessGroup, необязательно): Группа процессов для работы. Если None, будет использована группа процессов по умолчанию. Значение по умолчанию None.
Возвращает

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

Примечание

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

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

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

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

Вызов scatter_object_list() с GPU-тензорами не очень хорошо поддерживается и является неэффективным, так как это приводит к передаче данных с GPU на CPU, так как тензоры будут сериализованы. Пожалуйста, рассмотрите использование 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, input_list, op=<RedOpType.SUM: 0>, group=None, async_op=False) [source]

Производит уменьшение, затем рассеивание списка тензоров по всем процессам в группе.

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

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

torch.distributed.reduce_scatter_tensor(output, input, op=<RedOpType.SUM: 0>, group=None, async_op=False) [source]

Производит уменьшение, затем рассеивание тензора по всем рангам в группе.

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

Дескриптор асинхронной работы, если 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_tensor(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_tensor(tensor_out, tensor_in)
>>> tensor_out
tensor([0, 2], device='cuda:0') # Rank 0
tensor([4, 6], device='cuda:1') # Rank 1

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

Бэкенд Gloo не поддерживает этот API.

torch.distributed.all_to_all_single(output, input, output_split_sizes=None, input_split_sizes=None, group=None, async_op=False) [source]

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

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

Параметры
  • output (Тензор) – Собраный конкатенированный выходной тензор.
  • input (Тензор) – Входной тензор для рассеивания.
  • output_split_sizes – (список[Целое число], необязательно): Размеры разделения выходного тензора по размеру 0, если указано None или пусто, размер 0 тензора output должен делиться равномерно на world_size.
  • input_split_sizes – (список[Целое число], необязательно): Размеры разделения входного тензора по размеру 0, если указано None или пусто, размер 0 тензора input должен делиться равномерно на world_size.
  • group (ProcessGroup, необязательно) – Группа процессов для работы. Если None, используется группа по умолчанию.
  • async_op (bool, необязательно) – Является ли операция асинхронной.
Возвращает

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

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

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) [source]

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

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

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

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

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

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=None, async_op=False, device_ids=None) [source]

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

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

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

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

torch.distributed.monitored_barrier(group=None, timeout=None, wait_all_ranks=False) [source]

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

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

Примечание

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

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

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.ReduceOp

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

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

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

PREMUL_SUM умножает входные данные на заданный скаляр локально перед редукцией. PREMUL_SUM доступен только с бэкендом NCCL и только для версий NCCL 2.11 и выше. Пользователи должны использовать torch.distributed._make_nccl_premul_sum.

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

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

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

class torch.distributed.reduce_op

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

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

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

Обратите внимание, что вы можете использовать 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)

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

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

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

Функции для нескольких GPU будут устаревать. Если вам необходимо их использовать, пожалуйста, обратитесь к нашей документации позже.

Если у вас есть более одного GPU на каждом узле, при использовании бэкэндов NCCL и Gloo, broadcast_multigpu() all_reduce_multigpu() reduce_multigpu() all_gather_multigpu() и reduce_scatter_multigpu() поддерживают распределенные коллективные операции между несколькими GPU на каждом узле. Эти функции могут потенциально улучшить общую производительность распределенного обучения и легко используются путем передачи списка тензоров. Каждый тензор в переданном списке тензоров должен находиться на отдельном устройстве GPU хоста, на котором вызывается функция. Обратите внимание, что длина списка тензоров должна быть одинаковой для всех распределенных процессов. Также обратите внимание, что в настоящее время функции коллективных операций для нескольких GPU поддерживаются только бэкендом NCCL.

Например, если система, которую мы используем для распределенного обучения, имеет 2 узла, каждый из которых имеет 8 GPU. На каждом из 16 GPU есть тензор, который мы хотим всеобъединить. Следующий код может служить ссылкой:

Код, выполняющийся на узле 0

import torch
import torch.distributed as dist

dist.init_process_group(backend="nccl",
                        init_method="file:///distributed_test",
                        world_size=2,
                        rank=0)
tensor_list = []
for dev_idx in range(torch.cuda.device_count()):
    tensor_list.append(torch.FloatTensor([1]).cuda(dev_idx))

dist.all_reduce_multigpu(tensor_list)

Код, выполняющийся на узле 1

import torch
import torch.distributed as dist

dist.init_process_group(backend="nccl",
                        init_method="file:///distributed_test",
                        world_size=2,
                        rank=1)
tensor_list = []
for dev_idx in range(torch.cuda.device_count()):
    tensor_list.append(torch.FloatTensor([1]).cuda(dev_idx))

dist.all_reduce_multigpu(tensor_list)

После вызова все 16 тензоров на двух узлах будут иметь значение всехобъединения 16

torch.distributed.broadcast_multigpu(tensor_list, src, group=None, async_op=False, src_tensor=0) [source]

Транслирует тензор во всю группу с тензорами нескольких GPU на каждом узле.

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

В настоящее время поддерживаются только бэкэнды nccl и gloo, тензоры должны быть только тензорами GPU

Параметры
  • tensor_list (Список[Tensor]) – Тензоры, участвующие в коллективной операции. Если src ранг, то указанный src_tensor элемент tensor_list (tensor_list[src_tensor]) будет транслирован во все остальные тензоры (на разных GPU) в исходном процессе и все тензоры в tensor_list других процессов, не являющихся исходными. Также необходимо убедиться, что len(tensor_list) одинаков для всех распределенных процессов, вызывающих эту функцию.
  • src (int) – Истоковый ранг.
  • group (ProcessGroup, необязательно) – Группа процессов для работы. Если None, будет использоваться группа по умолчанию.
  • async_op (bool, необязательно) – Является ли эта операция асинхронной
  • src_tensor (int, необязательно) – Ранг исходного тензора внутри tensor_list
Возвращает

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

torch.distributed.all_reduce_multigpu(tensor_list, op=<RedOpType.SUM: 0>, group=None, async_op=False) [source]

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

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

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

В настоящее время поддерживаются только бэкэнды nccl и gloo; тензоры должны быть только тензорами графического процессора

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

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

torch.distributed.reduce_multigpu(tensor_list, dst, op=<RedOpType.SUM: 0>, group=None, async_op=False, dst_tensor=0) [source]

Снижает данные тензора на нескольких графических процессорах на всех машинах. Каждый тензор в tensor_list должен находиться на отдельном графическом процессоре

Только графический процессор tensor_list[dst_tensor] в процессе с рангом dst получит окончательный результат.

В настоящее время поддерживается только бэкэнд nccl; тензоры должны быть только тензорами графического процессора

Параметры
  • tensor_list (List[Tensor]) – Входные и выходные тензоры графического процессора коллектива. Функция работает на месте. Вы также должны убедиться, что len(tensor_list) одинаковый для всех распределенных процессов, вызывающих эту функцию.
  • dst (int) – Ранг назначения
  • op (необязательно) – Одно из значений из torch.distributed.ReduceOp перечисления. Указывает операцию, используемую для поэлементного снижения.
  • group (ProcessGroup, необязательно) – Группа процессов, на которой нужно работать. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, необязательно) – Является ли эта операция асинхронной
  • dst_tensor (int, необязательно) – Ранг тензора назначения в tensor_list
Возвращает

Дескриптор асинхронной работы, если async_op установлено в True. В противном случае None

torch.distributed.all_gather_multigpu(output_tensor_lists, input_tensor_list, group=None, async_op=False) [source]

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

В настоящее время поддерживается только бэкэнд nccl; тензоры должны быть только тензорами графического процессора

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

Параметры
  • output_tensor_lists (List[List[Tensor]]) –

    Выходные списки. Он должен содержать тензоры правильного размера на каждом графическом процессоре, используемые для вывода коллектива, например, output_tensor_lists[i] содержит результат all_gather, который находится на графическом процессоре input_tensor_list[i].

    Обратите внимание, что каждый элемент output_tensor_lists имеет размер world_size * len(input_tensor_list), так как функция собирает результат с каждого отдельного графического процессора в группе. Чтобы интерпретировать каждый элемент output_tensor_lists[i], обратите внимание, что input_tensor_list[j] ранга k появится в output_tensor_lists[i][k * world_size + j].

    Также обратите внимание, что len(output_tensor_lists), и размер каждого элемента в output_tensor_lists (каждый элемент является списком, следовательно, len(output_tensor_lists[i])) должны быть одинаковыми для всех распределенных процессов, вызывающих эту функцию.

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

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

torch.distributed.reduce_scatter_multigpu(output_tensor_list, input_tensor_lists, op=<RedOpType.SUM: 0>, group=None, async_op=False) [source]

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

Каждый тензор в output_tensor_list должен находиться на отдельном графическом процессоре, как и каждый список тензоров в input_tensor_lists.

Параметры
  • output_tensor_list (List[Tensor]) –

    Выходные тензоры (на разных графических процессорах) для получения результата операции.

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

  • input_tensor_lists (List[List[Tensor]]) –

    Входные списки. Он должен содержать тензоры правильного размера на каждом графическом процессоре, используемые для входных данных коллектива, например, input_tensor_lists[i] содержит вход reduce_scatter, который находится на графическом процессоре output_tensor_list[i].

    Обратите внимание, что каждый элемент input_tensor_lists имеет размер world_size * len(output_tensor_list), так как функция рассеивает результат с каждого отдельного графического процессора в группе. Чтобы интерпретировать каждый элемент input_tensor_lists[i], обратите внимание, что output_tensor_list[j] ранга k получает результат reduce-scatter от input_tensor_lists[i][k * world_size + j].

    Также обратите внимание, что len(input_tensor_lists), и размер каждого элемента в input_tensor_lists (каждый элемент является списком, следовательно, len(input_tensor_lists[i])) должны быть одинаковыми для всех распределенных процессов, вызывающих эту функцию.

  • group (ProcessGroup, необязательно) – Группа процессов, на которой нужно работать. Если None, будет использоваться группа процессов по умолчанию.
  • async_op (bool, необязательно) – Является ли эта операция асинхронной.
Возвращает

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

Бэкэнды сторонних разработчиков

Помимо встроенных бэкэндов GLOO/MPI/NCCL, PyTorch distributed поддерживает бэкэнды сторонних разработчиков с помощью механизма регистрации во время выполнения. Для получения ссылок по разработке бэкэнда сторонних разработчиков с помощью C++ Extension, пожалуйста, обратитесь к Tutorials - Custom C++ and CUDA Extensions и test/cpp_extensions/cpp_c10d_extension.cpp. Возможности бэкэндов сторонних разработчиков определяются их собственными реализациями.

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

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

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

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

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

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

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

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

Этот модуль будет устаревать в пользу torchrun.

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

В обоих случаях, для распределённого обучения на одном узле или на нескольких узлах, эта утилита запустит заданное количество процессов на каждом узле (--nproc-per-node). При использовании для обучения на графическом процессоре это число должно быть меньше или равно количеству графических процессоров в текущей системе (nproc_per_node), и каждый процесс будет работать с одним графическим процессором от 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. Эта утилита и распределённое обучение на нескольких процессах (на одном узле или на нескольких узлах) с использованием графического процессора в настоящее время достигает наилучшей производительности только при использовании распределённого бэкенда NCCL. Таким образом, бэкенд NCCL рекомендуется использовать для обучения на графическом процессоре.

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

Обработка аргумента local_rank

>>> import argparse
>>> parser = argparse.ArgumentParser()
>>> parser.add_argument("--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
>>>    ...

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

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

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

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

Убедитесь, что аргумент device_ids установлен в качестве единственного идентификатора устройства графического процессора, на котором будет работать ваш код. Обычно это локальный ранг процесса. Другими словами, 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 НЕ является глобально уникальным: он уникален только для каждого процесса на компьютере. Поэтому не используйте его для принятия решения о том, например, следует ли записывать в сетевую файловую систему. См. https://github.com/pytorch/pytorch/issues/12042 для примера того, как могут возникнуть проблемы, если вы этого не сделаете правильно.

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

Пакет Пакет многопоточности - torch.multiprocessing также предоставляет функцию spawn в torch.multiprocessing.spawn(). Эта вспомогательная функция может быть использована для запуска нескольких процессов. Она работает путём передачи функции, которую нужно запустить, и запускает N процессов для её выполнения. Это также можно использовать для распределённого обучения с несколькими процессами.

Для ссылок на её использование обратитесь к Пример PyTorch - Реализация ImageNet

Обратите внимание, что эта функция требует Python 3.4 или более поздней версии.

Отладка приложений torch.distributed

Отладка распределённых приложений может быть сложной задачей из-за труднопонимаемых зависаний, сбоев или несогласованного поведения между рангами. torch.distributed предоставляет набор инструментов для помощи в отладке обучающих приложений в самообслуживаемом формате:

Мониторируемая преграда

Начиная с версии 1.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() из-за неиспользуемых параметров в модели. В настоящее время find_unused_parameters=True необходимо передать в torch.nn.parallel.DistributedDataParallel() инициализации, если существуют параметры, которые могут быть неиспользуемыми в прямом проходе, а начиная с версии 1.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 запустит дополнительные проверки согласованности и синхронизации при каждом коллективном вызове, сделанном пользователем, как напрямую, так и косвенно (например, DDP allreduce). Это делается путем создания оберточной группы процессов, которая обволакивает все группы процессов, возвращаемые 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().

Ведение журнала

В дополнение к явной поддержке отладки с помощью 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

Отслеживание (также известное как Все)

В распределенном модуле есть собственный тип исключения, производный от RuntimeError, называемый torch.distributed.DistBackendError. Это исключение выбрасывается, когда возникает ошибка, специфичная для бэкенда. Например, если используется бэкенд NCCL, и пользователь пытается использовать графический процессор, недоступный для библиотеки NCCL.

class torch.distributed.DistBackendError

Исключение, выбрасываемое при возникновении ошибки бэкенда в распределенном модуле

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

Тип исключения DistBackendError — экспериментальная функция, и его поведение может измениться.

© 2024, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://pytorch.org/docs/2.1/distributed.html

Spec-Zone.ru

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