Пакет распределённой коммуникации - torch.distributed
Примечание
Обратитесь к Обзору распределённых вычислений PyTorch для краткого введения во все возможности, связанные с распределённым обучением.
Бэкэнды
torch.distributed поддерживает три встроенных бэкэнда, каждый из которых обладает различными возможностями. В таблице ниже показано, какие функции доступны для использования с тензорами CPU/CUDA. MPI поддерживает CUDA только если реализация, используемая для построения PyTorch, её поддерживает.
Бэкэнд |
|
|
| |||
|---|---|---|---|---|---|---|
Устройство | 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.- Возвращаемый тип
-
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 основных способа инициализации группы процессов:
-
- Явно указать
store,rank, иworld_size. - Указать
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 (str или Backend, необязательно) – Используемый бэкенд. В зависимости от конфигурации во время сборки, допустимые значения включают
Примечание
Чтобы включить
backend == Backend.MPI, PyTorch необходимо собрать из исходного кода на системе, поддерживающей MPI.Примечание
Поддержка нескольких бэкендов находится в стадии разработки. В настоящее время, когда бэкенд не указан, создаются оба бэкенда
glooиnccl. Бэкендglooбудет использоваться для коллективов с тензорами CPU, а бэкендncclбудет использоваться для коллективов с тензорами CUDA. Пользовательский бэкенд может быть указан путём передачи строки в формате «<тип_устройства>:<имя_бэкенда>,<тип_устройства>:<имя_бэкенда>», например «cpu:gloo,cuda:custom_backend».
-
torch.distributed.is_initialized()[source] -
Проверка инициализации стандартной группы процессов.
- Тип возвращаемого значения
-
torch.distributed.is_mpi_available()[source] -
Проверка доступности бэкенда MPI.
- Тип возвращаемого значения
-
torch.distributed.is_nccl_available()[source] -
Проверка доступности бэкенда NCCL.
- Тип возвращаемого значения
-
torch.distributed.is_gloo_available()[source] -
Проверка доступности бэкенда Gloo.
- Тип возвращаемого значения
-
torch.distributed.is_torchelastic_launched()[source] -
Проверка, был ли данный процесс запущен с
torch.distributed.elastic(иначе torchelastic). Существование переменной окруженияTORCHELASTIC_RUN_IDиспользуется как признак для определения того, был ли текущий процесс запущен с torchelastic. Это разумный признак, посколькуTORCHELASTIC_RUN_IDсопоставляется с идентификатором встречи, который всегда имеет значение, отличное от нуля, указывающее идентификатор задания для целей обнаружения узлов.- Тип возвращаемого значения
В настоящее время поддерживаются три метода инициализации:
Инициализация по 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”
-
name (str) – Имя бэкэнда
Примечание
Поддержка сторонних бэкэндов экспериментальна и может быть изменена.
-
-
torch.distributed.get_backend(group=None)[source] -
Возвращает бэкэнд заданной группы процессов.
- Parameters
-
group (ProcessGroup, optional) – Группа процессов для работы. По умолчанию используется основная главная группа процессов. Если указана другая группа, вызывающий процесс должен быть частью
group. - Returns
-
Бэкэнд данной группы процессов в виде строчной строки.
- Return type
-
torch.distributed.get_rank(group=None)[source] -
Возвращает ранг текущего процесса в заданной
groupили в группе по умолчанию, если она не была указана.Ранг — уникальный идентификатор, присваиваемый каждому процессу в распределённой группе процессов. Они всегда являются последовательными целыми числами от 0 до
world_size.- Parameters
-
group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использована группа процессов по умолчанию.
- Returns
-
Ранг группы процессов -1, если процесс не входит в группу
- Return type
-
torch.distributed.get_world_size(group=None)[source] -
Возвращает количество процессов в текущей группе процессов
- Parameters
-
group (ProcessGroup, optional) – Группа процессов для работы. Если None, будет использована группа процессов по умолчанию.
- Returns
-
Размер группы процессов -1, если процесс не входит в группу
- Return type
Распределённый хранилище ключей-значений
Распределённый пакет поставляется с распределённым хранилищем ключей-значений, которое можно использовать для обмена информацией между процессами в группе, а также для инициализации распределённого пакета в 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 -
Реализация хранилища, которая использует файл для хранения пар ключ-значение.
- Параметры
- Пример::
-
>>> 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.- Параметры
- Пример::
-
>>> 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(), приведёт к исключению.- Параметры
- Пример::
-
>>> 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является пустой строкой.- Параметры
- Пример::
-
>>> 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")
-
torch.distributed.Store.wait(*args, **kwargs) -
Перегруженный метод.
- 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"])
- 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 и не присоединяются к барьеру.
-
ranks (список[целое]) – Список рангов членов группы. Если
- Возвращает
-
Хэндл распределённой группы, который можно передать в вызовы коллективов, или 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 - Тип возвращаемого значения
Примечание: вызов этой функции для группы по умолчанию возвращает идентичность
-
torch.distributed.get_global_rank(group, group_rank)[source] -
Перевод ранга группы в глобальный ранг.
group_rankдолжен быть частьюgroup, иначе это вызывает RuntimeError.- Параметры
-
- group (ProcessGroup) – ProcessGroup для поиска глобального ранга.
- group_rank (int) – Ранг группы для запроса.
- Возвращает
-
Глобальный ранг
group_rankотносительноgroup - Тип возвращаемого значения
Примечание: вызов этой функции для группы по умолчанию возвращает идентичность
-
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, если не входит в группу
- Тип возвращаемого значения
isend() и irecv() возвращают объекты запроса распределённых операций при использовании. В общем случае тип этого объекта не определён, так как их нельзя создавать вручную, но они гарантированно поддерживают два метода:
-
is_completed()- возвращает True, если операция завершена -
wait()- заблокирует процесс до завершения операции.is_completed()гарантированно вернёт True после возвращения.
-
torch.distributed.isend(tensor, dst, group=None, tag=0)[source] -
Асинхронная отправка тензора.
Предупреждение
Изменение
tensorдо завершения запроса приводит к неопределённому поведению.Предупреждение
tagне поддерживается с бэкэндом NCCL.- Параметры
- Возвращает
-
Объект запроса распределённой операции. 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) – Тег для сопоставления отправки и получения.
-
op (Callable) – Функция для отправки данных или получения данных от процесса партнёра. Тип
Синхронные и асинхронные коллективные операции
Каждая функция коллективной операции поддерживает два вида операций в зависимости от значения флага 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) – Является ли эта операция асинхронной
-
tensor (Тензор) – Данные для отправки, если
- Возвращаемое значение
-
Обработчик асинхронной работы, если 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.
-
object_list (List[Any]) – Список входных объектов для трансляции. Каждый объект должен быть сериализуемым. Только объекты на ранге
- Возвращаемое значение
-
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_gatherAPI, входные тензоры в этом API должны иметь одинаковый размер на всех рангах. - group (ProcessGroup, необязательно) – Процесс-группа для работы. Если None, используется группа по умолчанию.
- async_op (bool, необязательно) – Будет ли эта операция асинхронной
-
output_tensor (Тензор) – Выходной тензор для размещения элементов тензоров со всех рангов. Он должен быть правильно размечен, чтобы иметь одну из следующих форм: (i) конкатенацию всех входных тензоров по первичному измерению; для определения «конкатенации» см.
- Возвращает
-
Дескриптор асинхронной работы, если 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().
- Параметры
- Возвращает
-
Дескриптор асинхронной работы, если 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=Truemonitored_barrierбудет собирать все ошибочные ранги и выбросить ошибку, содержащую информацию обо всех ошибочных рангах.
-
group (ProcessGroup, необязательно) – Группа процессов для работы. Если
- Возвращает
-
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
-
tensor_list (Список[Tensor]) – Тензоры, участвующие в коллективной операции. Если
- Возвращает
-
Дескриптор асинхронной работы, если 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, необязательно) – Является ли эта операция асинхронной
-
tensor_list (List[Tensor]) – Список входных и выходных тензоров коллектива. Функция работает на месте и требует, чтобы каждый тензор был тензором графического процессора на разных графических процессорах. Вы также должны убедиться, что
- Возвращает
-
Дескриптор асинхронной работы, если 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
-
tensor_list (List[Tensor]) – Входные и выходные тензоры графического процессора коллектива. Функция работает на месте. Вы также должны убедиться, что
- Возвращает
-
Дескриптор асинхронной работы, если 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).
Как использовать этот модуль:
- Распределённое обучение на одном узле с несколькими процессами
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: (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)
- Чтобы узнать, какие дополнительные аргументы предлагает этот модуль:
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.
|
| Эффективный уровень журнала |
|---|---|---|
| игнорируется | Ошибка |
| игнорируется | Предупреждение |
| игнорируется | Информация |
|
| Отладка |
|
| Отслеживание (также известное как Все) |
В распределенном модуле есть собственный тип исключения, производный от 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