Пакет распределённой коммуникации - 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 поддерживает 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-thrashing», возникающие при управлении несколькими потоками выполнения, репликами модели или 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, 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"), к которой также можно получить доступ через атрибуты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.
-
torch.distributed.is_initialized()[source] -
Проверка, инициализирована ли стандартная группа процессов
- Тип возвращаемого значения:
-
torch.distributed.is_mpi_available()[source] -
Проверка доступности бэкенда MPI.
- Тип возвращаемого значения:
-
torch.distributed.is_nccl_available()[source] -
Проверка доступности бэкенда NCCL.
- Тип возвращаемого значения:
-
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)
Инициализация переменными окружения
Этот метод считывает конфигурацию из переменных окружения, что позволяет полностью настроить способ получения информации. Устанавливаемые переменные:
-
MASTER_PORT— обязательно; должна быть свободный порт на машине с рангом 0 -
MASTER_ADDR— обязательно (кроме ранга 0); адрес узла с рангом 0 -
WORLD_SIZE— обязательно; может быть установлено здесь или в вызове функции init -
RANK— обязательно; может быть установлено здесь или в вызове функции init
Машина с рангом 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)[source] -
Регистрирует новый бэкенд с заданным именем и функцией создания.
Этот метод класса используется сторонними
ProcessGroupрасширениями для регистрации новых бэкендов.- Параметры:
-
-
name (строка) — имя бэкенда стороннего
ProcessGroupрасширения. Оно должно соответствовать имени вinit_process_group(). -
func (функция) — обработчик функции, который создаёт бэкенд. Функция должна быть реализована в расширении бэкенда и принимает четыре аргумента, включая
store,rank,world_size, иtimeout. -
extended_api (bool, необязательно) — указывает, поддерживает ли бэкенд расширенную структуру аргументов. По умолчанию:
False. Если установлено значениеTrue, бэкенд получит экземплярc10d::DistributedBackendOptions, и объект параметров группы процессов, как определено реализацией бэкенда.
-
name (строка) — имя бэкенда стороннего
Примечание
Поддержка сторонних бэкендов экспериментальная и может быть изменена.
-
-
torch.distributed.get_backend(group=None)[source] -
Возвращает бэкенд заданной группы процессов.
- Параметры:
-
group (ProcessGroup, необязательно) — группа процессов для работы. По умолчанию используется основная группа процессов. Если указана другая группа, вызывающий процесс должен быть частью
group. - Возвращает:
-
Бэкенд заданной группы процессов в виде строчной строки.
- Тип возвращаемого значения:
-
torch.distributed.get_rank(group=None)[source] -
Возвращает ранг текущего процесса в заданной
groupили в группе по умолчанию, если не было указано.Ранг — это уникальный идентификатор, присваиваемый каждому процессу в группе распределённых процессов. Они всегда являются последовательными целыми числами от 0 до
world_size.- Параметры:
-
group (ProcessGroup, необязательно) — группа процессов для работы. Если None, будет использована группа процессов по умолчанию.
- Возвращает:
-
Ранг группы процессов -1, если процесс не входит в группу
- Тип возвращаемого значения:
-
torch.distributed.get_world_size(group=None)[source] -
Возвращает количество процессов в текущей группе процессов.
- Параметры:
-
group (ProcessGroup, необязательно) — группа процессов для работы. Если None, будет использована группа процессов по умолчанию.
- Возвращает:
-
Размер мира группы процессов -1, если процесс не входит в группу
- Тип возвращаемого значения:
Распределённый хранилище ключей-значений
Распределённый пакет поставляется с распределённым хранилищем ключей-значений, которое можно использовать для обмена информацией между процессами в группе, а также для инициализации распределённого пакета в 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.
- Пример::
-
>>> 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выбросит исключение.- Параметры:
-
ключи (список) – Список ключей, по которым ожидать, пока они будут установлены в хранилище.
- Пример::
-
>>> 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.- Параметры:
-
- ключи (список) – Список ключей, по которым ожидать, пока они будут установлены в хранилище.
- таймаут (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приведет к исключению.- Параметры:
-
ключ (строка) – Ключ, который нужно удалить из хранилища
- Возвращаемое значение:
-
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().- Параметры:
-
таймаут (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)[source] -
Создает новую распределённую группу.
Эта функция требует, чтобы все процессы в основной группе (т.е. все процессы, которые являются частью распределённой задачи) вошли в эту функцию, даже если они не будут членами группы. Кроме того, группы должны создаваться в одном и том же порядке во всех процессах.
Предупреждение
Использование нескольких групп процессов с бэкендом
NCCLодновременно небезопасно, и пользователь должен выполнить явную синхронизацию в своем приложении, чтобы гарантировать использование только одной группы процессов за раз. Это означает, что коллективы из одной группы процессов должны завершить выполнение на устройстве (а не только в очереди, поскольку выполнение CUDA асинхронно), прежде чем коллективы из другой группы процессов будут вставлены в очередь. См. Использование нескольких NCCL-коммуникаторов одновременно для получения дополнительных сведений.- Параметры:
-
-
ранги (список[целое число]) – Список рангов членов группы. Если
None, будет установлен для всех рангов. По умолчаниюNone. -
таймаут (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 на поврежденных данных. Только одна из этих двух переменных окружения должна быть установлена. -
бэкенд (строка или Бэкенд, необязательно) – Бэкенд для использования. В зависимости от конфигурации во время компиляции допустимые значения
glooиnccl. По умолчанию использует тот же бэкенд, что и глобальная группа. Это поле должно быть задано строкой в нижнем регистре (например,"gloo"), которое также можно получить через атрибутыBackend(например,Backend.GLOO). ЕслиNoneпередано, используется бэкенд, соответствующий группе процессов по умолчанию. По умолчаниюNone. -
pg_options (ProcessGroupOptions, необязательно) – параметры группы процессов, указывающие какие дополнительные параметры необходимо передать во время создания конкретных групп процессов. i.e. для бэкенда
nccl,is_high_priority_streamможно указать, чтобы группа процессов подбирала потоки CUDA высокого приоритета.
-
ранги (список[целое число]) – Список рангов членов группы. Если
- Возвращаемое значение:
-
Дескриптор распределённой группы, который может быть передан в вызовы коллективов.
Точечная коммуникация
-
torch.distributed.send(tensor, dst, group=None, tag=0)[source] -
Синхронно отправляет тензор.
- Параметры:
- Тип возвращаемого значения:
-
Work
-
torch.distributed.recv(tensor, src=None, group=None, tag=0)[source] -
Синхронно получает тензор.
- Параметры:
-
- tensor (Тензор) – Тензор для заполнения полученными данными.
- src (int, опционально) – Ранг отправителя. Получит от любого процесса, если не указано.
- group (ProcessGroup, опционально) – Процесс-группа для работы. Если None, используется группа по умолчанию.
- tag (int, опционально) – Метка для сопоставления приёма с отправкой на удалённом узле.
- Возвращаемое значение:
-
Ранг отправителя -1, если не входит в группу
- Тип возвращаемого значения:
-
Work
isend() и irecv() возвращают объекты распределённых запросов при использовании. В общем случае, тип этого объекта не определён, так как их никогда не нужно создавать вручную, но они гарантированно поддерживают два метода:
-
is_completed()- возвращает True, если операция завершена -
wait()- заблокирует процесс до завершения операции.is_completed()гарантированно вернёт True, как только вернётся.
-
torch.distributed.isend(tensor, dst, group=None, tag=0)[source] -
Асинхронно отправляет тензор.
Предупреждение
Изменение
tensorдо завершения запроса приводит к неопределённому поведению.- Параметры:
- Возвращаемое значение:
-
Объект распределённого запроса. None, если не входит в группу
- Тип возвращаемого значения:
-
Work
-
torch.distributed.irecv(tensor, src=None, group=None, tag=0)[source] -
Асинхронно получает тензор.
- Параметры:
-
- tensor (Тензор) – Тензор для заполнения полученными данными.
- src (int, опционально) – Ранг отправителя. Получит от любого процесса, если не указано.
- group (ProcessGroup, опционально) – Процесс-группа для работы. Если None, используется группа по умолчанию.
- tag (int, опционально) – Метка для сопоставления приёма с отправкой на удалённом узле.
- Возвращаемое значение:
-
Объект распределённого запроса. None, если не входит в группу
- Тип возвращаемого значения:
-
Work
Синхронные и асинхронные коллективные операции
Каждая коллективная операция поддерживает два типа операций, в зависимости от флага 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, за исключением операций peer-to-peer. Примечание: поскольку мы продолжаем принимать 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, необязательно) – Группа процессов для работы. Если None, используется группа по умолчанию.
- async_op (bool, необязательно) – Является ли эта операция асинхронной.
-
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должны быть сериализуемыми (picklable) для трансляции.- Параметры:
-
-
object_list (List[Any]) – Список входных объектов для трансляции. Каждый объект должен быть сериализуем (picklable). Только объекты на
srcранге будут транслироваться, но каждый ранг должен предоставить списки одинакового размера. -
src (int) – Ранг источника, из которого транслировать
object_list. -
group – (ProcessGroup, необязательно): Группа процессов для работы. Если None, используется группа по умолчанию. Значение по умолчанию
None. -
device (
torch.device, необязательно) – Если не None, объекты сериализуются и преобразуются в тензоры, которые перемещаются наdeviceперед трансляцией. Значение по умолчаниюNone.
-
object_list (List[Any]) – Список входных объектов для трансляции. Каждый объект должен быть сериализуем (picklable). Только объекты на
- Возвращает:
-
None. Если ранг входит в группу,object_listбудет содержать транслированные объекты сsrcранга.
Примечание
Для групп процессов на основе NCCL внутренние тензорные представления объектов должны быть перемещены на устройство GPU перед началом связи. В этом случае используемое устройство задается
torch.cuda.current_device()и пользователь несет ответственность за обеспечение того, чтобы у каждого ранга было отдельное устройство GPU, используяtorch.cuda.set_device().Примечание
Обратите внимание, что этот API немного отличается от коллектива
all_gather(), так как он не предоставляет дескриптор асинхронной работы и, следовательно, будет блокирующим вызовом.Предупреждение
broadcast_object_list()неявно использует модульpickle, который известен как небезопасный. Возможна конструкция вредоночных данных pickle, которые выполнят произвольный код во время распаковки. Используйте эту функцию только с данными, которым вы доверяете.- Пример::
-
>>> # 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=<torch.distributed.distributed_c10d.ReduceOp object>, group=None, async_op=False)[source] -
Сводит данные тензора по всем машинам таким образом, что все получают окончательный результат.
После вызова
tensorбудет побито идентичен на всех процессах.Поддерживаются сложные тензоры.
- Параметры:
-
- tensor (Тензор) – Вход и выход коллектива. Функция работает in-place.
-
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=<torch.distributed.distributed_c10d.ReduceOp object>, group=None, async_op=False)[source] -
Сводит данные тензора по всем машинам.
Только процесс с рангом
dstполучит окончательный результат.- Параметры:
-
- tensor (Тензор) – Вход и выход коллектива. Функция работает in-place.
- 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 (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 (Tensor) – Выходной тензор для размещения элементов тензоров со всех рангов. Он должен быть правильно размечен, чтобы иметь одну из следующих форм: (i) конкатенацию всех входных тензоров по главной размерности; для определения «конкатенации» см.
torch.cat(); (ii) стопку всех входных тензоров по главной размерности; для определения «стека» см.torch.stack(). Примеры ниже могут лучше объяснить поддерживаемые формы выходных данных. -
input_tensor (Tensor) – Тензор, который должен быть собран с текущего ранга. В отличие от API
all_gather, входные тензоры в этом API должны иметь одинаковый размер на всех рангах. - group (ProcessGroup, необязательно) – Процесс-группа для работы. Если None, будет использоваться группа по умолчанию.
- async_op (bool, необязательно) – Является ли эта операция асинхронной
-
output_tensor (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Предупреждение
Глобальный бэкенд не поддерживает этот API.
-
torch.distributed.all_gather_object(object_list, obj, group=None)[source] -
Сборка пиклируемых объектов из всей группы в список. Аналогично
all_gather(), но можно передавать объекты Python. Обратите внимание, что объект должен быть пиклируемым, чтобы его можно было собрать.- Параметры:
-
- object_list (list[Any]) – Выходной список. Он должен быть правильно размечен, как размер группы для этого коллектива, и будет содержать выходные данные.
- object (Any) – Пиклируемый объект 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, о котором известно, что он небезопасен. Возможно создание вредоносных пиклированных данных, которые выполнят произвольный код во время распиклирования. Вызывайте эту функцию только с данными, которым доверяете.- Пример::
-
>>> # 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 (Tensor) – Входной тензор.
- gather_list (list[Tensor], необязательно) – Список тензоров соответствующего размера для использования со собранными данными (по умолчанию 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 (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, о котором известно, что он небезопасен. Возможно создание вредоносных пиклированных данных, которые выполнят произвольный код во время распиклирования. Вызывайте эту функцию только с данными, которым доверяете.- Пример::
-
>>> # 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 (целое число) – Исходный ранг (по умолчанию 0)
- group (ProcessGroup, необязательно) – Группа процессов для обработки. Если None, будет использована группа по умолчанию.
- async_op (булево, необязательно) – Указывает, будет ли операция асинхронной.
- Возвращает:
-
Дескриптор асинхронной работы, если 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 (Список[Любой]) – Непустой список, первый элемент которого будет хранить объект, рассеянный на данном ранге.
-
scatter_object_input_list (Список[Любой]) – Список входных объектов для рассеивания. Каждый объект должен быть пиклируемым. Только объекты на ранге
srcбудут рассеяны, и аргумент может бытьNoneдля рангов, не являющихся исходными. -
src (целое число) – Исходный ранг, с которого рассеивать
scatter_object_input_list. -
group – (ProcessGroup, необязательно): Группа процессов для обработки. Если None, используется группа по умолчанию. По умолчанию
None.
- Возвращает:
-
None. Если ранг принадлежит группе,scatter_object_output_listбудет иметь свой первый элемент, установленный на рассеянный объект для этого ранга.
Примечание
Этот API отличается от коллектива scatter тем, что не предоставляет
async_opдескриптор и поэтому является блокирующим вызовом.Предупреждение
scatter_object_list()использует модульpickleнеявно, что известно как небезопасный. Возможна конструкция вредоносного пиклируемого данных, которая выполнит произвольный код во время распаковки. Вызывайте эту функцию только с данными, которым вы доверяете.- Пример::
-
>>> # 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=<torch.distributed.distributed_c10d.ReduceOp object>, group=None, async_op=False)[source] -
Сначала выполняет reduce, а затем scatter списка тензоров для всех процессов в группе.
- Параметры:
-
- output (Тензор) – Результирующий тензор.
- input_list (список[Тензор]) – Список тензоров для reduce и scatter.
-
op (необязательно) – Одно из значений из перечисления
torch.distributed.ReduceOp. Определяет операцию, используемую для поэлементных reduce. - group (ProcessGroup, необязательно) – Группа процессов для обработки. Если None, используется группа по умолчанию.
- async_op (булево, необязательно) – Указывает, будет ли операция асинхронной.
- Возвращает:
-
Дескриптор асинхронной работы, если async_op установлено в True. None, если не async_op или если процесс не принадлежит группе.
-
torch.distributed.reduce_scatter_tensor(output, input, op=<torch.distributed.distributed_c10d.ReduceOp object>, group=None, async_op=False)[source] -
Сначала выполняет reduce, а затем scatter тензора для всех рангов в группе.
- Параметры:
-
- output (Тензор) – Результирующий тензор. Он должен иметь одинаковый размер на всех рангах.
-
input (Тензор) – Входной тензор для reduce и scatter. Его размер должен быть размером тензора output умноженным на размер вселенной. Входной тензор может иметь одну из следующих форм: (i) конкатенация тензоров output вдоль первичного измерения или (ii) стопка тензоров output вдоль первичного измерения. Для определения “конкатенации” см.
torch.cat(). Для определения “стопки” см.torch.stack(). - group (ProcessGroup, необязательно) – Группа процессов для обработки. Если None, используется группа по умолчанию.
- async_op (булево, необязательно) – Указывает, будет ли операция асинхронной.
- Возвращает:
-
Дескриптор асинхронной работы, если 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Предупреждение
Глобальный бэкэнд не поддерживает этот API.
-
torch.distributed.all_to_all(output_tensor_list, input_tensor_list, group=None, async_op=False)[source] -
Каждый процесс рассеивает список входных тензоров всем процессам в группе и возвращает собранный список тензоров в списке output.
Поддерживаются сложные тензоры.
- Параметры:
-
- output_tensor_list (список[Тензор]) – Список тензоров, собираемых по одному на ранг.
- input_tensor_list (список[Тензор]) – Список тензоров для рассеивания по одному на ранг.
- group (ProcessGroup, необязательно) – Группа процессов для обработки. Если None, используется группа по умолчанию.
- async_op (булево, необязательно) – Указывает, будет ли операция асинхронной.
- Возвращает:
-
Дескриптор асинхронной работы, если async_op установлено в True. None, если не async_op или если процесс не принадлежит группе.
Предупреждение
all_to_allэкспериментальная функция и может быть изменена.Примеры
>>> input = torch.arange(4) + rank * 4 >>> input = list(input.chunk(4)) >>> input [tensor([0]), tensor([1]), tensor([2]), tensor([3])] # Rank 0 [tensor([4]), tensor([5]), tensor([6]), tensor([7])] # Rank 1 [tensor([8]), tensor([9]), tensor([10]), tensor([11])] # Rank 2 [tensor([12]), tensor([13]), tensor([14]), tensor([15])] # Rank 3 >>> output = list(torch.empty([4], dtype=torch.int64).chunk(4)) >>> dist.all_to_all(output, input) >>> output [tensor([0]), tensor([4]), tensor([8]), tensor([12])] # Rank 0 [tensor([1]), tensor([5]), tensor([9]), tensor([13])] # Rank 1 [tensor([2]), tensor([6]), tensor([10]), tensor([14])] # Rank 2 [tensor([3]), tensor([7]), tensor([11]), tensor([15])] # Rank 3
>>> # Essentially, it is similar to following operation: >>> scatter_list = input >>> gather_list = output >>> for i in range(world_size): >>> dist.scatter(gather_list[i], scatter_list if i == rank else [], src = i)
>>> input tensor([0, 1, 2, 3, 4, 5]) # Rank 0 tensor([10, 11, 12, 13, 14, 15, 16, 17, 18]) # Rank 1 tensor([20, 21, 22, 23, 24]) # Rank 2 tensor([30, 31, 32, 33, 34, 35, 36]) # Rank 3 >>> input_splits [2, 2, 1, 1] # Rank 0 [3, 2, 2, 2] # Rank 1 [2, 1, 1, 1] # Rank 2 [2, 2, 2, 1] # Rank 3 >>> output_splits [2, 3, 2, 2] # Rank 0 [2, 2, 1, 2] # Rank 1 [1, 2, 1, 2] # Rank 2 [1, 2, 1, 1] # Rank 3 >>> input = list(input.split(input_splits)) >>> input [tensor([0, 1]), tensor([2, 3]), tensor([4]), tensor([5])] # Rank 0 [tensor([10, 11, 12]), tensor([13, 14]), tensor([15, 16]), tensor([17, 18])] # Rank 1 [tensor([20, 21]), tensor([22]), tensor([23]), tensor([24])] # Rank 2 [tensor([30, 31]), tensor([32, 33]), tensor([34, 35]), tensor([36])] # Rank 3 >>> output = ... >>> dist.all_to_all(output, input) >>> output [tensor([0, 1]), tensor([10, 11, 12]), tensor([20, 21]), tensor([30, 31])] # Rank 0 [tensor([2, 3]), tensor([13, 14]), tensor([22]), tensor([32, 33])] # Rank 1 [tensor([4]), tensor([15, 16]), tensor([23]), tensor([34, 35])] # Rank 2 [tensor([5]), tensor([17, 18]), tensor([24]), tensor([36])] # Rank 3
>>> # Another example with tensors of torch.cfloat type. >>> input = torch.tensor([1+1j, 2+2j, 3+3j, 4+4j], dtype=torch.cfloat) + 4 * rank * (1+1j) >>> input = list(input.chunk(4)) >>> input [tensor([1+1j]), tensor([2+2j]), tensor([3+3j]), tensor([4+4j])] # Rank 0 [tensor([5+5j]), tensor([6+6j]), tensor([7+7j]), tensor([8+8j])] # Rank 1 [tensor([9+9j]), tensor([10+10j]), tensor([11+11j]), tensor([12+12j])] # Rank 2 [tensor([13+13j]), tensor([14+14j]), tensor([15+15j]), tensor([16+16j])] # Rank 3 >>> output = list(torch.empty([4], dtype=torch.int64).chunk(4)) >>> dist.all_to_all(output, input) >>> output [tensor([1+1j]), tensor([5+5j]), tensor([9+9j]), tensor([13+13j])] # Rank 0 [tensor([2+2j]), tensor([6+6j]), tensor([10+10j]), tensor([14+14j])] # Rank 1 [tensor([3+3j]), tensor([7+7j]), tensor([11+11j]), tensor([15+15j])] # Rank 2 [tensor([4+4j]), tensor([8+8j]), tensor([12+12j]), tensor([16+16j])] # Rank 3
-
torch.distributed.barrier(group=None, async_op=False, device_ids=None)[source] -
Синхронизирует все процессы.
Этот коллектив блокирует процессы до тех пор, пока вся группа не войдёт в эту функцию, если async_op равен False, или если обработчик асинхронной работы вызывается в wait().
- Параметры:
-
- group (ProcessGroup, необязательно) – Группа процессов, над которой нужно работать. Если None, будет использоваться группа по умолчанию.
- async_op (bool, необязательно) – Является ли эта операция асинхронной
- device_ids ([int], необязательно) – Список идентификаторов устройств/GPU. Действителен только для бэкенда NCCL.
- Возвращаемое значение:
-
Обработчик асинхронной работы, если 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 для профилирования коллективной коммуникации и точечно-точечных коммуникационных API, упомянутых здесь. Все встроенные бэкенды (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 есть тензор, который мы хотим выполнить all-reduce. Следующий код может служить ссылкой:
Код, выполняющийся на узле 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 тензоров на двух узлах будут иметь значение all-reduce, равное 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=<torch.distributed.distributed_c10d.ReduceOp object>, group=None, async_op=False)[source] -
Производит сокращение данных тензора по всем машинам таким образом, что все получают окончательный результат. Эта функция сокращает количество тензоров на каждом узле, при этом каждый тензор находится на разных GPU. Таким образом, входной тензор в списке тензоров должен быть тензором GPU. Кроме того, каждый тензор в списке тензоров должен находиться на разных GPU.
После вызова все
tensorвtensor_listбудут точно такими же во всех процессах.Поддерживаются сложные тензоры.
В настоящее время поддерживаются только бэкэнды nccl и gloo, тензоры должны быть только тензорами GPU
- Параметры:
-
-
tensor_list (Список[Tensor]) – Список входных и выходных тензоров коллектива. Функция работает на месте и требует, чтобы каждый тензор был тензором GPU на разных GPU. Вы также должны убедиться, что
len(tensor_list)одинаково для всех распределённых процессов, вызывающих эту функцию. -
op (необязательно) – Одно из значений из
torch.distributed.ReduceOpперечисления. Указывает операцию, используемую для поэлементных сокращений. -
group (ProcessGroup, необязательно) – Группа процессов для работы. Если
None, будет использоваться группа процессов по умолчанию. - async_op (bool, необязательно) – Является ли эта операция асинхронной операцией.
-
tensor_list (Список[Tensor]) – Список входных и выходных тензоров коллектива. Функция работает на месте и требует, чтобы каждый тензор был тензором GPU на разных GPU. Вы также должны убедиться, что
- Возвращаемые значения:
-
Дескриптор асинхронной работы, если async_op установлен в True. None, если не async_op или если не принадлежит группе
-
torch.distributed.reduce_multigpu(tensor_list, dst, op=<torch.distributed.distributed_c10d.ReduceOp object>, group=None, async_op=False, dst_tensor=0)[source] -
Сокращает данные тензора на нескольких GPU по всем машинам. Каждый тензор в
tensor_listдолжен находиться на отдельном GPUТолько GPU
tensor_list[dst_tensor]в процессе с рангомdstполучит окончательный результат.В настоящее время поддерживается только бэкенд nccl, тензоры должны быть только тензорами GPU
- Параметры:
-
-
tensor_list (Список[Tensor]) – Входные и выходные тензоры GPU коллектива. Функция работает на месте. Вы также должны убедиться, что
len(tensor_list)одинаково для всех распределённых процессов, вызывающих эту функцию. - dst (int) – Ранг назначения
-
op (необязательно) – Одно из значений из
torch.distributed.ReduceOpперечисления. Указывает операцию, используемую для поэлементных сокращений. - group (ProcessGroup, необязательно) – Группа процессов для работы. Если None, будет использоваться группа процессов по умолчанию.
- async_op (bool, необязательно) – Является ли эта операция асинхронной операцией.
-
dst_tensor (int, необязательно) – Ранг тензора назначения в
tensor_list
-
tensor_list (Список[Tensor]) – Входные и выходные тензоры GPU коллектива. Функция работает на месте. Вы также должны убедиться, что
- Возвращаемые значения:
-
Дескриптор асинхронной работы, если async_op установлен в True. Иначе None
-
torch.distributed.all_gather_multigpu(output_tensor_lists, input_tensor_list, group=None, async_op=False)[source] -
Собрать тензоры из всей группы в список. Каждый тензор в
tensor_listдолжен находиться на отдельном GPU.В настоящее время поддерживается только бэкэнд nccl; тензоры должны быть только тензорами GPU.
Поддерживаются сложные тензоры.
- Параметры:
-
-
output_tensor_lists (Список[Список[Тензор]]) –
Выходные списки. Он должен содержать тензоры правильного размера на каждом GPU, которые будут использоваться для вывода коллектива, например,
output_tensor_lists[i]содержит результат all_gather, находящийся на GPUinput_tensor_list[i].Обратите внимание, что каждый элемент
output_tensor_listsимеет размерworld_size * len(input_tensor_list), так как функция собирает результат с каждого отдельного GPU в группе. Для интерпретации каждого элемента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 (Список[Тензор]) – Список тензоров (на разных GPU) для трансляции из текущего процесса. Обратите внимание, что
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=<torch.distributed.distributed_c10d.ReduceOp object>, group=None, async_op=False)[source] -
Сведение и рассеивание списка тензоров по всей группе. В настоящее время поддерживается только бэкэнд nccl.
Каждый тензор в
output_tensor_listдолжен находиться на отдельном GPU, как и каждый список тензоров вinput_tensor_lists.- Параметры:
-
-
output_tensor_list (Список[Тензор]) –
Выходные тензоры (на разных GPU) для получения результата операции.
Обратите внимание, что
len(output_tensor_list)должен быть одинаковым для всех распределенных процессов, вызывающих эту функцию. -
input_tensor_lists (Список[Список[Тензор]]) –
Входные списки. Он должен содержать тензоры правильного размера на каждом GPU, которые будут использоваться для входных данных коллектива, например,
input_tensor_lists[i]содержит входные данные reduce_scatter, находящиеся на GPUoutput_tensor_list[i].Обратите внимание, что каждый элемент
input_tensor_listsимеет размерworld_size * len(output_tensor_list), так как функция рассеивает результат с каждого отдельного GPU в группе. Для интерпретации каждого элемента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++, обратитесь к 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.
Утилита может использоваться для распределённого обучения на одном узле, в котором будет запущено один или несколько процессов на узел. Утилита может использоваться для обучения как на CPU, так и на GPU. Если утилита используется для обучения на GPU, каждый распределённый процесс будет работать с одним GPU. Это может значительно улучшить производительность обучения на одном узле. Она также может использоваться для распределённого обучения на нескольких узлах, запуская несколько процессов на каждом узле, что также улучшает производительность распределённого обучения на нескольких узлах. Это особенно полезно для систем с несколькими интерфейсами Infiniband, поддерживающими прямое взаимодействие с GPU, так как все они могут быть использованы для агрегированной пропускной способности связи.
В обоих случаях, для распределённого обучения на одном узле или на нескольких узлах, эта утилита запустит заданное количество процессов на узел (--nproc_per_node). При использовании для обучения на GPU это число должно быть меньше или равно количеству GPU в текущей системе (nproc_per_node), и каждый процесс будет работать с одним GPU от 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. Эта утилита и распределённое обучение с несколькими процессами (на одном или нескольких узлах) с использованием GPU в настоящее время достигает наилучшей производительности только с использованием распределённого бэкенда NCCL. Таким образом, бэкенд NCCL рекомендуется для использования при обучении на GPU.
2. В вашей обучающей программе необходимо обработать аргумент командной строки: --local_rank=LOCAL_PROCESS_RANK, который будет предоставлен этим модулем. Если ваша обучающая программа использует GPU, убедитесь, что ваш код работает только на устройстве GPU с LOCAL_PROCESS_RANK. Это можно сделать, используя:
Обработку аргумента local_rank
>>> import argparse
>>> parser = argparse.ArgumentParser()
>>> parser.add_argument("--local_rank", type=int)
>>> args = parser.parse_args()
Установите устройство на local rank, используя:
>>> 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(). Если ваша обучающая программа использует GPU для обучения и вы хотите использовать модуль torch.nn.parallel.DistributedDataParallel(), вот как его настроить.
>>> model = torch.nn.parallel.DistributedDataParallel(model, >>> device_ids=[args.local_rank], >>> output_device=args.local_rank)
Пожалуйста, убедитесь, что аргумент device_ids установлен в единственный ID устройства GPU, на котором ваш код будет работать. Обычно это локальный ранк процесса. Другими словами, device_ids должен быть [args.local_rank], а output_device должен быть args.local_rank, чтобы использовать эту утилиту.
5. Ещё один способ передать local_rank подпроцессам через переменную окружения LOCAL_RANK. Это поведение включено, когда вы запускаете скрипт с --use_env=True. Вы должны изменить пример подпроцесса выше, чтобы заменить args.local_rank на os.environ['LOCAL_RANK']; загрузчик не передаст --local_rank при указании этого флага.
Предупреждение
local_rank НЕ является глобально уникальным: он уникален только для каждого процесса на машине. Поэтому не используйте его для принятия решений, например, о записи в сетевую файловую систему. См. https://github.com/pytorch/pytorch/issues/12042 для примера того, как могут возникнуть проблемы, если вы не сделаете это правильно.
Утилита создания процессов
Пакет Пакет Multiprocessing - 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.
|
| Эффективный уровень журнала |
|---|---|---|
| игнорируется | Ошибка |
| игнорируется | Предупреждение |
| игнорируется | Информация |
|
| Отладка |
|
| Отслеживание (т. е. Все) |
© 2024, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://pytorch.org/docs/1.13/distributed.html