{/ Эта страница автоматически генерируется из SKILL.md навыка с помощью website/scripts/generate-skill-docs.py. Редактируйте исходный SKILL.md, а не эту страницу. /}
Pytorch Fsdp
Экспертное руководство по Fully Sharded Data Parallel обучению с PyTorch FSDP — шардирование параметров, смешанная точность, выгрузка на CPU, FSDP2
Метаданные навыка
| Источник | Опционально — установка: hermes skills install official/mlops/pytorch-fsdp |
| Путь | optional-skills/mlops/pytorch-fsdp |
| Версия | 1.0.0 |
| Автор | Orchestra Research |
| Лицензия | MIT |
| Зависимости | torch>=2.0, transformers |
| Платформы | linux, macos |
| Теги | Distributed Training, PyTorch, FSDP, Data Parallel, Sharding, Mixed Precision, CPU Offloading, FSDP2, Large-Scale Training |
Справочник: полный SKILL.mdℹ️ Info
Ниже приведено полное определение навыка, которое Hermes загружает при его активации. Это то, что агент видит в качестве инструкций, когда навык активен.
Навык Pytorch-Fsdp
ℹ️ Info
Ниже приведено полное определение навыка, которое Hermes загружает при его активации. Это то, что агент видит в качестве инструкций, когда навык активен.
Всесторонняя помощь по разработке с pytorch-fsdp, сгенерированная из официальной документации.
Когда использовать этот навык
Этот навык следует активировать, когда: - Вы работаете с pytorch-fsdp - Вы задаете вопросы о возможностях или API pytorch-fsdp - Вы реализуете решения на pytorch-fsdp - Вы отлаживаете код pytorch-fsdp - Вы изучаете лучшие практики pytorch-fsdp
Краткий справочник
Распространенные шаблоны
Шаблон 1: Generic Join Context Manager# Создано: 06 июня 2025 | Последнее обновление: 06 июня 2025 Generic join context manager облегчает распределенное обучение на неравномерных входных данных. На этой странице описан API соответствующих классов: Join, Joinable и JoinHook. Учебное пособие см. в разделе Distributed Training with Uneven Inputs Using the Join Context Manager. class torch.distributed.algorithms.Join(joinables, enable=True, throw_on_early_termination=False, kwargs)[source]# Этот класс определяет generic join context manager, который позволяет вызывать пользовательские хуки после того, как процесс присоединился. Эти хуки должны затенять коллективные коммуникации неприсоединившихся процессов, чтобы предотвратить зависания и ошибки и обеспечить алгоритмическую корректность. Обратитесь к JoinHook за подробностями об определении хука. Предупреждение Контекстный менеджер требует, чтобы каждый участвующий Joinable вызывал метод notify_join_context() перед своими собственными коллективными коммуникациями на каждой итерации для обеспечения корректности. Предупреждение Контекстный менеджер требует, чтобы все атрибуты process_group в объектах JoinHook были одинаковыми. Если существует несколько объектов JoinHook, то используется устройство первого из них. Информация о группе процессов и устройстве используется для проверки неприсоединившихся процессов и для уведомления процессов о необходимости выбросить исключение, если включен throw_on_early_termination, и то и другое с использованием all-reduce. Параметры joinables (List[Joinable]) – список участвующих Joinable; их хуки перебираются в заданном порядке. enable (bool) – флаг, включающий обнаружение неравномерных входных данных; установка False отключает функциональность контекстного менеджера и должна устанавливаться только тогда, когда пользователь знает, что входные данные не будут неравномерными (по умолчанию: True). throw_on_early_termination (bool) – флаг, управляющий тем, следует ли выбрасывать исключение при обнаружении неравномерных входных данных (по умолчанию: False). Пример: >>> import os >>> import torch >>> import torch.distributed as dist >>> import torch.multiprocessing as mp >>> import torch.nn.parallel.DistributedDataParallel as DDP >>> import torch.distributed.optim.ZeroRedundancyOptimizer as ZeRO >>> from torch.distributed.algorithms.join import Join >>> >>> # На каждом порожденном рабочем >>> def worker(rank): >>> dist.init_process_group("nccl", rank=rank, world_size=2) >>> model = DDP(torch.nn.Linear(1, 1).to(rank), device_ids=[rank]) >>> optim = ZeRO(model.parameters(), torch.optim.Adam, lr=0.01) >>> # Ранг 1 получает на один вход больше, чем ранг 0 >>> inputs = [torch.tensor([1.]).to(rank) for _ in range(10 + rank)] >>> with Join([model, optim]): >>> for input in inputs: >>> loss = model(input).sum() >>> loss.backward() >>> optim.step() >>> # Все ранги достигают этого места без зависания/ошибки static notify_join_context(joinable)[source]# Уведомляет join context manager о том, что вызывающий процесс еще не присоединился. Затем, если throw_on_early_termination=True, проверяет, были ли обнаружены неравномерные входные данные (т.е. если один процесс уже присоединился), и выбрасывает исключение, если это так. Этот метод должен вызываться из объекта Joinable перед его коллективными коммуникациями на каждой итерации. Например, он должен вызываться в начале прямого прохода в DistributedDataParallel. Только первый объект Joinable, переданный в контекстный менеджер, выполняет коллективные коммуникации в этом методе; для остальных этот метод является пустым. Параметры joinable (Joinable) – объект Joinable, вызывающий этот метод. Возвращает Асинхронный дескриптор работы для all-reduce, предназначенный для уведомления контекстного менеджера о том, что процесс еще не присоединился, если joinable является первым, переданным в контекстный менеджер; в противном случае None. class torch.distributed.algorithms.Joinable[source]# Определяет абстрактный базовый класс для joinable классов. Joinable класс (наследующий от Joinable) должен реализовывать join_hook(), который возвращает экземпляр JoinHook, а также join_device() и join_process_group(), которые возвращают информацию об устройстве и группе процессов соответственно. abstract property join_device: device# Возвращает устройство, с которого выполнять коллективные коммуникации, необходимые join context manager. abstract join_hook(kwargs)[source]# Возвращает экземпляр JoinHook для данного Joinable. Параметры kwargs (dict) – словарь, содержащий любые ключевые аргументы для изменения поведения join hook во время выполнения; все экземпляры Joinable, использующие один и тот же join context manager, получают одно и то же значение kwargs. Тип возвращаемого значения JoinHook abstract property join_process_group: Any# Возвращает группу процессов для коллективных коммуникаций, необходимых самому join context manager. class torch.distributed.algorithms.JoinHook[source]# Определяет join hook, который предоставляет две точки входа в join context manager. Точки входа: основной хук, который вызывается повторно, пока существует неприсоединившийся процесс, и пост-хук, который вызывается один раз, когда все процессы присоединились. Чтобы реализовать join hook для generic join context manager, определите класс, наследующий от JoinHook, и переопределите main_hook() и post_hook() по мере необходимости. main_hook()[source]# Вызывайте этот хук, пока существует неприсоединившийся процесс, чтобы затенять коллективные коммуникации в итерации обучения. Итерация обучения, т.е. один прямой проход, обратный проход и шаг оптимизатора. post_hook(is_last_joiner)[source]# Вызывайте хук после того, как все процессы присоединились. Ему передается дополнительный аргумент bool is_last_joiner, который указывает, является ли ранг одним из последних присоединившихся. Параметры is_last_joiner (bool) – True, если ранг является одним из последних присоединившихся; в противном случае False.
Join
Шаблон 2: Пакет распределенной коммуникации - torch.distributed# Создано: 12 июля 2017 | Последнее обновление: 04 сентября 2025 Примечание Пожалуйста, обратитесь к PyTorch Distributed Overview для краткого введения во все функции, связанные с распределенным обучением. Бэкенды# torch.distributed поддерживает четыре встроенных бэкенда, каждый с разными возможностями. В таблице ниже показано, какие функции доступны для использования с CPU или GPU для каждого бэкенда. Для NCCL GPU означает CUDA GPU, а для XCCL — XPU GPU. MPI поддерживает CUDA только в том случае, если реализация, используемая для сборки PyTorch, поддерживает его. Бэкенд gloo mpi nccl xccl Устройство CPU GPU CPU GPU CPU GPU CPU GPU send ✓ ✘ ✓? ✘ ✓ ✘ ✓ recv ✓ ✘ ✓? ✘ ✓ ✘ ✓ broadcast ✓ ✓ ✓? ✘ ✓ ✘ ✓ all_reduce ✓ ✓ ✓? ✘ ✓ ✘ ✓ reduce ✓ ✓ ✓? ✘ ✓ ✘ ✓ all_gather ✓ ✓ ✓? ✘ ✓ ✘ ✓ gather ✓ ✓ ✓? ✘ ✓ ✘ ✓ scatter ✓ ✓ ✓? ✘ ✓ ✘ ✓ reduce_scatter ✓ ✓ ✘ ✘ ✘ ✓ ✘ ✓ all_to_all ✓ ✓ ✓? ✘ ✓ ✘ ✓ barrier ✓ ✘ ✓? ✘ ✓ ✘ ✓ Бэкенды, поставляемые с PyTorch# Пакет распределенных вычислений PyTorch поддерживает Linux (стабильно), MacOS (стабильно) и Windows (прототип). По умолчанию для Linux бэкенды Gloo и NCCL собираются и включаются в состав распределенного пакета PyTorch (NCCL только при сборке с CUDA). MPI — это опциональный бэкенд, который можно включить только при сборке PyTorch из исходного кода (например, сборка PyTorch на хосте с установленным MPI). Примечание Начиная с PyTorch v1.8, Windows поддерживает все бэкенды коллективных коммуникаций, кроме NCCL. Если аргумент init_method функции init_process_group() указывает на файл, он должен соответствовать следующей схеме: Локальная файловая система, init_method="file:///d:/tmp/some_file" Общая файловая система, init_method="file://////{machine_name}/{share_folder_name}/some_file" Как и на платформе Linux, вы можете включить TcpStore, установив переменные окружения MASTER_ADDR и MASTER_PORT. Какой бэкенд использовать?# Раньше нас часто спрашивали: «Какой бэкенд мне следует использовать?». Эмпирическое правило Используйте бэкенд NCCL для распределенного обучения с CUDA GPU. Используйте бэкенд XCCL для распределенного обучения с XPU GPU. Используйте бэкенд Gloo для распределенного обучения с CPU. Хосты с GPU и InfiniBand Используйте NCCL, так как это единственный бэкенд, который в настоящее время поддерживает InfiniBand и GPUDirect. Хосты с GPU и Ethernet Используйте NCCL, так как он в настоящее время обеспечивает наилучшую производительность распределенного обучения на GPU, особенно для многопроцессорного одноузлового или многоузлового распределенного обучения. Если у вас возникли проблемы с NCCL, используйте Gloo в качестве запасного варианта. (Обратите внимание, что Gloo в настоящее время работает медленнее, чем NCCL для GPU.) Хосты с CPU и InfiniBand Если ваш InfiniBand поддерживает IP over IB, используйте Gloo, в противном случае используйте MPI. Мы планируем добавить поддержку InfiniBand для Gloo в будущих релизах. Хосты с CPU и Ethernet Используйте Gloo, если у вас нет особых причин использовать MPI. Общие переменные окружения# Выбор сетевого интерфейса# По умолчанию бэкенды NCCL и Gloo пытаются найти правильный сетевой интерфейс. Если автоматически обнаруженный интерфейс неверен, вы можете переопределить его с помощью следующих переменных окружения (применимо к соответствующему бэкенду): NCCL_SOCKET_IFNAME, например export NCCL_SOCKET_IFNAME=eth0 GLOO_SOCKET_IFNAME, например export GLOO_SOCKET_IFNAME=eth0 Если вы используете бэкенд Gloo, вы можете указать несколько интерфейсов, разделив их запятой, например: export GLOO_SOCKET_IFNAME=eth0,eth1,eth2,eth3. Бэкенд будет распределять операции между этими интерфейсами по круговому принципу. Важно, чтобы все процессы указывали одинаковое количество интерфейсов в этой переменной. Другие переменные окружения NCCL# Отладка — в случае сбоя NCCL вы можете установить NCCL_DEBUG=INFO, чтобы вывести явное предупреждающее сообщение, а также базовую информацию об инициализации NCCL. Вы также можете использовать NCCL_DEBUG_SUBSYS для получения более подробной информации о конкретном аспекте NCCL. Например, NCCL_DEBUG_SUBSYS=COLL выведет журналы коллективных вызовов, что может быть полезно при отладке зависаний, особенно вызванных несоответствием типа коллектива или размера сообщения. В случае сбоя обнаружения топологии будет полезно установить NCCL_DEBUG_SUBSYS=GRAPH для просмотра подробного результата обнаружения и сохранения его в качестве справки, если потребуется дальнейшая помощь от команды NCCL. Настройка производительности — NCCL выполняет автоматическую настройку на основе обнаружения топологии, чтобы избавить пользователей от необходимости ручной настройки. В некоторых системах на основе сокетов пользователи все же могут попробовать настроить NCCL_SOCKET_NTHREADS и NCCL_NSOCKS_PERTHREAD для увеличения пропускной способности сокетной сети. Эти две переменные окружения были предварительно настроены NCCL для некоторых облачных провайдеров, таких как AWS или GCP. Полный список переменных окружения NCCL см. в официальной документации NVIDIA NCCL. Вы можете дополнительно настроить коммуникаторы NCCL с помощью torch.distributed.ProcessGroupNCCL.NCCLConfig и torch.distributed.ProcessGroupNCCL.Options. Узнайте больше о них, используя help (например, help(torch.distributed.ProcessGroupNCCL.NCCLConfig)) в интерпретаторе. Основы# Пакет torch.distributed предоставляет поддержку PyTorch и примитивы коммуникации для многопроцессорного параллелизма на нескольких вычислительных узлах, работающих на одной или нескольких машинах. Класс torch.nn.parallel.DistributedDataParallel() построен на этой функциональности, чтобы обеспечить синхронное распределенное обучение в качестве обертки вокруг любой модели PyTorch. Это отличается от видов параллелизма, предоставляемых пакетом Multiprocessing package - torch.multiprocessing и torch.nn.DataParallel(), тем, что поддерживает несколько машин, соединенных сетью, и тем, что пользователь должен явно запускать отдельную копию основного обучающего скрипта для каждого процесса. В синхронном случае на одной машине torch.distributed или обертка torch.nn.parallel.DistributedDataParallel() все еще могут иметь преимущества перед другими подходами к параллелизму данных, включая torch.nn.DataParallel(): Каждый процесс поддерживает свой собственный оптимизатор и выполняет полный шаг оптимизации на каждой итерации. Хотя это может показаться избыточным, поскольку градиенты уже были собраны и усреднены по процессам и, следовательно, одинаковы для каждого процесса, это означает, что шаг широковещательной рассылки параметров не требуется, что сокращает время передачи тензоров между узлами. Каждый процесс содержит независимый интерпретатор Python, что устраняет дополнительные накладные расходы на интерпретатор и «GIL-трэшинг», возникающие при управлении несколькими потоками выполнения, репликами модели или GPU из одного процесса Python. Это особенно важно для моделей, которые активно используют среду выполнения Python, включая модели с рекуррентными слоями или множеством мелких компонентов. Инициализация# Пакет необходимо инициализировать с помощью функции torch.distributed.init_process_group() или torch.distributed.device_mesh.init_device_mesh() перед вызовом любых других методов. Обе блокируются до тех пор, пока все процессы не присоединятся. Предупреждение Инициализация не является потокобезопасной. Создание группы процессов должно выполняться из одного потока, чтобы предотвратить непоследовательное назначение «UUID» между рангами и предотвратить состояния гонки во время инициализации, которые могут привести к зависаниям. torch.distributed.is_available()[source]# Возвращает True, если пакет распределенных вычислений доступен. В противном случае torch.distributed не предоставляет никаких других API. В настоящее время torch.distributed доступен на Linux, MacOS и Windows. Установите USE_DISTRIBUTED=1, чтобы включить его при сборке PyTorch из исходного кода. В настоящее время значение по умолчанию — USE_DISTRIBUTED=1 для Linux и Windows, USE_DISTRIBUTED=0 для MacOS. Тип возвращаемого значения bool torch.distributed.init_process_group(backend=None, init_method=None, timeout=None, world_size=-1, rank=-1, store=None, group_name='', pg_options=None, device_id=None)[source]# Инициализирует группу процессов по умолчанию. Это также инициализирует пакет распределенных вычислений. Есть 2 основных способа инициализации группы процессов: Явно указать store, rank и world_size. Указать init_method (строку URL), которая указывает, где/как обнаруживать пиров. Опционально указать rank и world_size или закодировать все необходимые параметры в URL и опустить их. Если не указано ни то, ни другое, init_method считается равным "env://". Параметры backend (str или Backend, опционально) – Используемый бэкенд. В зависимости от конфигурации сборки допустимые значения включают mpi, gloo, nccl, ucc, xccl или зарегистрированные сторонним плагином. Начиная с версии 2.6, если backend не указан, c10d будет использовать бэкенд, зарегистрированный для типа устройства, указанного в kwargs device_id (если он предоставлен). Известные регистрации по умолчанию на сегодня: nccl для cuda, gloo для cpu, xccl для xpu. Если не указаны ни backend, ни device_id, c10d обнаружит ускоритель на машине времени выполнения и использует бэкенд, зарегистрированный для этого обнаруженного ускорителя (или cpu). Это поле может быть задано в виде строки в нижнем регистре (например, "gloo"), к которой также можно получить доступ через атрибуты Backend (например, Backend.GLOO). При использовании нескольких процессов на машине с бэкендом nccl каждый процесс должен иметь эксклюзивный доступ к каждому используемому GPU, поскольку совместное использование GPU между процессами может привести к взаимоблокировке или недопустимому использованию NCCL. Бэкенд ucc является экспериментальным. Бэкенд по умолчанию для устройства можно запросить с помощью get_default_backend_for_device(). 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, опционально) – Тайм-аут для операций, выполняемых в группе процессов. Значение по умолчанию — 10 минут для NCCL и 30 минут для других бэкендов. Это продолжительность, после которой коллективы будут асинхронно прерваны, и процесс завершится с ошибкой. Это сделано потому, что выполнение CUDA является асинхронным, и больше небезопасно продолжать выполнение пользовательского кода, поскольку неудачные асинхронные операции NCCL могут привести к тому, что последующие операции CUDA будут выполняться с поврежденными данными. Когда установлен TORCH_NCCL_BLOCKING_WAIT, процесс будет блокироваться и ждать этого тайм-аута. group_name (str, опционально, устарело) – Имя группы. Этот аргумент игнорируется. pg_options (ProcessGroupOptions, опционально) – Параметры группы процессов, указывающие, какие дополнительные параметры необходимо передать при создании конкретных групп процессов. На данный момент единственная поддерживаемая опция — ProcessGroupNCCL.Options для бэкенда nccl; is_high_priority_stream может быть указан, чтобы бэкенд nccl мог использовать потоки CUDA с высоким приоритетом, когда ожидают вычислительные ядра. Другие доступные опции для настройки nccl см. https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/api/types.html#ncclconfig-t device_id (torch.device | int, опционально) – одно конкретное устройство, с которым будет работать этот процесс, что позволяет выполнять оптимизацию, специфичную для бэкенда. В настоящее время это имеет два эффекта, только под NCCL: коммуникатор формируется немедленно (вызов ncclCommInit сразу, а не обычный ленивый вызов), и подгруппы будут использовать ncclCommSplit, когда это возможно, чтобы избежать ненужных накладных расходов на создание группы. Если вы хотите узнать об ошибке инициализации NCCL на раннем этапе, вы также можете использовать это поле. Если указано int, API предполагает, что будет использоваться тип ускорителя во время компиляции. Примечание Чтобы включить backend == Backend.MPI, PyTorch должен быть собран из исходного кода в системе, поддерживающей MPI. Примечание Поддержка нескольких бэкендов является экспериментальной. В настоящее время, если бэкенд не указан, будут созданы бэкенды gloo и nccl. Бэкенд gloo будет использоваться для коллективов с тензорами CPU, а бэкенд nccl — для коллективов с тензорами CUDA. Пользовательский бэкенд может быть указан путем передачи строки в формате "<device_type>:<backend_name>,<device_type>:<backend_name>", например "cpu:gloo,cuda:custom_backend". torch.distributed.device_mesh.init_device_mesh(device_type, mesh_shape, , mesh_dim_names=None, backend_override=None)[source]# Инициализирует DeviceMesh на основе параметров device_type, mesh_shape и mesh_dim_names. Это создает DeviceMesh с n-мерным массивом, где n — длина mesh_shape. Если указан mesh_dim_names, каждое измерение помечается как mesh_dim_names[i]. Примечание init_device_mesh следует модели программирования SPMD, что означает, что одна и та же программа Python PyTorch выполняется на всех процессах/рангах в кластере. Убедитесь, что mesh_shape (размеры nD массива, описывающего расположение устройств) идентичен во всех рангах. Несогласованный mesh_shape может привести к зависанию. Примечание Если группа процессов не найдена, init_device_mesh инициализирует группу/группы процессов, необходимые для распределенных коммуникаций, «за кулисами». Параметры device_type (str) – Тип устройства сетки. В настоящее время поддерживается: "cpu", "cuda/cuda-like", "xpu". Передача типа устройства с индексом GPU, например "cuda:0", не допускается. mesh_shape (Tuple[int]) – Кортеж, определяющий размеры многомерного массива, описывающего расположение устройств. mesh_dim_names (Tuple[str], опционально) – Кортеж имен измерений сетки для назначения каждому измерению многомерного массива, описывающего расположение устройств. Его длина должна совпадать с длиной mesh_shape. Каждая строка в mesh_dim_names должна быть уникальной. backend_override (Dict[int | str, tuple[str, Options] | str | Options], опционально) – Переопределения для некоторых или всех ProcessGroups, которые будут созданы для каждого измерения сетки. Каждый ключ может быть либо индексом измерения, либо его именем (если указан mesh_dim_names). Каждое значение может быть кортежем, содержащим имя бэкенда и его параметры, или только одним из этих двух компонентов (в этом случае другой будет установлен в значение по умолчанию). Возвращает Объект DeviceMesh, представляющий расположение устройств. Тип возвращаемого значения DeviceMesh Пример: >>> from torch.distributed.device_mesh import init_device_mesh >>> >>> mesh_1d = init_device_mesh("cuda", mesh_shape=(8,)) >>> mesh_2d = init_device_mesh("cuda", mesh_shape=(2, 8), mesh_dim_names=("dp", "tp")) torch.distributed.is_initialized()[source]# Проверяет, была ли инициализирована группа процессов по умолчанию. Тип возвращаемого значения bool torch.distributed.is_mpi_available()[source]# Проверяет, доступен ли бэкенд MPI. Тип возвращаемого значения bool torch.distributed.is_nccl_available()[source]# Проверяет, доступен ли бэкенд NCCL. Тип возвращаемого значения bool torch.distributed.is_gloo_available()[source]# Проверяет, доступен ли бэкенд Gloo. Тип возвращаемого значения bool torch.distributed.distributed_c10d.is_xccl_available()[source]# Проверяет, доступен ли бэкенд XCCL. Тип возвращаемого значения bool torch.distributed.is_torchelastic_launched()[source]# Проверяет, был ли этот процесс запущен с помощью torch.distributed.elastic (также известного как torchelastic). Наличие переменной окружения TORCHELASTIC_RUN_ID используется в качестве прокси для определения того, был ли текущий процесс запущен с помощью torchelastic. Это разумный прокси, поскольку TORCHELASTIC_RUN_ID соответствует идентификатору рандеву, который всегда является ненулевым значением, указывающим идентификатор задания для целей обнаружения пиров. Тип возвращаемого значения bool torch.distributed.get_default_backend_for_device(device)[source]# Возвращает бэкенд по умолчанию для данного устройства. Параметры device (Union[str, torch.device]) – Устройство, для которого нужно получить бэкенд по умолчанию. Возвращает Бэкенд по умолчанию для данного устройства в виде строки в нижнем регистре. Тип возвращаемого значения str В настоящее время поддерживаются три метода инициализации: TCP инициализация# Есть два способа инициализации с использованием TCP, оба требуют сетевого адреса, доступного из всех процессов, и желаемого world_size. Первый способ требует указания адреса, принадлежащего процессу с рангом 0. Этот метод инициализации требует, чтобы все процессы вручную указали ранги. Обратите внимание, что многоадресный адрес больше не поддерживается в последней версии пакета распределенных вычислений. group_name также устарел. import torch.distributed as dist # Используйте адрес одной из машин dist.init_process_group(backend, init_method='tcp://10.1.1.20:23456', rank=args.rank, world_size=4) Инициализация через общую файловую систему# Другой метод инициализации использует файловую систему, которая является общей и видна со всех машин в группе, вместе с желаемым world_size. URL должен начинаться с file:// и содержать путь к несуществующему файлу (в существующем каталоге) в общей файловой системе. Файловая инициализация автоматически создаст этот файл, если он не существует, но не удалит его. Поэтому вы несете ответственность за то, чтобы файл был очищен перед следующим вызовом init_process_group() с тем же путем/именем файла. Обратите внимание, что автоматическое назначение рангов больше не поддерживается в последней версии пакета распределенных вычислений, и group_name также устарел. Предупреждение Этот метод предполагает, что файловая система поддерживает блокировку с помощью fcntl — большинство локальных систем и NFS поддерживают ее. Предупреждение Этот метод всегда будет создавать файл и приложит все усилия для его очистки и удаления в конце программы. Другими словами, каждая инициализация с помощью файлового метода init будет нуждаться в совершенно новом пустом файле для успешной инициализации. Если тот же файл, использованный в предыдущей инициализации (который не был очищен), используется снова, это неожиданное поведение и может часто вызывать взаимоблокировки и сбои. Поэтому, даже если этот метод приложит все усилия для очистки файла, если автоудаление окажется неудачным, вы несете ответственность за то, чтобы файл был удален в конце обучения, чтобы предотвратить повторное использование того же файла в следующий раз. Это особенно важно, если вы планируете вызывать init_process_group() несколько раз с одним и тем же именем файла. Другими словами, если файл не удален/не очищен, и вы снова вызываете init_process_group() для этого файла, ожидаются сбои. Эмпирическое правило здесь заключается в том, чтобы убедиться, что файл не существует или пуст каждый раз при вызове init_process_group(). import torch.distributed as dist # ранг всегда должен быть указан dist.init_process_group(backend, init_method='file:///mnt/nfs/sharedfile', world_size=4, rank=args.rank) Инициализация через переменные окружения# Этот метод будет читать конфигурацию из переменных окружения, что позволяет полностью настроить способ получения информации. Устанавливаемые переменные: MASTER_PORT — обязательно; должен быть свободным портом на машине с рангом 0 MASTER_ADDR — обязательно (кроме ранга 0); адрес узла с рангом 0 WORLD_SIZE — обязательно; может быть установлен здесь или в вызове функции init RANK — обязательно; может быть установлен здесь или в вызове функции init Машина с рангом 0 будет использоваться для установки всех соединений. Это метод по умолчанию, то есть init_method не нужно указывать (или он может быть env://). Улучшение времени инициализации# TORCH_GLOO_LAZY_INIT — устанавливает соединения по требованию, а не использует полную сетку, что может значительно улучшить время инициализации для операций, отличных от all2all. Пост-инициализация# После запуска torch.distributed.init_process_group() можно использовать следующие функции. Чтобы проверить, была ли уже инициализирована группа процессов, используйте torch.distributed.is_initialized(). class torch.distributed.Backend(name)[source]# Класс, подобный перечислению, для бэкендов. Доступные бэкенды: GLOO, NCCL, UCC, MPI, XCCL и другие зарегистрированные бэкенды. Значения этого класса представляют собой строки в нижнем регистре, например "gloo". К ним можно получить доступ как к атрибутам, например Backend.NCCL. Этот класс можно вызывать напрямую для разбора строки, например Backend(backend_str) проверит, действителен ли backend_str, и вернет разобранную строку в нижнем регистре, если это так. Он также принимает строки в верхнем регистре, например Backend("GLOO") вернет "gloo". Примечание Запись Backend.UNDEFINED присутствует, но используется только как начальное значение некоторых полей. Пользователям не следует использовать ее напрямую или предполагать ее существование. classmethod register_backend(name, func, extended_api=False, devices=None)[source]# Регистрирует новый бэкенд с заданным именем и функцией создания экземпляра. Этот метод класса используется расширением ProcessGroup сторонних разработчиков для регистрации новых бэкендов. Параметры name (str) – Имя бэкенда расширения ProcessGroup. Оно должно совпадать с именем в init_process_group(). func (function) – Обработчик функции, создающий экземпляр бэкенда. Функция должна быть реализована в расширении бэкенда и принимать четыре аргумента: store, rank, world_size и timeout. extended_api (bool, опционально) – Поддерживает ли бэкенд расширенную структуру аргументов. По умолчанию: False. Если установлено True, бэкенд получит экземпляр c10d::DistributedBackendOptions и объект параметров группы процессов, определенный реализацией бэкенда. device (str или list of str, опционально) – Тип устройства, поддерживаемый этим бэкендом, например "cpu", "cuda" и т.д. Если None, предполагается поддержка как "cpu", так и "cuda". Примечание Эта поддержка стороннего бэкенда является экспериментальной и может быть изменена. torch.distributed.get_backend(group=None)[source]# Возвращает бэкенд данной группы процессов. Параметры group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. По умолчанию используется основная группа процессов. Если указана другая конкретная группа, вызывающий процесс должен быть частью group. Возвращает Бэкенд данной группы процессов в виде строки в нижнем регистре. Тип возвращаемого значения Backend torch.distributed.get_rank(group=None)[source]# Возвращает ранг текущего процесса в предоставленной группе, иначе по умолчанию. Ранг — это уникальный идентификатор, присваиваемый каждому процессу в распределенной группе процессов. Они всегда являются последовательными целыми числами от 0 до world_size. Параметры group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. Возвращает Ранг процесса или -1, если не является частью группы. Тип возвращаемого значения int torch.distributed.get_world_size(group=None)[source]# Возвращает количество процессов в текущей группе процессов. Параметры group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. Возвращает Размер мира группы процессов или -1, если не является частью группы. Тип возвращаемого значения int Завершение работы# Важно очищать ресурсы при выходе с помощью вызова destroy_process_group(). Самый простой шаблон — уничтожить каждую группу процессов и бэкенд, вызвав destroy_process_group() со значением по умолчанию None для аргумента group, в точке обучающего скрипта, где коммуникации больше не нужны, обычно ближе к концу main(). Вызов должен быть выполнен один раз на процесс-тренер, а не на внешнем уровне запуска процессов. Если destroy_process_group() не вызывается всеми рангами в pg в течение времени тайм-аута, особенно когда в приложении есть несколько групп процессов, например для N-D параллелизма, возможны зависания при выходе. Это связано с тем, что деструктор ProcessGroupNCCL вызывает ncclCommAbort, который должен вызываться коллективно, но порядок вызова деструктора ProcessGroupNCCL, если он вызывается сборщиком мусора Python, не является детерминированным. Вызов destroy_process_group() помогает, гарантируя, что ncclCommAbort вызывается в согласованном порядке между рангами, и избегает вызова ncclCommAbort во время деструктора ProcessGroupNCCL. Повторная инициализация# destroy_process_group также можно использовать для уничтожения отдельных групп процессов. Один из вариантов использования — отказоустойчивое обучение, когда группа процессов может быть уничтожена, а затем во время выполнения инициализирована новая. В этом случае крайне важно синхронизировать процессы-тренеры с помощью каких-либо средств, отличных от примитивов torch.distributed, после вызова destroy и перед последующей инициализацией. Это поведение в настоящее время не поддерживается/не тестируется из-за сложности достижения такой синхронизации и считается известной проблемой. Пожалуйста, создайте issue на GitHub или RFC, если этот вариант использования блокирует вас. Группы# По умолчанию коллективы работают с группой по умолчанию (также называемой миром) и требуют, чтобы все процессы вошли в вызов распределенной функции. Однако некоторые рабочие нагрузки могут выиграть от более детальной коммуникации. Здесь в игру вступают распределенные группы. Функция new_group() может использоваться для создания новых групп с произвольными подмножествами всех процессов. Она возвращает непрозрачный дескриптор группы, который может быть передан в качестве аргумента group всем коллективам (коллективы — это распределенные функции для обмена информацией в определенных хорошо известных шаблонах программирования). torch.distributed.new_group(ranks=None, timeout=None, backend=None, pg_options=None, use_local_synchronization=False, group_desc=None, device_id=None)[source]# Создает новую распределенную группу. Эта функция требует, чтобы все процессы в основной группе (т.е. все процессы, являющиеся частью распределенного задания) вошли в эту функцию, даже если они не будут членами группы. Кроме того, группы должны создаваться в одном и том же порядке во всех процессах. Предупреждение Безопасное параллельное использование: При использовании нескольких групп процессов с бэкендом NCCL пользователь должен обеспечить глобально согласованный порядок выполнения коллективов между рангами. Если несколько потоков в процессе выдают коллективы, необходима явная синхронизация для обеспечения согласованного порядка. При использовании асинхронных вариантов API коммуникации torch.distributed возвращается объект work, и ядро коммуникации помещается в отдельный поток CUDA, что позволяет перекрывать коммуникацию и вычисления. После того как одна или несколько асинхронных операций были выданы в одной группе процессов, они должны быть синхронизированы с другими потоками cuda путем вызова work.wait() перед использованием другой группы процессов. См. Using multiple NCCL communicators concurrently <https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/communicators.html#using-multiple-nccl-communicators-concurrently> для получения дополнительной информации. Параметры ranks (list[int]) – Список рангов членов группы. Если None, будет установлено на все ранги. По умолчанию None. timeout (timedelta, опционально) – см. init_process_group для подробностей и значения по умолчанию. backend (str или Backend, опционально) – Используемый бэкенд. В зависимости от конфигурации сборки допустимыми значениями являются gloo и nccl. По умолчанию используется тот же бэкенд, что и у глобальной группы. Это поле должно быть задано в виде строки в нижнем регистре (например, "gloo"), к которой также можно получить доступ через атрибуты Backend (например, Backend.GLOO). Если передано None, будет использоваться бэкенд, соответствующий группе процессов по умолчанию. По умолчанию None. pg_options (ProcessGroupOptions, опционально) – Параметры группы процессов, указывающие, какие дополнительные параметры необходимо передать при создании конкретных групп процессов. Например, для бэкенда nccl может быть указан is_high_priority_stream, чтобы группа процессов могла использовать потоки cuda с высоким приоритетом. Другие доступные опции для настройки nccl см. https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/api/types.html#ncclconfig-tuse_local_synchronization (bool, опционально): выполнить локальный барьер группы в конце создания группы процессов. Это отличается тем, что ранги, не являющиеся членами, не должны вызывать API и не присоединяются к барьеру. group_desc (str, опционально) – строка для описания группы процессов. device_id (torch.device, опционально) – одно конкретное устройство для «привязки» этого процесса. Вызов new_group попытается немедленно инициализировать бэкенд коммуникации для устройства, если это поле задано. Возвращает Дескриптор распределенной группы, который может быть передан коллективным вызовам или GroupMember.NON_GROUP_MEMBER, если ранг не входит в ranks. N.B. use_local_synchronization не работает с MPI. N.B. Хотя use_local_synchronization=True может быть значительно быстрее в больших кластерах и небольших группах процессов, следует соблюдать осторожность, поскольку это изменяет поведение кластера: ранги, не являющиеся членами, не присоединяются к group barrier(). N.B. use_local_synchronization=True может привести к взаимоблокировкам, когда каждый ранг создает несколько перекрывающихся групп процессов. Чтобы избежать этого, убедитесь, что все ранги следуют одному и тому же глобальному порядку создания. torch.distributed.get_group_rank(group, global_rank)[source]# Преобразует глобальный ранг в ранг группы. global_rank должен быть частью group, иначе возникает RuntimeError. Параметры group (ProcessGroup) – ProcessGroup, в которой нужно найти относительный ранг. global_rank (int) – Глобальный ранг для запроса. Возвращает Ранг группы для global_rank относительно group. Тип возвращаемого значения int N.B. вызов этой функции для группы процессов по умолчанию возвращает идентичность. torch.distributed.get_global_rank(group, group_rank)[source]# Преобразует ранг группы в глобальный ранг. group_rank должен быть частью group, иначе возникает RuntimeError. Параметры group (ProcessGroup) – ProcessGroup, из которой нужно получить глобальный ранг. group_rank (int) – Ранг группы для запроса. Возвращает Глобальный ранг group_rank относительно group. Тип возвращаемого значения int N.B. вызов этой функции для группы процессов по умолчанию возвращает идентичность. torch.distributed.get_process_group_ranks(group)[source]# Получить все ранги, связанные с group. Параметры group (Optional[ProcessGroup]) – ProcessGroup, из которой нужно получить все ранги. Если None, будет использоваться группа процессов по умолчанию. Возвращает Список глобальных рангов, упорядоченных по рангу группы. Тип возвращаемого значения list[int] DeviceMesh# DeviceMesh — это абстракция более высокого уровня, которая управляет группами процессов (или коммуникаторами NCCL). Она позволяет пользователю легко создавать меж- и внутриузловые группы процессов, не беспокоясь о том, как правильно настроить ранги для различных подгрупп процессов, и помогает легко управлять этими распределенными группами процессов. Функция init_device_mesh() может использоваться для создания нового DeviceMesh с формой сетки, описывающей топологию устройств. class torch.distributed.device_mesh.DeviceMesh(device_type, mesh, , mesh_dim_names=None, backend_override=None, _init_backend=True)[source]# DeviceMesh представляет собой сетку устройств, где расположение устройств может быть представлено в виде n-мерного массива, и каждое значение n-мерного массива является глобальным идентификатором рангов группы процессов по умолчанию. DeviceMesh может использоваться для настройки N-мерных соединений устройств в кластере и управления ProcessGroups для N-мерного параллелизма. Коммуникации могут происходить по каждому измерению DeviceMesh отдельно. DeviceMesh уважает устройство, которое пользователь уже выбрал (т.е. если пользователь вызвал torch.cuda.set_device до инициализации DeviceMesh), и выберет/установит устройство для текущего процесса, если пользователь не установил устройство заранее. Обратите внимание, что ручной выбор устройства должен происходить ДО инициализации DeviceMesh. DeviceMesh также может использоваться в качестве контекстного менеджера при совместном использовании с API DTensor. Примечание DeviceMesh следует модели программирования SPMD, что означает, что одна и та же программа Python PyTorch выполняется на всех процессах/рангах в кластере. Поэтому пользователи должны убедиться, что массив mesh (который описывает расположение устройств) идентичен во всех рангах. Несогласованная сетка приведет к молчаливому зависанию. Параметры device_type (str) – Тип устройства сетки. В настоящее время поддерживается: "cpu", "cuda/cuda-like". mesh (ndarray) – Многомерный массив или целочисленный тензор, описывающий расположение устройств, где идентификаторы являются глобальными идентификаторами группы процессов по умолчанию. Возвращает Объект DeviceMesh, представляющий расположение устройств. Тип возвращаемого значения DeviceMesh Следующая программа выполняется на каждом процессе/ранге в манере SPMD. В этом примере у нас есть 2 хоста с 4 GPU каждый. Редукция по первому измерению сетки уменьшит по столбцам (0, 4),.. и (3, 7), редукция по второму измерению сетки уменьшит по строкам (0, 1, 2, 3) и (4, 5, 6, 7). Пример: >>> from torch.distributed.device_mesh import DeviceMesh >>> >>> # Инициализация сетки устройств как (2, 4) для представления топологии >>> # кросс-хоста (dim 0) и внутри хоста (dim 1). >>> mesh = DeviceMesh(device_type="cuda", mesh=[[0, 1, 2, 3],[4, 5, 6, 7]]) static from_group(group, device_type, mesh=None, , mesh_dim_names=None)[source]# Создает DeviceMesh с device_type из существующей ProcessGroup или списка существующих ProcessGroup. Созданная сетка устройств имеет количество измерений, равное количеству переданных групп. Например, если передана одна группа процессов, результирующий DeviceMesh является 1D сеткой. Если передан список из 2 групп процессов, результирующий DeviceMesh является 2D сеткой. Если передано более одной группы, то обязательны аргументы mesh и mesh_dim_names. Порядок переданных групп процессов определяет топологию сетки. Например, первая группа процессов будет 0-м измерением DeviceMesh. Переданный тензор mesh должен иметь то же количество измерений, что и количество переданных групп процессов, и порядок измерений в тензоре mesh должен соответствовать порядку в переданных группах процессов. Параметры group (ProcessGroup or list[ProcessGroup]) – существующая ProcessGroup или список существующих ProcessGroup. device_type (str) – Тип устройства сетки. В настоящее время поддерживается: "cpu", "cuda/cuda-like". Передача типа устройства с индексом GPU, например "cuda:0", не допускается. mesh (torch.Tensor or ArrayLike, опционально) – Многомерный массив или целочисленный тензор, описывающий расположение устройств, где идентификаторы являются глобальными идентификаторами группы процессов по умолчанию. По умолчанию None. mesh_dim_names (tuple[str], опционально) – Кортеж имен измерений сетки для назначения каждому измерению многомерного массива, описывающего расположение устройств. Его длина должна совпадать с длиной mesh_shape. Каждая строка в mesh_dim_names должна быть уникальной. По умолчанию None. Возвращает Объект DeviceMesh, представляющий расположение устройств. Тип возвращаемого значения DeviceMesh get_all_groups()[source]# Возвращает список ProcessGroups для всех измерений сетки. Возвращает Список объектов ProcessGroup. Тип возвращаемого значения list[torch.distributed.distributed_c10d.ProcessGroup] get_coordinate()[source]# Возвращает относительные индексы этого ранга относительно всех измерений сетки. Если этот ранг не является частью сетки, возвращает None. Тип возвращаемого значения Optional[list[int]] get_group(mesh_dim=None)[source]# Возвращает единственную ProcessGroup, указанную mesh_dim, или, если mesh_dim не указан и DeviceMesh является 1-мерным, возвращает единственную ProcessGroup в сетке. Параметры mesh_dim (str/python:int, опционально) – может быть именем измерения сетки или индексом None. (измерения сетки. По умолчанию) – Возвращает Объект ProcessGroup. Тип возвращаемого значения ProcessGroup get_local_rank(mesh_dim=None)[source]# Возвращает локальный ранг данного mesh_dim DeviceMesh. Параметры mesh_dim (str/python:int, опционально) – может быть именем измерения сетки или индексом None. (измерения сетки. По умолчанию) – Возвращает Целое число, обозначающее локальный ранг. Тип возвращаемого значения int Следующая программа выполняется на каждом процессе/ранге в манере SPMD. В этом примере у нас есть 2 хоста с 4 GPU каждый. Вызов mesh_2d.get_local_rank(mesh_dim=0) на рангах 0, 1, 2, 3 вернет 0. Вызов mesh_2d.get_local_rank(mesh_dim=0) на рангах 4, 5, 6, 7 вернет 1. Вызов mesh_2d.get_local_rank(mesh_dim=1) на рангах 0, 4 вернет 0. Вызов mesh_2d.get_local_rank(mesh_dim=1) на рангах 1, 5 вернет 1. Вызов mesh_2d.get_local_rank(mesh_dim=1) на рангах 2, 6 вернет 2. Вызов mesh_2d.get_local_rank(mesh_dim=1) на рангах 3, 7 вернет 3. Пример: >>> from torch.distributed.device_mesh import DeviceMesh >>> >>> # Инициализация сетки устройств как (2, 4) для представления топологии >>> # кросс-хоста (dim 0), и внутри хоста (dim 1). >>> mesh = DeviceMesh(device_type="cuda", mesh=[[0, 1, 2, 3],[4, 5, 6, 7]]) get_rank()[source]# Возвращает текущий глобальный ранг. Тип возвращаемого значения int Точка-точка коммуникация# torch.distributed.send(tensor, dst=None, group=None, tag=0, group_dst=None)[source]# Отправляет тензор синхронно. Предупреждение tag не поддерживается бэкендом NCCL. Параметры tensor (Tensor) – Тензор для отправки. dst (int) – Целевой ранг в глобальной группе процессов (независимо от аргумента group). Целевой ранг не должен совпадать с рангом текущего процесса. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. tag (int, опционально) – Тег для сопоставления send с удаленным recv. group_dst (int, опционально) – Целевой ранг в group. Недопустимо указывать и dst, и group_dst. torch.distributed.recv(tensor, src=None, group=None, tag=0, group_src=None)[source]# Получает тензор синхронно. Предупреждение tag не поддерживается бэкендом NCCL. Параметры tensor (Tensor) – Тензор для заполнения полученными данными. src (int, опционально) – Исходный ранг в глобальной группе процессов (независимо от аргумента group). Если не указано, будет получать от любого процесса. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. tag (int, опционально) – Тег для сопоставления recv с удаленным send. group_src (int, опционально) – Целевой ранг в group. Недопустимо указывать и src, и group_src. Возвращает Ранг отправителя или -1, если не является частью группы. Тип возвращаемого значения int isend() и irecv() возвращают распределенные объекты запросов при использовании. В общем, тип этого объекта не указан, так как они никогда не должны создаваться вручную, но гарантируется, что они поддерживают два метода: is_completed() — возвращает True, если операция завершена. wait() — блокирует процесс до завершения операции. is_completed() гарантированно возвращает True после своего возврата. torch.distributed.isend(tensor, dst=None, group=None, tag=0, group_dst=None)[source]# Отправляет тензор асинхронно. Предупреждение Изменение тензора до завершения запроса приводит к неопределенному поведению. Предупреждение tag не поддерживается бэкендом NCCL. В отличие от send, который является блокирующим, isend допускает src == dst rank, т.е. отправку самому себе. Параметры tensor (Tensor) – Тензор для отправки. dst (int) – Целевой ранг в глобальной группе процессов (независимо от аргумента group). group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. tag (int, опционально) – Тег для сопоставления send с удаленным recv. group_dst (int, опционально) – Целевой ранг в group. Недопустимо указывать и dst, и group_dst. Возвращает Объект распределенного запроса. None, если не является частью группы. Тип возвращаемого значения Optional[Work] torch.distributed.irecv(tensor, src=None, group=None, tag=0, group_src=None)[source]# Получает тензор асинхронно. Предупреждение tag не поддерживается бэкендом NCCL. В отличие от recv, который является блокирующим, irecv допускает src == dst rank, т.е. получение от самого себя. Параметры tensor (Tensor) – Тензор для заполнения полученными данными. src (int, опционально) – Исходный ранг в глобальной группе процессов (независимо от аргумента group). Если не указано, будет получать от любого процесса. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. tag (int, опционально) – Тег для сопоставления recv с удаленным send. group_src (int, опционально) – Целевой ранг в group. Недопустимо указывать и src, и group_src. Возвращает Объект распределенного запроса. None, если не является частью группы. Тип возвращаемого значения Optional[Work] torch.distributed.send_object_list(object_list, dst=None, group=None, device=None, group_dst=None, use_batch=False)[source]# Отправляет picklable объекты в object_list синхронно. Аналогично send(), но можно передавать объекты Python. Обратите внимание, что все объекты в object_list должны быть picklable для отправки. Параметры object_list (List[Any]) – Список входных объектов для отправки. Каждый объект должен быть picklable. Получатель должен предоставить списки равных размеров. dst (int) – Целевой ранг для отправки object_list. Целевой ранг основан на глобальной группе процессов (независимо от аргумента group). group (Optional[ProcessGroup]) – (ProcessGroup, опционально): Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. По умолчанию None. device (torch.device, опционально) – Если не None, объекты сериализуются и преобразуются в тензоры, которые перемещаются на устройство перед отправкой. По умолчанию None. group_dst (int, опционально) – Целевой ранг в group. Должен быть указан один из dst и group_dst, но не оба. use_batch (bool, опционально) – Если True, использовать пакетные p2p операции вместо обычных операций send. Это позволяет избежать инициализации 2-ранговых коммуникаторов и использует существующие коммуникаторы всей группы. См. batch_isend_irecv для использования и предположений. По умолчанию False. Возвращает None. Примечание Для групп процессов на основе NCCL внутренние тензорные представления объектов должны быть перемещены на устройство GPU до начала коммуникации. В этом случае используемое устройство задается torch.cuda.current_device(), и пользователь несет ответственность за то, чтобы оно было установлено так, чтобы каждый ранг имел отдельный GPU, с помощью torch.cuda.set_device(). Предупреждение Объектные коллективы имеют ряд серьезных ограничений по производительности и масштабируемости. См. Object collectives для подробностей. Предупреждение send_object_list() неявно использует модуль pickle, который, как известно, небезопасен. Можно создать вредоносные данные pickle, которые будут выполнять произвольный код при распаковке. Вызывайте эту функцию только с доверенными данными. Предупреждение Вызов send_object_list() с тензорами GPU плохо поддерживается и неэффективен, так как требует передачи GPU -> CPU, поскольку тензоры будут упакованы pickle. Вместо этого рассмотрите возможность использования send(). Пример::>>> # Примечание: Инициализация группы процессов опущена на каждом ранге. >>> import torch.distributed as dist >>> # Предполагается, что бэкенд не NCCL >>> device = torch.device("cpu") >>> if dist.get_rank() == 0: >>> # Предполагается world_size = 2. >>> objects = ["foo", 12, {1: 2}] # любой picklable объект >>> dist.send_object_list(objects, dst=1, device=device) >>> else: >>> objects = [None, None, None] >>> dist.recv_object_list(objects, src=0, device=device) >>> objects ['foo', 12, {1: 2}] torch.distributed.recv_object_list(object_list, src=None, group=None, device=None, group_src=None, use_batch=False)[source]# Получает picklable объекты в object_list синхронно. Аналогично recv(), но может получать объекты Python. Параметры object_list (List[Any]) – Список объектов для получения. Должен предоставлять список размеров, равный размеру отправляемого списка. src (int, опционально) – Исходный ранг, от которого получать object_list. Исходный ранг основан на глобальной группе процессов (независимо от аргумента group). Если установлено None, будет получать от любого ранга. По умолчанию None. group (Optional[ProcessGroup]) – (ProcessGroup, опционально): Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. По умолчанию None. device (torch.device, опционально) – Если не None, получает на этом устройстве. По умолчанию None. group_src (int, опционально) – Целевой ранг в group. Недопустимо указывать и src, и group_src. use_batch (bool, опционально) – Если True, использовать пакетные p2p операции вместо обычных операций send. Это позволяет избежать инициализации 2-ранговых коммуникаторов и использует существующие коммуникаторы всей группы. См. batch_isend_irecv для использования и предположений. По умолчанию False. Возвращает Ранг отправителя. -1, если ранг не является частью группы. Если ранг является частью группы, object_list будет содержать отправленные объекты от src ранга. Примечание Для групп процессов на основе NCCL внутренние тензорные представления объектов должны быть перемещены на устройство GPU до начала коммуникации. В этом случае используемое устройство задается torch.cuda.current_device(), и пользователь несет ответственность за то, чтобы оно было установлено так, чтобы каждый ранг имел отдельный GPU, с помощью torch.cuda.set_device(). Предупреждение Объектные коллективы имеют ряд серьезных ограничений по производительности и масштабируемости. См. Object collectives для подробностей. Предупреждение recv_object_list() неявно использует модуль pickle, который, как известно, небезопасен. Можно создать вредоносные данные pickle, которые будут выполнять произвольный код при распаковке. Вызывайте эту функцию только с доверенными данными. Предупреждение Вызов recv_object_list() с тензорами GPU плохо поддерживается и неэффективен, так как требует передачи GPU -> CPU, поскольку тензоры будут упакованы pickle. Вместо этого рассмотрите возможность использования recv(). Пример::>>> # Примечание: Инициализация группы процессов опущена на каждом ранге. >>> import torch.distributed as dist >>> # Предполагается, что бэкенд не NCCL >>> device = torch.device("cpu") >>> if dist.get_rank() == 0: >>> # Предполагается world_size = 2. >>> objects = ["foo", 12, {1: 2}] # любой picklable объект >>> dist.send_object_list(objects, dst=1, device=device) >>> else: >>> objects = [None, None, None] >>> dist.recv_object_list(objects, src=0, device=device) >>> objects ['foo', 12, {1: 2}] torch.distributed.batch_isend_irecv(p2p_op_list)[source]# Отправляет или получает пакет тензоров асинхронно и возвращает список запросов. Обрабатывает каждую из операций в p2p_op_list и возвращает соответствующие запросы. В настоящее время поддерживаются бэкенды NCCL, Gloo и UCC. Параметры p2p_op_list (list[torch.distributed.distributed_c10d.P2POp]) – Список операций точка-точка (тип каждого оператора — torch.distributed.P2POp). Порядок isend/irecv в списке важен и должен соответствовать соответствующим isend/irecv на удаленном конце. Возвращает Список объектов распределенного запроса, возвращенных вызовом соответствующей операции в op_list. Тип возвращаемого значения list[torch.distributed.distributed_c10d.Work] Примеры >>> send_tensor = torch.arange(2, dtype=torch.float32) + 2 * rank >>> recv_tensor = torch.randn(2, dtype=torch.float32) >>> send_op = dist.P2POp(dist.isend, send_tensor, (rank + 1) % world_size) >>> recv_op = dist.P2POp(... dist.irecv, recv_tensor, (rank - 1 + world_size) % world_size... ) >>> reqs = batch_isend_irecv([send_op, recv_op]) >>> for req in reqs: >>> req.wait() >>> recv_tensor tensor([2, 3]) # Ранг 0 tensor([0, 1]) # Ранг 1 Примечание Обратите внимание, что при использовании этого API с бэкендом NCCL PG пользователи должны установить текущее устройство GPU с помощью torch.cuda.set_device, иначе это приведет к неожиданным зависаниям. Кроме того, если этот API является первым коллективным вызовом в группе, переданной в dist.P2POp, все ранги группы должны участвовать в этом вызове API; в противном случае поведение не определено. Если этот вызов API не является первым коллективным вызовом в группе, разрешены пакетные P2P операции, включающие только подмножество рангов группы. class torch.distributed.P2POp(op, tensor, peer=None, group=None, tag=0, group_peer=None)[source]# Класс для создания операций точка-точка для batch_isend_irecv. Этот класс создает тип операции P2P, буфер коммуникации, ранг пира, группу процессов и тег. Экземпляры этого класса будут переданы в batch_isend_irecv для коммуникаций точка-точка. Параметры op (Callable) – Функция для отправки данных или получения данных от пирового процесса. Тип op — torch.distributed.isend или torch.distributed.irecv. tensor (Tensor) – Тензор для отправки или получения. peer (int, опционально) – Целевой или исходный ранг. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. tag (int, опционально) – Тег для сопоставления send с recv. group_peer (int, опционально) – Целевой или исходный ранг. Синхронные и асинхронные коллективные операции# Каждая функция коллективной операции поддерживает следующие два вида операций, в зависимости от настройки флага async_op, передаваемого в коллектив: Синхронная операция — режим по умолчанию, когда async_op установлен в False. Когда функция возвращается, гарантируется, что коллективная операция выполнена. В случае операций CUDA не гарантируется, что операция CUDA завершена, поскольку операции CUDA являются асинхронными. Для коллективов CPU любые последующие вызовы функций, использующие выходные данные коллективного вызова, будут вести себя ожидаемо. Для коллективов CUDA вызовы функций, использующие выходные данные в том же потоке CUDA, будут вести себя ожидаемо. Пользователи должны заботиться о синхронизации в сценарии работы с разными потоками. Подробнее о семантике CUDA, такой как синхронизация потоков, см. CUDA Semantics. См. скрипт ниже для примеров различий в этой семантике для операций CPU и CUDA. Асинхронная операция — когда async_op установлен в True. Функция коллективной операции возвращает объект распределенного запроса. В общем, вам не нужно создавать его вручную, и гарантируется, что он поддерживает два метода: is_completed() — в случае коллективов CPU возвращает True, если завершено. В случае операций CUDA возвращает True, если операция была успешно поставлена в очередь в поток CUDA и выходные данные могут быть использованы в потоке по умолчанию без дополнительной синхронизации. wait() — в случае коллективов CPU будет блокировать процесс до завершения операции. В случае коллективов CUDA будет блокировать текущий активный поток CUDA до завершения операции (но не будет блокировать CPU). get_future() — возвращает объект torch.C.Future. Поддерживается для NCCL, а также для большинства операций на GLOO и MPI, за исключением операций точка-точка. Примечание: по мере того, как мы продолжаем внедрять Futures и объединять API, вызов get_future() может стать избыточным. Пример Следующий код может служить справочным материалом по семантике для операций CUDA при использовании распределенных коллективов. Он показывает явную необходимость синхронизации при использовании выходных данных коллектива в разных потоках CUDA: # Код выполняется на каждом ранге. 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 гарантирует, что операция поставлена в очередь, но не обязательно завершена. handle.wait() # Использование результата в нестандартном потоке. with torch.cuda.stream(s): s.wait_stream(torch.cuda.default_stream()) output.add(100) if rank == 0: # если явный вызов wait_stream был опущен, вывод ниже будет # недетерминированно 1 или 101, в зависимости от того, перезаписал ли allreduce # значение после завершения add. print(output) Коллективные функции# torch.distributed.broadcast(tensor, src=None, group=None, async_op=False, group_src=None)[source]# Рассылает тензор всей группе. tensor должен иметь одинаковое количество элементов во всех процессах, участвующих в коллективе. Параметры tensor (Tensor) – Данные для отправки, если src является рангом текущего процесса, и тензор для сохранения полученных данных в противном случае. src (int) – Исходный ранг в глобальной группе процессов (независимо от аргумента group). group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. group_src (int) – Исходный ранг в group. Должен быть указан один из group_src и src, но не оба. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. torch.distributed.broadcast_object_list(object_list, src=None, group=None, device=None, group_src=None)[source]# Рассылает picklable объекты в object_list всей группе. Аналогично broadcast(), но можно передавать объекты Python. Обратите внимание, что все объекты в object_list должны быть picklable для рассылки. Параметры object_list (List[Any]) – Список входных объектов для рассылки. Каждый объект должен быть picklable. Только объекты на src ранге будут разосланы, но каждый ранг должен предоставить списки равных размеров. src (int) – Исходный ранг, от которого рассылать object_list. Исходный ранг основан на глобальной группе процессов (независимо от аргумента group). group (Optional[ProcessGroup]) – (ProcessGroup, опционально): Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. По умолчанию None. device (torch.device, опционально) – Если не None, объекты сериализуются и преобразуются в тензоры, которые перемещаются на устройство перед рассылкой. По умолчанию None. group_src (int) – Исходный ранг в group. Нельзя указывать один из group_src и src, но не оба. Возвращает None. Если ранг является частью группы, object_list будет содержать разосланные объекты от src ранга. Примечание Для групп процессов на основе NCCL внутренние тензорные представления объектов должны быть перемещены на устройство GPU до начала коммуникации. В этом случае используемое устройство задается torch.cuda.current_device(), и пользователь несет ответственность за то, чтобы оно было установлено так, чтобы каждый ранг имел отдельный GPU, с помощью torch.cuda.set_device(). Примечание Обратите внимание, что этот API немного отличается от коллектива broadcast(), поскольку он не предоставляет дескриптор async_op и, следовательно, будет блокирующим вызовом. Предупреждение Объектные коллективы имеют ряд серьезных ограничений по производительности и масштабируемости. См. Object collectives для подробностей. Предупреждение broadcast_object_list() неявно использует модуль pickle, который, как известно, небезопасен. Можно создать вредоносные данные pickle, которые будут выполнять произвольный код при распаковке. Вызывайте эту функцию только с доверенными данными. Предупреждение Вызов broadcast_object_list() с тензорами GPU плохо поддерживается и неэффективен, так как требует передачи GPU -> CPU, поскольку тензоры будут упакованы pickle. Вместо этого рассмотрите возможность использования broadcast(). Пример::>>> # Примечание: Инициализация группы процессов опущена на каждом ранге. >>> import torch.distributed as dist >>> if dist.get_rank() == 0: >>> # Предполагается world_size = 3. >>> objects = ["foo", 12, {1: 2}] # любой picklable объект >>> else: >>> objects = [None, None, None] >>> # Предполагается, что бэкенд не NCCL >>> device = torch.device("cpu") >>> dist.broadcast_object_list(objects, src=0, device=device) >>> objects ['foo', 12, {1: 2}] torch.distributed.all_reduce(tensor, op=<RedOpType.SUM: 0>, group=None, async_op=False)[source]# Уменьшает данные тензора на всех машинах таким образом, что все получают окончательный результат. После вызова тензор будет побитово идентичен во всех процессах. Поддерживаются комплексные тензоры. Параметры tensor (Tensor) – Вход и выход коллектива. Функция работает на месте. op (опционально) – Одно из значений из перечисления torch.distributed.ReduceOp. Указывает операцию, используемую для поэлементных редукций. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Примеры >>> # Все тензоры ниже имеют тип torch.int64. >>> # У нас 2 группы процессов, 2 ранга. >>> device = torch.device(f"cuda:{rank}") >>> tensor = torch.arange(2, dtype=torch.int64, device=device) + 1 + 2 * rank >>> tensor tensor([1, 2], device='cuda:0') # Ранг 0 tensor([3, 4], device='cuda:1') # Ранг 1 >>> dist.all_reduce(tensor, op=ReduceOp.SUM) >>> tensor tensor([4, 6], device='cuda:0') # Ранг 0 tensor([4, 6], device='cuda:1') # Ранг 1 >>> # Все тензоры ниже имеют тип torch.cfloat. >>> # У нас 2 группы процессов, 2 ранга. >>> tensor = torch.tensor(... [1 + 1j, 2 + 2j], dtype=torch.cfloat, device=device... ) + 2 * rank * (1 + 1j) >>> tensor tensor([1.+1.j, 2.+2.j], device='cuda:0') # Ранг 0 tensor([3.+3.j, 4.+4.j], device='cuda:1') # Ранг 1 >>> dist.all_reduce(tensor, op=ReduceOp.SUM) >>> tensor tensor([4.+4.j, 6.+6.j], device='cuda:0') # Ранг 0 tensor([4.+4.j, 6.+6.j], device='cuda:1') # Ранг 1 torch.distributed.reduce(tensor, dst=None, op=<RedOpType.SUM: 0>, group=None, async_op=False, group_dst=None)[source]# Уменьшает данные тензора на всех машинах. Только процесс с рангом dst получит окончательный результат. Параметры tensor (Tensor) – Вход и выход коллектива. Функция работает на месте. dst (int) – Целевой ранг в глобальной группе процессов (независимо от аргумента group). op (опционально) – Одно из значений из перечисления torch.distributed.ReduceOp. Указывает операцию, используемую для поэлементных редукций. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. group_dst (int) – Целевой ранг в group. Должен быть указан один из group_dst и dst, но не оба. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. torch.distributed.all_gather(tensor_list, tensor, group=None, async_op=False)[source]# Собирает тензоры от всей группы в список. Поддерживаются комплексные тензоры и тензоры неравного размера. Параметры tensor_list (list[Tensor]) – Выходной список. Должен содержать тензоры правильного размера для использования в качестве вывода коллектива. Поддерживаются тензоры неравного размера. tensor (Tensor) – Тензор для рассылки от текущего процесса. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Примеры >>> # Все тензоры ниже имеют dtype torch.int64. >>> # У нас 2 группы процессов, 2 ранга. >>> device = torch.device(f"cuda:{rank}") >>> tensor_list = [... torch.zeros(2, dtype=torch.int64, device=device) for _ in range(2)... ] >>> tensor_list [tensor([0, 0], device='cuda:0'), tensor([0, 0], device='cuda:0')] # Ранг 0 [tensor([0, 0], device='cuda:1'), tensor([0, 0], device='cuda:1')] # Ранг 1 >>> tensor = torch.arange(2, dtype=torch.int64, device=device) + 1 + 2 * rank >>> tensor tensor([1, 2], device='cuda:0') # Ранг 0 tensor([3, 4], device='cuda:1') # Ранг 1 >>> dist.all_gather(tensor_list, tensor) >>> tensor_list [tensor([1, 2], device='cuda:0'), tensor([3, 4], device='cuda:0')] # Ранг 0 [tensor([1, 2], device='cuda:1'), tensor([3, 4], device='cuda:1')] # Ранг 1 >>> # Все тензоры ниже имеют dtype torch.cfloat. >>> # У нас 2 группы процессов, 2 ранга. >>> tensor_list = [... torch.zeros(2, dtype=torch.cfloat, device=device) for _ in range(2)... ] >>> tensor_list [tensor([0.+0.j, 0.+0.j], device='cuda:0'), tensor([0.+0.j, 0.+0.j], device='cuda:0')] # Ранг 0 [tensor([0.+0.j, 0.+0.j], device='cuda:1'), tensor([0.+0.j, 0.+0.j], device='cuda:1')] # Ранг 1 >>> tensor = torch.tensor(... [1 + 1j, 2 + 2j], dtype=torch.cfloat, device=device... ) + 2 * rank * (1 + 1j) >>> tensor tensor([1.+1.j, 2.+2.j], device='cuda:0') # Ранг 0 tensor([3.+3.j, 4.+4.j], device='cuda:1') # Ранг 1 >>> dist.all_gather(tensor_list, tensor) >>> tensor_list [tensor([1.+1.j, 2.+2.j], device='cuda:0'), tensor([3.+3.j, 4.+4.j], device='cuda:0')] # Ранг 0 [tensor([1.+1.j, 2.+2.j], device='cuda:1'), tensor([3.+3.j, 4.+4.j], device='cuda:1')] # Ранг 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, опционально) – Должна ли эта операция быть асинхронной. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Примеры >>> # Все тензоры ниже имеют dtype torch.int64 и находятся на устройствах CUDA. >>> # У нас два ранга. >>> 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') # Ранг 0 tensor([3, 4], device='cuda:1') # Ранг 1 >>> # Вывод в форме конкатенации >>> 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') # Ранг 0 tensor([1, 2, 3, 4], device='cuda:1') # Ранг 1 >>> # Вывод в форме стека >>> 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') # Ранг 0 tensor([[1, 2], [3, 4]], device='cuda:1') # Ранг 1 torch.distributed.all_gather_object(object_list, obj, group=None)[source]# Собирает picklable объекты от всей группы в список. Аналогично all_gather(), но можно передавать объекты Python. Обратите внимание, что объект должен быть picklable для сбора. Параметры object_list (list[Any]) – Выходной список. Он должен быть правильно размера как размер группы для этого коллектива и будет содержать вывод. obj (Any) – Pickable объект 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(). Предупреждение Объектные коллективы имеют ряд серьезных ограничений по производительности и масштабируемости. См. Object collectives для подробностей. Предупреждение all_gather_object() неявно использует модуль pickle, который, как известно, небезопасен. Можно создать вредоносные данные pickle, которые будут выполнять произвольный код при распаковке. Вызывайте эту функцию только с доверенными данными. Предупреждение Вызов all_gather_object() с тензорами GPU плохо поддерживается и неэффективен, так как требует передачи GPU -> CPU, поскольку тензоры будут упакованы pickle. Вместо этого рассмотрите возможность использования all_gather(). Пример::>>> # Примечание: Инициализация группы процессов опущена на каждом ранге. >>> import torch.distributed as dist >>> # Предполагается world_size = 3. >>> gather_objects = ["foo", 12, {1: 2}] # любой picklable объект >>> 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=None, group=None, async_op=False, group_dst=None)[source]# Собирает список тензоров в одном процессе. Эта функция требует, чтобы все тензоры были одинакового размера на каждом процессе. Параметры tensor (Tensor) – Входной тензор. gather_list (list[Tensor], опционально) – Список соответствующим образом размещенных тензоров одинакового размера для использования в качестве собранных данных (по умолчанию None, должен быть указан на целевом ранге). dst (int, опционально) – Целевой ранг в глобальной группе процессов (независимо от аргумента group). (Если и dst, и group_dst равны None, по умолчанию используется глобальный ранг 0). group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. group_dst (int, опционально) – Целевой ранг в group. Недопустимо указывать и dst, и group_dst. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Примечание Обратите внимание, что все тензоры в gather_list должны быть одинакового размера. Пример::>>> # У нас 2 группы процессов, 2 ранга. >>> tensor_size = 2 >>> device = torch.device(f'cuda:{rank}') >>> tensor = torch.ones(tensor_size, device=device) + rank >>> if dist.get_rank() == 0: >>> gather_list = [torch.zeros_like(tensor, device=device) for i in range(2)] >>> else: >>> gather_list = None >>> dist.gather(tensor, gather_list, dst=0) >>> # Ранг 0 получает собранные данные. >>> gather_list [tensor([1., 1.], device='cuda:0'), tensor([2., 2.], device='cuda:0')] # Ранг 0 None # Ранг 1 torch.distributed.gather_object(obj, object_gather_list=None, dst=None, group=None, group_dst=None)[source]# Собирает picklable объекты от всей группы в одном процессе. Аналогично gather(), но можно передавать объекты Python. Обратите внимание, что объект должен быть picklable для сбора. Параметры obj (Any) – Входной объект. Должен быть picklable. object_gather_list (list[Any]) – Выходной список. На dst ранге он должен быть правильно размера как размер группы для этого коллектива и будет содержать вывод. Должен быть None на рангах, не являющихся dst. (по умолчанию None). dst (int, опционально) – Целевой ранг в глобальной группе процессов (независимо от аргумента group). (Если и dst, и group_dst равны None, по умолчанию используется глобальный ранг 0). group (Optional[ProcessGroup]) – (ProcessGroup, опционально): Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. По умолчанию None. group_dst (int, опционально) – Целевой ранг в group. Недопустимо указывать и dst, и group_dst. Возвращает None. На dst ранге object_gather_list будет содержать вывод коллектива. Примечание Обратите внимание, что этот API немного отличается от коллектива gather, поскольку он не предоставляет дескриптор async_op и, следовательно, будет блокирующим вызовом. Примечание Для групп процессов на основе NCCL внутренние тензорные представления объектов должны быть перемещены на устройство GPU до начала коммуникации. В этом случае используемое устройство задается torch.cuda.current_device(), и пользователь несет ответственность за то, чтобы оно было установлено так, чтобы каждый ранг имел отдельный GPU, с помощью torch.cuda.set_device(). Предупреждение Объектные коллективы имеют ряд серьезных ограничений по производительности и масштабируемости. См. Object collectives для подробностей. Предупреждение gather_object() неявно использует модуль pickle, который, как известно, небезопасен. Можно создать вредоносные данные pickle, которые будут выполнять произвольный код при распаковке. Вызывайте эту функцию только с доверенными данными. Предупреждение Вызов gather_object() с тензорами GPU плохо поддерживается и неэффективен, так как требует передачи GPU -> CPU, поскольку тензоры будут упакованы pickle. Вместо этого рассмотрите возможность использования gather(). Пример::>>> # Примечание: Инициализация группы процессов опущена на каждом ранге. >>> import torch.distributed as dist >>> # Предполагается world_size = 3. >>> gather_objects = ["foo", 12, {1: 2}] # любой picklable объект >>> output = [None for _ in gather_objects] >>> dist.gather_object(... gather_objects[dist.get_rank()],... output if dist.get_rank() == 0 else None,... dst=0... ) >>> # На ранге 0 >>> output ['foo', 12, {1: 2}] torch.distributed.scatter(tensor, scatter_list=None, src=None, group=None, async_op=False, group_src=None)[source]# Разбрасывает список тензоров всем процессам в группе. Каждый процесс получит ровно один тензор и сохранит его данные в аргументе tensor. Поддерживаются комплексные тензоры. Параметры tensor (Tensor) – Выходной тензор. scatter_list (list[Tensor]) – Список тензоров для разбрасывания (по умолчанию None, должен быть указан на исходном ранге). src (int) – Исходный ранг в глобальной группе процессов (независимо от аргумента group). (Если и src, и group_src равны None, по умолчанию используется глобальный ранг 0). group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. group_src (int, опционально) – Исходный ранг в group. Недопустимо указывать и src, и group_src. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Примечание Обратите внимание, что все тензоры в scatter_list должны быть одинакового размера. Пример::>>> # Примечание: Инициализация группы процессов опущена на каждом ранге. >>> import torch.distributed as dist >>> tensor_size = 2 >>> device = torch.device(f'cuda:{rank}') >>> output_tensor = torch.zeros(tensor_size, device=device) >>> if dist.get_rank() == 0: >>> # Предполагается world_size = 2. >>> # Только тензоры, все они должны быть одинакового размера. >>> t_ones = torch.ones(tensor_size, device=device) >>> t_fives = torch.ones(tensor_size, device=device) * 5 >>> scatter_list = [t_ones, t_fives] >>> else: >>> scatter_list = None >>> dist.scatter(output_tensor, scatter_list, src=0) >>> # Ранг i получает scatter_list[i]. >>> output_tensor tensor([1., 1.], device='cuda:0') # Ранг 0 tensor([5., 5.], device='cuda:1') # Ранг 1 torch.distributed.scatter_object_list(scatter_object_output_list, scatter_object_input_list=None, src=None, group=None, group_src=None)[source]# Разбрасывает picklable объекты в scatter_object_input_list всей группе. Аналогично scatter(), но можно передавать объекты Python. На каждом ранге разбросанный объект будет сохранен как первый элемент scatter_object_output_list. Обратите внимание, что все объекты в scatter_object_input_list должны быть picklable для разбрасывания. Параметры scatter_object_output_list (List[Any]) – Непустой список, первый элемент которого будет хранить объект, разбросанный на этот ранг. scatter_object_input_list (List[Any], опционально) – Список входных объектов для разбрасывания. Каждый объект должен быть picklable. Только объекты на src ранге будут разбросаны, и аргумент может быть None для рангов, не являющихся src. src (int) – Исходный ранг, от которого разбрасывать scatter_object_input_list. Исходный ранг основан на глобальной группе процессов (независимо от аргумента group). (Если и src, и group_src равны None, по умолчанию используется глобальный ранг 0). group (Optional[ProcessGroup]) – (ProcessGroup, опционально): Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. По умолчанию None. group_src (int, опционально) – Исходный ранг в group. Недопустимо указывать и src, и group_src. Возвращает None. Если ранг является частью группы, scatter_object_output_list будет иметь свой первый элемент, установленный на разбросанный объект для этого ранга. Примечание Обратите внимание, что этот API немного отличается от коллектива scatter, поскольку он не предоставляет дескриптор async_op и, следовательно, будет блокирующим вызовом. Предупреждение Объектные коллективы имеют ряд серьезных ограничений по производительности и масштабируемости. См. Object collectives для подробностей. Предупреждение scatter_object_list() неявно использует модуль pickle, который, как известно, небезопасен. Можно создать вредоносные данные pickle, которые будут выполнять произвольный код при распаковке. Вызывайте эту функцию только с доверенными данными. Предупреждение Вызов scatter_object_list() с тензорами GPU плохо поддерживается и неэффективен, так как требует передачи GPU -> CPU, поскольку тензоры будут упакованы pickle. Вместо этого рассмотрите возможность использования scatter(). Пример::>>> # Примечание: Инициализация группы процессов опущена на каждом ранге. >>> import torch.distributed as dist >>> if dist.get_rank() == 0: >>> # Предполагается world_size = 3. >>> objects = ["foo", 12, {1: 2}] # любой picklable объект >>> else: >>> # Может быть любым списком на рангах, не являющихся src, элементы не используются. >>> objects = [None, None, None] >>> output_list = [None] >>> dist.scatter_object_list(output_list, objects, src=0) >>> # Ранг i получает objects[i]. Например, на ранге 2: >>> output_list [{1: 2}] torch.distributed.reduce_scatter(output, input_list, op=<RedOpType.SUM: 0>, group=None, async_op=False)[source]# Уменьшает, а затем разбрасывает список тензоров всем процессам в группе. Параметры output (Tensor) – Выходной тензор. input_list (list[Tensor]) – Список тензоров для уменьшения и разбрасывания. op (опционально) – Одно из значений из перечисления torch.distributed.ReduceOp. Указывает операцию, используемую для поэлементных редукций. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. torch.distributed.reduce_scatter_tensor(output, input, op=<RedOpType.SUM: 0>, group=None, async_op=False)[source]# Уменьшает, а затем разбрасывает тензор всем рангам в группе. Параметры output (Tensor) – Выходной тензор. Он должен быть одинакового размера во всех рангах. input (Tensor) – Входной тензор для уменьшения и разбрасывания. Его размер должен быть равен размеру выходного тензора, умноженному на world_size. Входной тензор может иметь одну из следующих форм: (i) конкатенация выходных тензоров по основному измерению, или (ii) стек выходных тензоров по основному измерению. Для определения «конкатенации» см. torch.cat(). Для определения «стека» см. torch.stack(). group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Примеры >>> # Все тензоры ниже имеют dtype torch.int64 и находятся на устройствах CUDA. >>> # У нас два ранга. >>> device = torch.device(f"cuda:{rank}") >>> tensor_out = torch.zeros(2, dtype=torch.int64, device=device) >>> # Ввод в форме конкатенации >>> tensor_in = torch.arange(world_size * 2, dtype=torch.int64, device=device) >>> tensor_in tensor([0, 1, 2, 3], device='cuda:0') # Ранг 0 tensor([0, 1, 2, 3], device='cuda:1') # Ранг 1 >>> dist.reduce_scatter_tensor(tensor_out, tensor_in) >>> tensor_out tensor([0, 2], device='cuda:0') # Ранг 0 tensor([4, 6], device='cuda:1') # Ранг 1 >>> # Ввод в форме стека >>> tensor_in = torch.reshape(tensor_in, (world_size, 2)) >>> tensor_in tensor([[0, 1], [2, 3]], device='cuda:0') # Ранг 0 tensor([[0, 1], [2, 3]], device='cuda:1') # Ранг 1 >>> dist.reduce_scatter_tensor(tensor_out, tensor_in) >>> tensor_out tensor([0, 2], device='cuda:0') # Ранг 0 tensor([4, 6], device='cuda:1') # Ранг 1 torch.distributed.all_to_all_single(output, input, output_split_sizes=None, input_split_sizes=None, group=None, async_op=False)[source]# Разделяет входной тензор, а затем разбрасывает разделенный список всем процессам в группе. Позже полученные тензоры конкатенируются от всех процессов в группе и возвращаются как один выходной тензор. Поддерживаются комплексные тензоры. Параметры output (Tensor) – Собранный конкатенированный выходной тензор. input (Tensor) – Входной тензор для разбрасывания. output_split_sizes – (list[Int], опционально): Размеры разделения вывода для dim 0. Если указано None или пусто, dim 0 выходного тензора должен делиться поровну на world_size. input_split_sizes – (list[Int], опционально): Размеры разделения ввода для dim 0. Если указано None или пусто, dim 0 входного тензора должен делиться поровну на world_size. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Предупреждение all_to_all_single является экспериментальным и может быть изменен. Примеры >>> input = torch.arange(4) + rank * 4 >>> input tensor([0, 1, 2, 3]) # Ранг 0 tensor([4, 5, 6, 7]) # Ранг 1 tensor([8, 9, 10, 11]) # Ранг 2 tensor([12, 13, 14, 15]) # Ранг 3 >>> output = torch.empty([4], dtype=torch.int64) >>> dist.all_to_all_single(output, input) >>> output tensor([0, 4, 8, 12]) # Ранг 0 tensor([1, 5, 9, 13]) # Ранг 1 tensor([2, 6, 10, 14]) # Ранг 2 tensor([3, 7, 11, 15]) # Ранг 3 >>> # По сути, это похоже на следующую операцию: >>> scatter_list = list(input.chunk(world_size)) >>> gather_list = list(output.chunk(world_size)) >>> for i in range(world_size): >>> dist.scatter(gather_list[i], scatter_list if i == rank else [], src = i) >>> # Другой пример с неравномерным разделением >>> input tensor([0, 1, 2, 3, 4, 5]) # Ранг 0 tensor([10, 11, 12, 13, 14, 15, 16, 17, 18]) # Ранг 1 tensor([20, 21, 22, 23, 24]) # Ранг 2 tensor([30, 31, 32, 33, 34, 35, 36]) # Ранг 3 >>> input_splits [2, 2, 1, 1] # Ранг 0 [3, 2, 2, 2] # Ранг 1 [2, 1, 1, 1] # Ранг 2 [2, 2, 2, 1] # Ранг 3 >>> output_splits [2, 3, 2, 2] # Ранг 0 [2, 2, 1, 2] # Ранг 1 [1, 2, 1, 2] # Ранг 2 [1, 2, 1, 1] # Ранг 3 >>> output =... >>> dist.all_to_all_single(output, input, output_splits, input_splits) >>> output tensor([ 0, 1, 10, 11, 12, 20, 21, 30, 31]) # Ранг 0 tensor([ 2, 3, 13, 14, 22, 32, 33]) # Ранг 1 tensor([ 4, 15, 16, 23, 34, 35]) # Ранг 2 tensor([ 5, 17, 18, 24, 36]) # Ранг 3 >>> # Другой пример с тензорами типа torch.cfloat. >>> input = torch.tensor(... [1 + 1j, 2 + 2j, 3 + 3j, 4 + 4j], dtype=torch.cfloat... ) + 4 * rank * (1 + 1j) >>> input tensor([1+1j, 2+2j, 3+3j, 4+4j]) # Ранг 0 tensor([5+5j, 6+6j, 7+7j, 8+8j]) # Ранг 1 tensor([9+9j, 10+10j, 11+11j, 12+12j]) # Ранг 2 tensor([13+13j, 14+14j, 15+15j, 16+16j]) # Ранг 3 >>> output = torch.empty([4], dtype=torch.int64) >>> dist.all_to_all_single(output, input) >>> output tensor([1+1j, 5+5j, 9+9j, 13+13j]) # Ранг 0 tensor([2+2j, 6+6j, 10+10j, 14+14j]) # Ранг 1 tensor([3+3j, 7+7j, 11+11j, 15+15j]) # Ранг 2 tensor([4+4j, 8+8j, 12+12j, 16+16j]) # Ранг 3 torch.distributed.all_to_all(output_tensor_list, input_tensor_list, group=None, async_op=False)[source]# Разбрасывает список входных тензоров всем процессам в группе и возвращает собранный список тензоров в выходном списке. Поддерживаются комплексные тензоры. Параметры output_tensor_list (list[Tensor]) – Список тензоров для сбора по одному на ранг. input_tensor_list (list[Tensor]) – Список тензоров для разбрасывания по одному на ранг. group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Предупреждение all_to_all является экспериментальным и может быть изменен. Примеры >>> input = torch.arange(4) + rank * 4 >>> input = list(input.chunk(4)) >>> input [tensor([0]), tensor([1]), tensor([2]), tensor([3])] # Ранг 0 [tensor([4]), tensor([5]), tensor([6]), tensor([7])] # Ранг 1 [tensor([8]), tensor([9]), tensor([10]), tensor([11])] # Ранг 2 [tensor([12]), tensor([13]), tensor([14]), tensor([15])] # Ранг 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])] # Ранг 0 [tensor([1]), tensor([5]), tensor([9]), tensor([13])] # Ранг 1 [tensor([2]), tensor([6]), tensor([10]), tensor([14])] # Ранг 2 [tensor([3]), tensor([7]), tensor([11]), tensor([15])] # Ранг 3 >>> # По сути, это похоже на следующую операцию: >>> 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]) # Ранг 0 tensor([10, 11, 12, 13, 14, 15, 16, 17, 18]) # Ранг 1 tensor([20, 21, 22, 23, 24]) # Ранг 2 tensor([30, 31, 32, 33, 34, 35, 36]) # Ранг 3 >>> input_splits [2, 2, 1, 1] # Ранг 0 [3, 2, 2, 2] # Ранг 1 [2, 1, 1, 1] # Ранг 2 [2, 2, 2, 1] # Ранг 3 >>> output_splits [2, 3, 2, 2] # Ранг 0 [2, 2, 1, 2] # Ранг 1 [1, 2, 1, 2] # Ранг 2 [1, 2, 1, 1] # Ранг 3 >>> input = list(input.split(input_splits)) >>> input [tensor([0, 1]), tensor([2, 3]), tensor([4]), tensor([5])] # Ранг 0 [tensor([10, 11, 12]), tensor([13, 14]), tensor([15, 16]), tensor([17, 18])] # Ранг 1 [tensor([20, 21]), tensor([22]), tensor([23]), tensor([24])] # Ранг 2 [tensor([30, 31]), tensor([32, 33]), tensor([34, 35]), tensor([36])] # Ранг 3 >>> output =... >>> dist.all_to_all(output, input) >>> output [tensor([0, 1]), tensor([10, 11, 12]), tensor([20, 21]), tensor([30, 31])] # Ранг 0 [tensor([2, 3]), tensor([13, 14]), tensor([22]), tensor([32, 33])] # Ранг 1 [tensor([4]), tensor([15, 16]), tensor([23]), tensor([34, 35])] # Ранг 2 [tensor([5]), tensor([17, 18]), tensor([24]), tensor([36])] # Ранг 3 >>> # Другой пример с тензорами типа torch.cfloat. >>> 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])] # Ранг 0 [tensor([5+5j]), tensor([6+6j]), tensor([7+7j]), tensor([8+8j])] # Ранг 1 [tensor([9+9j]), tensor([10+10j]), tensor([11+11j]), tensor([12+12j])] # Ранг 2 [tensor([13+13j]), tensor([14+14j]), tensor([15+15j]), tensor([16+16j])] # Ранг 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])] # Ранг 0 [tensor([2+2j]), tensor([6+6j]), tensor([10+10j]), tensor([14+14j])] # Ранг 1 [tensor([3+3j]), tensor([7+7j]), tensor([11+11j]), tensor([15+15j])] # Ранг 2 [tensor([4+4j]), tensor([8+8j]), tensor([12+12j]), tensor([16+16j])] # Ранг 3 torch.distributed.barrier(group=None, async_op=False, device_ids=None)[source]# Синхронизирует все процессы. Этот коллектив блокирует процессы до тех пор, пока вся группа не войдет в эту функцию, если async_op равен False, или если для асинхронного дескриптора работы не вызван wait(). Параметры group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. async_op (bool, опционально) – Должна ли эта операция быть асинхронной. device_ids ([int], опционально) – Список идентификаторов устройств/GPU. Ожидается только один идентификатор. Возвращает Асинхронный дескриптор работы, если async_op установлен в True. None, если не async_op или не является частью группы. Примечание ProcessGroupNCCL теперь блокирует поток CPU до завершения коллектива barrier. Примечание ProcessGroupNCCL реализует barrier как all_reduce тензора из 1 элемента. Устройство должно быть выбрано для размещения этого тензора. Выбор устройства осуществляется путем проверки в следующем порядке: (1) первое устройство, переданное в аргумент device_ids barrier, если не None, (2) устройство, переданное в init_process_group, если не None, (3) устройство, которое было впервые использовано с этой группой процессов, если был выполнен другой коллектив с тензорными входами, (4) индекс устройства, указанный глобальным рангом по модулю количества локальных устройств. torch.distributed.monitored_barrier(group=None, timeout=None, wait_all_ranks=False)[source]# Синхронизирует процессы аналогично torch.distributed.barrier, но учитывает настраиваемый тайм-аут. Он может сообщать о рангах, которые не прошли этот барьер в течение заданного тайм-аута. В частности, для ненулевых рангов будет блокироваться до тех пор, пока не будет обработан send/recv от ранга 0. Ранг 0 будет блокироваться до тех пор, пока все send/recv от других рангов не будут обработаны, и сообщит о сбоях для рангов, которые не ответили вовремя. Обратите внимание, что если один ранг не достигает monitored_barrier (например, из-за зависания), все остальные ранги потерпят неудачу в monitored_barrier. Этот коллектив будет блокировать все процессы/ранги в группе до тех пор, пока вся группа не выйдет из функции успешно, что делает его полезным для отладки и синхронизации. Однако он может повлиять на производительность и должен использоваться только для отладки или сценариев, требующих полных точек синхронизации на стороне хоста. Для целей отладки этот барьер может быть вставлен перед коллективными вызовами приложения, чтобы проверить, не десинхронизированы ли какие-либо ранги. Примечание Обратите внимание, что этот коллектив поддерживается только бэкендом GLOO. Параметры group (ProcessGroup, опционально) – Группа процессов, с которой нужно работать. Если None, будет использоваться группа процессов по умолчанию. timeout (datetime.timedelta, опционально) – Тайм-аут для monitored_barrier. Если None, будет использоваться тайм-аут группы процессов по умолчанию. wait_all_ranks (bool, опционально) – Собирать ли все неудачные ранги или нет. По умолчанию это False, и monitored_barrier на ранге 0 вызовет исключение на первом же обнаруженном неудачном ранге для быстрого отказа. Установка wait_all_ranks=True заставит monitored_barrier собрать все неудачные ранги и выдать ошибку, содержащую информацию обо всех неудачных рангах. Возвращает None. Пример::>>> # Примечание: Инициализация группы процессов опущена на каждом ранге. >>> import torch.distributed as dist >>> if dist.get_rank()!= 1: >>> dist.monitored_barrier() # Вызывает исключение, указывающее, что >>> # ранг 1 не вызвал monitored_barrier. >>> # Пример с wait_all_ranks=True >>> if dist.get_rank() == 0: >>> dist.monitored_barrier(wait_all_ranks=True) # Вызывает исключение >>> # указывающее, что ранги 1, 2,... world_size - 1 не вызвали >>> # monitored_barrier. class torch.distributed.Work# Объект Work представляет дескриптор ожидающей асинхронной операции в распределенном пакете PyTorch. Он возвращается неблокирующими коллективными операциями, такими как dist.all_reduce(tensor, async_op=True). block_current_stream(self: torch.C._distributed_c10d.Work) → None# Блокирует текущий активный поток GPU до завершения операции. Для коллективов на основе GPU это эквивалентно синхронизации. Для коллективов, инициированных CPU, таких как с Gloo, это заблокирует поток CUDA до завершения операции. В любом случае это возвращается немедленно. Чтобы проверить, была ли операция успешной, вы должны асинхронно проверить результат объекта Work. boxed(self: torch._C._distributed_c10d.Work) → object# exception(self: torch._C._distributed_c10d.Work) → std::__exception_ptr::exception_ptr# get_future(self: torch._C._distributed_c10d.Work) → torch.Future# Возвращает Объект torch.futures.Future, связанный с завершением Work. Например, объект future может быть получен с помощью fut = process_group.allreduce(tensors).get_future(). Пример::Ниже приведен пример простого хука коммуникации DDP allreduce, который использует API get_future для получения Future, связанного с завершением allreduce. >>> def allreduce(process_group: dist.ProcessGroup, bucket: dist.GradBucket): -> torch.futures.Future >>> group_to_use = process_group if process_group is not None else torch.distributed.group.WORLD >>> tensor = bucket.buffer().div(group_to_use.size()) >>> return torch.distributed.all_reduce(tensor, group=group_to_use, async_op=True).get_future() >>> ddp_model.register_comm_hook(state=None, hook=allreduce) Предупреждение API get_future поддерживает NCCL и частично бэкенды GLOO и MPI (без поддержки операций точка-точка, таких как send/recv) и вернет torch.futures.Future. В приведенном выше примере работа allreduce будет выполнена на GPU с использованием бэкенда NCCL, fut.wait() вернется после синхронизации соответствующих потоков NCCL с текущими потоками устройства PyTorch, чтобы обеспечить асинхронное выполнение CUDA, и не будет ждать завершения всей операции на GPU. Обратите внимание, что CUDAFuture не поддерживает флаг TORCH_NCCL_BLOCKING_WAIT или barrier() NCCL. Кроме того, если функция обратного вызова была добавлена с помощью fut.then(), она будет ждать, пока потоки NCCL WorkNCCL синхронизируются с выделенным потоком обратного вызова ProcessGroupNCCL, и вызовет обратный вызов встроенно после выполнения обратного вызова в потоке обратного вызова. fut.then() вернет другой CUDAFuture, который содержит возвращаемое значение обратного вызова и CUDAEvent, который записал поток обратного вызова. Для работы CPU fut.done() возвращает true, когда работа завершена и тензоры value() готовы. Для работы GPU fut.done() возвращает true только в том случае, если операция была поставлена в очередь. Для смешанной работы CPU-GPU (например, отправка тензоров GPU с GLOO) fut.done() возвращает true, когда тензоры прибыли на соответствующие узлы, но еще не обязательно синхронизированы на соответствующих GPU (аналогично работе GPU). get_future_result(self: torch._C._distributed_c10d.Work) → torch.Future# Возвращает Объект torch.futures.Future типа int, который соответствует типу перечисления WorkResult. Например, объект future может быть получен с помощью fut = process_group.allreduce(tensor).get_future_result(). Пример::пользователи могут использовать fut.wait() для блокирующего ожидания завершения работы и получения WorkResult с помощью fut.value(). Также пользователи могут использовать fut.then(call_back_func) для регистрации функции обратного вызова, которая будет вызвана при завершении работы, без блокировки текущего потока. Предупреждение API get_future_result поддерживает NCCL. is_completed(self: torch._C._distributed_c10d.Work) → bool# is_success(self: torch._C._distributed_c10d.Work) → bool# result(self: torch._C._distributed_c10d.Work) → list[torch.Tensor]# source_rank(self: torch._C._distributed_c10d.Work) → int# synchronize(self: torch._C._distributed_c10d.Work) → None# static unbox(arg0: object) → torch._C._distributed_c10d.Work# wait(self: torch._C._distributed_c10d.Work, timeout: datetime.timedelta = datetime.timedelta(0)) → bool# Возвращает true/false. Пример:: try:work.wait(timeout) except:# некоторая обработка Предупреждение В обычных случаях пользователям не нужно устанавливать тайм-аут. Вызов wait() эквивалентен вызову synchronize(): позволяет текущему потоку блокироваться до завершения работы NCCL. Однако, если установлен тайм-аут, он будет блокировать поток CPU до завершения работы NCCL или истечения тайм-аута. Если тайм-аут истечет, будет выброшено исключение. class torch.distributed.ReduceOp# Класс, подобный перечислению, для доступных операций редукции: SUM, PRODUCT, MIN, MAX, BAND, BOR, BXOR и PREMUL_SUM. Редукции BAND, BOR и BXOR недоступны при использовании бэкенда NCCL. AVG делит значения на world_size перед суммированием по рангам. 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(). Этот класс не поддерживает свойство members. class torch.distributed.reduce_op# Устаревший класс, подобный перечислению, для операций редукции: SUM, PRODUCT, MIN и MAX. Рекомендуется использовать ReduceOp. Распределенное хранилище ключ-значение# Пакет распределенных вычислений поставляется с распределенным хранилищем ключ-значение, которое может использоваться для обмена информацией между процессами в группе, а также для инициализации пакета распределенных вычислений в torch.distributed.init_process_group() (путем явного создания хранилища в качестве альтернативы указанию init_method). Есть 3 варианта хранилищ ключ-значение: TCPStore, FileStore и HashStore. class torch.distributed.Store# Базовый класс для всех реализаций хранилищ, таких как 3, предоставляемые PyTorch distributed: (TCPStore, FileStore и HashStore). init(self: torch._C._distributed_c10d.Store) → None# add(self: torch._C._distributed_c10d.Store, arg0: str, arg1: SupportsInt) → int# Первый вызов add для данного ключа создает счетчик, связанный с ключом в хранилище, инициализированный значением amount. Последующие вызовы add с тем же ключом увеличивают счетчик на указанную величину. Вызов add() с ключом, который уже был установлен в хранилище с помощью set(), приведет к исключению. Параметры key (str) – Ключ в хранилище, счетчик которого будет увеличен. amount (int) – Величина, на которую будет увеличен счетчик. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> # Используем TCPStore в качестве примера, другие типы хранилищ также могут использоваться >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.add("first_key", 1) >>> store.add("first_key", 6) >>> # Должно вернуть 7 >>> store.get("first_key") append(self: torch._C._distributed_c10d.Store, arg0: str, arg1: str) → None# Добавляет пару ключ-значение в хранилище на основе предоставленного ключа и значения. Если ключ не существует в хранилище, он будет создан. Параметры key (str) – Ключ, который нужно добавить в хранилище. value (str) – Значение, связанное с ключом, которое нужно добавить в хранилище. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.append("first_key", "po") >>> store.append("first_key", "tato") >>> # Должно вернуть "potato" >>> store.get("first_key") check(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str]) → bool# Вызов для проверки, есть ли у заданного списка ключей значение, сохраненное в хранилище. Этот вызов немедленно возвращается в обычных случаях, но все еще страдает от некоторых крайних случаев взаимоблокировки, например, вызов check после того, как TCPStore был уничтожен. Вызов check() со списком ключей, которые нужно проверить, сохранены ли они в хранилище. Параметры keys (list[str]) – Ключи для запроса, сохранены ли они в хранилище. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> # Используем TCPStore в качестве примера, другие типы хранилищ также могут использоваться >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.add("first_key", 1) >>> # Должно вернуть 7 >>> store.check(["first_key"]) clone(self: torch._C._distributed_c10d.Store) → torch._C._distributed_c10d.Store# Клонирует хранилище и возвращает новый объект, указывающий на то же базовое хранилище. Возвращенное хранилище может использоваться одновременно с исходным объектом. Это предназначено для обеспечения безопасного способа использования хранилища из нескольких потоков путем клонирования одного хранилища на поток. compare_set(self: torch._C._distributed_c10d.Store, arg0: str, arg1: str, arg2: str) → bytes# Вставляет пару ключ-значение в хранилище на основе предоставленного ключа и выполняет сравнение между expected_value и desired_value перед вставкой. desired_value будет установлен только в том случае, если expected_value для ключа уже существует в хранилище или если expected_value является пустой строкой. Параметры key (str) – Ключ, который нужно проверить в хранилище. expected_value (str) – Значение, связанное с ключом, которое нужно проверить перед вставкой. desired_value (str) – Значение, связанное с ключом, которое нужно добавить в хранилище. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.set("key", "first_value") >>> store.compare_set("key", "first_value", "second_value") >>> # Должно вернуть "second_value" >>> store.get("key") delete_key(self: torch._C._distributed_c10d.Store, arg0: str) → bool# Удаляет пару ключ-значение, связанную с key, из хранилища. Возвращает true, если ключ был успешно удален, и false, если нет. Предупреждение API delete_key поддерживается только TCPStore и HashStore. Использование этого API с FileStore приведет к исключению. Параметры key (str) – Ключ, который нужно удалить из хранилища. Возвращает True, если ключ был удален, иначе False. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> # Используем TCPStore в качестве примера, HashStore также может использоваться >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.set("first_key") >>> # Это должно вернуть true >>> store.delete_key("first_key") >>> # Это должно вернуть false >>> store.delete_key("bad_key") get(self: torch._C._distributed_c10d.Store, arg0: str) → bytes# Извлекает значение, связанное с заданным ключом в хранилище. Если ключ отсутствует в хранилище, функция будет ждать тайм-аута, который определяется при инициализации хранилища, прежде чем выбросить исключение. Параметры key (str) – Функция вернет значение, связанное с этим ключом. Возвращает Значение, связанное с ключом, если ключ есть в хранилище. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.set("first_key", "first_value") >>> # Должно вернуть "first_value" >>> store.get("first_key") has_extended_api(self: torch._C._distributed_c10d.Store) → bool# Возвращает true, если хранилище поддерживает расширенные операции. multi_get(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str]) → list[bytes]# Извлекает все значения по ключам. Если какой-либо ключ в keys отсутствует в хранилище, функция будет ждать тайм-аута. Параметры keys (List[str]) – Ключи для извлечения из хранилища. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.set("first_key", "po") >>> store.set("second_key", "tato") >>> # Должно вернуть [b"po", b"tato"] >>> store.multi_get(["first_key", "second_key"]) multi_set(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str], arg1: collections.abc.Sequence[str]) → None# Вставляет список пар ключ-значение в хранилище на основе предоставленных ключей и значений. Параметры keys (List[str]) – Ключи для вставки. values (List[str]) – Значения для вставки. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.multi_set(["first_key", "second_key"], ["po", "tato"]) >>> # Должно вернуть b"po" >>> store.get("first_key") num_keys(self: torch._C._distributed_c10d.Store) → int# Возвращает количество ключей, установленных в хранилище. Обратите внимание, что это число обычно будет на единицу больше, чем количество ключей, добавленных с помощью set() и add(), поскольку один ключ используется для координации всех рабочих, использующих хранилище. Предупреждение При использовании с TCPStore num_keys возвращает количество ключей, записанных в базовый файл. Если хранилище уничтожено и другое хранилище создано с тем же файлом, исходные ключи будут сохранены. Возвращает Количество ключей, присутствующих в хранилище. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> # Используем TCPStore в качестве примера, другие типы хранилищ также могут использоваться >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.set("first_key", "first_value") >>> # Это должно вернуть 2 >>> store.num_keys() queue_len(self: torch._C._distributed_c10d.Store, arg0: str) → int# Возвращает длину указанной очереди. Если очередь не существует, возвращает 0. См. queue_push для более подробной информации. Параметры key (str) – Ключ очереди, для которой нужно получить длину. queue_pop(self: torch._C._distributed_c10d.Store, key: str, block: bool = True) → bytes# Извлекает значение из указанной очереди или ждет тайм-аута, если очередь пуста. См. queue_push для более подробной информации. Если block равен False, будет вызвано dist.QueueEmptyError, если очередь пуста. Параметры key (str) – Ключ очереди, из которой нужно извлечь. block (bool) – Блокировать ли ожидание ключа или немедленно вернуться. queue_push(self: torch._C._distributed_c10d.Store, arg0: str, arg1: str) → None# Помещает значение в указанную очередь. Использование одного и того же ключа для очередей и операций set/get может привести к неожиданному поведению. Операции wait/check поддерживаются для очередей. wait с очередями разбудит только одного ожидающего рабочего, а не всех. Параметры key (str) – Ключ очереди, в которую нужно поместить. value (str) – Значение для помещения в очередь. set(self: torch._C._distributed_c10d.Store, arg0: str, arg1: str) → None# Вставляет пару ключ-значение в хранилище на основе предоставленного ключа и значения. Если ключ уже существует в хранилище, он перезапишет старое значение новым предоставленным значением. Параметры key (str) – Ключ, который нужно добавить в хранилище. value (str) – Значение, связанное с ключом, которое нужно добавить в хранилище. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.set("first_key", "first_value") >>> # Должно вернуть "first_value" >>> store.get("first_key") set_timeout(self: torch._C._distributed_c10d.Store, arg0: datetime.timedelta) → None# Устанавливает тайм-аут хранилища по умолчанию. Этот тайм-аут используется во время инициализации и в wait() и get(). Параметры timeout (timedelta) – тайм-аут, который нужно установить в хранилище. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> # Используем TCPStore в качестве примера, другие типы хранилищ также могут использоваться >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> store.set_timeout(timedelta(seconds=10)) >>> # Это вызовет исключение через 10 секунд >>> store.wait(["bad_key"]) property timeout# Получает тайм-аут хранилища. wait(args, *kwargs)# Перегруженная функция. wait(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str]) -> None Ожидает, пока каждый ключ в keys не будет добавлен в хранилище. Если не все ключи установлены до истечения тайм-аута (установленного при инициализации хранилища), то wait вызовет исключение. Параметры keys (list) – Список ключей, которые нужно ждать, пока они не будут установлены в хранилище. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> # Используем TCPStore в качестве примера, другие типы хранилищ также могут использоваться >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> # Это вызовет исключение через 30 секунд >>> store.wait(["bad_key"]) wait(self: torch._C._distributed_c10d.Store, arg0: collections.abc.Sequence[str], arg1: datetime.timedelta) -> None Ожидает, пока каждый ключ в keys не будет добавлен в хранилище, и вызывает исключение, если ключи не были установлены в течение указанного тайм-аута. Параметры keys (list) – Список ключей, которые нужно ждать, пока они не будут установлены в хранилище. timeout (timedelta) – Время ожидания добавления ключей перед вызовом исключения. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> # Используем TCPStore в качестве примера, другие типы хранилищ также могут использоваться >>> store = dist.TCPStore("127.0.0.1", 0, 1, True, timedelta(seconds=30)) >>> # Это вызовет исключение через 10 секунд >>> store.wait(["bad_key"], timedelta(seconds=10)) class torch.distributed.TCPStore# Реализация распределенного хранилища ключ-значение на основе TCP. Серверное хранилище хранит данные, в то время как клиентские хранилища могут подключаться к серверному хранилищу по TCP и выполнять такие действия, как set() для вставки пары ключ-значение, get() для извлечения пары ключ-значение и т.д. Всегда должно быть инициализировано одно серверное хранилище, потому что клиентские хранилища будут ждать, пока сервер установит соединение. Параметры host_name (str) – Имя хоста или IP-адрес, на котором должно работать серверное хранилище. port (int) – Порт, на котором серверное хранилище должно прослушивать входящие запросы. world_size (int, опционально) – Общее количество пользователей хранилища (количество клиентов + 1 для сервера). По умолчанию None (None указывает на нефиксированное количество пользователей хранилища). is_master (bool, опционально) – True при инициализации серверного хранилища и False для клиентских хранилищ. По умолчанию False. timeout (timedelta, опционально) – Тайм-аут, используемый хранилищем во время инициализации и для таких методов, как get() и wait(). По умолчанию timedelta(seconds=300). wait_for_workers (bool, опционально) – Ожидать ли подключения всех рабочих к серверному хранилищу. Это применимо только тогда, когда world_size является фиксированным значением. По умолчанию True. multi_tenant (bool, опционально) – Если True, все экземпляры TCPStore в текущем процессе с тем же хостом/портом будут использовать один и тот же базовый TCPServer. По умолчанию False. master_listen_fd (int, опционально) – Если указано, базовый TCPServer будет прослушивать этот файловый дескриптор, который должен быть сокетом, уже привязанным к порту. Для привязки эфемерного порта мы рекомендуем установить порт в 0 и прочитать.port. По умолчанию None (означает, что сервер создает новый сокет и пытается привязать его к порту). use_libuv (bool, опционально) – Если True, использовать libuv для бэкенда TCPServer. По умолчанию True. Пример::>>> import torch.distributed as dist >>> from datetime import timedelta >>> # Запуск на процессе 1 (сервер) >>> server_store = dist.TCPStore("127.0.0.1", 1234, 2, True, timedelta(seconds=30)) >>> # Запуск на процессе 2 (клиент) >>> client_store = dist.TCPStore("127.0.0.1", 1234, 2, False) >>> # Используйте любой из методов хранилища от клиента или сервера после инициализации >>> server_store.set("first_key", "first_value") >>> client_store.get("first_key") init(self: torch._C._distributed_c10d.TCPStore, host_name: str, port: SupportsInt, world_size: SupportsInt | None = None, is_master: bool = False, timeout: datetime.timedelta = datetime.timedelta(seconds=300), wait_for_workers: bool = True, multi_tenant: bool = False, master_listen_fd: SupportsInt | None = None, use_libuv: bool = True) → None# Создает новый TCPStore. property host# Получает имя хоста, на котором хранилище прослушивает запросы. property libuvBackend# Возвращает True, если используется бэкенд libuv. property port# Получает номер порта, на котором хранилище прослушивает запросы. class torch.distributed.HashStore# Потокобезопасная реализация хранилища на основе базовой хэш-карты. Это хранилище может использоваться в рамках одного процесса (например, другими потоками), но не может использоваться между процессами. Пример::>>> import torch.distributed as dist >>> store = dist.HashStore() >>> # хранилище может использоваться из других потоков >>> # Используйте любой из методов хранилища после инициализации >>> store.set("first_key", "first_value") init(self: torch._C._distributed_c10d.HashStore) → None# Создает новый HashStore. class torch.distributed.FileStore# Реализация хранилища, которая использует файл для хранения базовых пар ключ-значение. Параметры file_name (str) – путь к файлу, в котором хранить пары ключ-значение. world_size (int, опционально) – Общее количество процессов, использующих хранилище. По умолчанию -1 (отрицательное значение указывает на нефиксированное количество пользователей хранилища). Пример::>>> import torch.distributed as dist >>> store1 = dist.FileStore("/tmp/filestore", 2) >>> store2 = dist.FileStore("/tmp/filestore", 2) >>> # Используйте любой из методов хранилища от клиента или сервера после инициализации >>> store1.set("first_key", "first_value") >>> store2.get("first_key") init(self: torch._C._distributed_c10d.FileStore, file_name: str, world_size: SupportsInt = -1) → None# Создает новый FileStore. property path# Получает путь к файлу, используемому FileStore для хранения пар ключ-значение. class torch.distributed.PrefixStore# Обертка вокруг любого из 3 хранилищ ключ-значение (TCPStore, FileStore и HashStore), которая добавляет префикс к каждому ключу, вставляемому в хранилище. Параметры prefix (str) – Строка префикса, которая добавляется к каждому ключу перед вставкой в хранилище. store (torch.distributed.store) – Объект хранилища, который образует базовое хранилище ключ-значение. init(self: torch._C._distributed_c10d.PrefixStore, prefix: str, store: torch._C._distributed_c10d.Store) → None# Создает новый PrefixStore. property underlying_store# Получает базовый объект хранилища, который оборачивает PrefixStore. Профилирование коллективной коммуникации# Обратите внимание, что вы можете использовать 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 на поток CPU) устарели. На сегодняшний день предпочтительной моделью программирования PyTorch Distributed является одно устройство на поток, как показано на примерах API в этом документе. Если вы разработчик бэкенда и хотите поддерживать несколько устройств на поток, пожалуйста, свяжитесь с мейнтейнерами PyTorch Distributed. Объектные коллективы# Предупреждение Объектные коллективы имеют ряд серьезных ограничений. Прочитайте дальше, чтобы определить, безопасно ли их использовать для вашего случая. Объектные коллективы — это набор операций, подобных коллективным, которые работают с произвольными объектами Python, если они могут быть упакованы pickle. Реализованы различные коллективные шаблоны (например, broadcast, all_gather, …), но каждый из них примерно следует этому шаблону: преобразовать входной объект в pickle (сырые байты), затем поместить его в байтовый тензор. Сообщить размер этого байтового тензора пирам (первая коллективная операция). Выделить тензор соответствующего размера для выполнения реального коллектива. Передать данные объекта (вторая коллективная операция). Преобразовать сырые данные обратно в Python (unpickle). Объектные коллективы иногда имеют неожиданные характеристики производительности или памяти, которые приводят к длительному времени выполнения или OOM, и поэтому их следует использовать с осторожностью. Вот некоторые распространенные проблемы. Асимметричное время pickle/unpickle — Упаковка объектов может быть медленной, в зависимости от количества, типа и размера объектов. Когда коллектив имеет тип fan-in (например, gather_object), принимающий ранг(и) должен распаковать в N раз больше объектов, чем отправляющий ранг(и) должен был упаковать, что может привести к тайм-ауту других рангов на их следующем коллективе. Неэффективная коммуникация тензоров — Тензоры должны отправляться через обычные коллективные API, а не через объектные коллективные API. Можно отправлять тензоры через объектные коллективные API, но они будут сериализованы и десериализованы (включая синхронизацию CPU и копирование с устройства на хост в случае тензоров не CPU), и почти в каждом случае, кроме отладки или кода устранения неполадок, стоит потрудиться переписать код для использования не-объектных коллективов. Неожиданные устройства тензоров — Если вы все еще хотите отправлять тензоры через объектные коллективы, есть еще один аспект, специфичный для тензоров cuda (и, возможно, других ускорителей). Если вы упакуете тензор, который в данный момент находится на cuda:3, а затем распакуете его, вы получите другой тензор на cuda:3 независимо от того, в каком процессе вы находитесь или какое устройство CUDA является «устройством по умолчанию» для этого процесса. С обычными коллективными API тензоров «выходные тензоры» всегда будут на одном и том же локальном устройстве, что обычно и ожидается. Распаковка тензора неявно активирует контекст CUDA, если это первый раз, когда GPU используется процессом, что может привести к значительному расходу памяти GPU. Этой проблемы можно избежать, переместив тензоры на CPU перед передачей их в качестве входных данных для объектного коллектива. Сторонние бэкенды# Помимо встроенных бэкендов GLOO/MPI/NCCL, PyTorch distributed поддерживает сторонние бэкенды через механизм регистрации во время выполнения. Справочную информацию о том, как разработать сторонний бэкенд через C++ Extension, см. в Tutorials - Custom C++ and CUDA Extensions и test/cpp_extensions/cpp_c10d_extension.cpp. Возможности сторонних бэкендов определяются их собственными реализациями. Новый бэкенд наследуется от c10d::ProcessGroup и регистрирует имя бэкенда и интерфейс создания экземпляра через torch.distributed.Backend.register_backend() при импорте. При ручном импорте этого бэкенда и вызове torch.distributed.init_process_group() с соответствующим именем бэкенда пакет torch.distributed будет работать на новом бэкенде. Предупреждение Поддержка стороннего бэкенда является экспериментальной и может быть изменена. Утилита запуска# Пакет torch.distributed также предоставляет утилиту запуска в torch.distributed.launch. Эта вспомогательная утилита может использоваться для запуска нескольких процессов на узел для распределенного обучения. Модуль torch.distributed.launch. 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 и все остальные аргументы вашего обучающего скрипта) Многоузловое многопроцессное распределенное обучение: (например, два узла) Узел 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 и все остальные аргументы вашего обучающего скрипта) Узел 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 и все остальные аргументы вашего обучающего скрипта) Чтобы узнать, какие опциональные аргументы предлагает этот модуль: 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", "--local_rank", type=int) >>> args = parser.parse_args() Установите ваше устройство на локальный ранг, используя либо >>> torch.cuda.set_device(args.local_rank) # перед запуском вашего кода или >>> with torch.cuda.device(args.local_rank): >>> # ваш код для запуска >>>... Изменено в версии 2.0.0: Лаунчер передает аргумент --local-rank=<rank> вашему скрипту. Начиная с PyTorch 2.0.0, предпочтительнее использовать дефисную форму --local-rank вместо ранее использовавшейся подчеркнутой --local_rank. Для обратной совместимости пользователям может потребоваться обрабатывать оба случая в коде разбора аргументов. Это означает включение как "--local-rank", так и "--local_rank" в парсер аргументов. Если указан только "--local_rank", лаунчер выдаст ошибку: "error: unrecognized arguments: –local-rank=<rank>". Для обучающего кода, поддерживающего только PyTorch 2.0.0+, включения "--local-rank" должно быть достаточно. 3. В вашей обучающей программе вы должны вызвать следующую функцию в начале для запуска распределенного бэкенда. Настоятельно рекомендуется использовать init_method=env://. Другие методы init (например, tcp://) могут работать, но env:// — это тот, который официально поддерживается этим модулем. >>> torch.distributed.init_process_group(backend='YOUR BACKEND', >>> init_method='env://') 4. В вашей обучающей программе вы можете использовать либо обычные распределенные функции, либо модуль torch.nn.parallel.DistributedDataParallel(). Если ваша обучающая программа использует GPU для обучения, и вы хотите использовать модуль torch.nn.parallel.DistributedDataParallel(), вот как его настроить. >>> model = torch.nn.parallel.DistributedDataParallel(model, >>> device_ids=[args.local_rank], >>> output_device=args.local_rank) Пожалуйста, убедитесь, что аргумент device_ids установлен на единственный идентификатор устройства GPU, с которым будет работать ваш код. Обычно это локальный ранг процесса. Другими словами, device_ids должен быть [args.local_rank], а output_device должен быть args.local_rank для использования этой утилиты. 5. Другой способ передать local_rank подпроцессам — через переменную окружения LOCAL_RANK. Это поведение включается, когда вы запускаете скрипт с --use-env=True. Вы должны скорректировать пример подпроцесса выше, заменив args.local_rank на os.environ['LOCAL_RANK']; лаунчер не будет передавать --local-rank, когда вы укажете этот флаг. Предупреждение local_rank НЕ является глобально уникальным: он уникален только для процесса на машине. Таким образом, не используйте его для принятия решения, например, о записи в сетевую файловую систему. См. pytorch/pytorch#12042 для примера того, как все может пойти не так, если вы не сделаете это правильно. Утилита порождения# Пакет Multiprocessing package - torch.multiprocessing также предоставляет функцию spawn в torch.multiprocessing.spawn(). Эта вспомогательная функция может использоваться для порождения нескольких процессов. Она работает путем передачи функции, которую вы хотите запустить, и порождает N процессов для ее выполнения. Это также может использоваться для многопроцессного распределенного обучения. Справочную информацию о том, как ее использовать, см. в PyTorch example - ImageNet implementation. Обратите внимание, что эта функция требует Python 3.4 или выше. Отладка приложений torch.distributed# Отладка распределенных приложений может быть сложной из-за трудно понимаемых зависаний, сбоев или несогласованного поведения между рангами. torch.distributed предоставляет набор инструментов для самостоятельной отладки обучающих приложений: Точка останова Python# Использовать отладчик Python в распределенной среде очень удобно, но, поскольку он не работает «из коробки», многие люди вообще его не используют. PyTorch предлагает настраиваемую обертку вокруг pdb, которая упрощает процесс. torch.distributed.breakpoint делает этот процесс легким. Внутренне он настраивает поведение точки останова pdb двумя способами, но в остальном ведет себя как обычный pdb. Прикрепляет отладчик только к одному рангу (указанному пользователем). Гарантирует, что все остальные ранги останавливаются, используя torch.distributed.barrier(), который будет снят, когда отлаживаемый ранг выдаст continue. Перенаправляет stdin из дочернего процесса таким образом, чтобы он подключался к вашему терминалу. Чтобы использовать его, просто вызовите torch.distributed.breakpoint(rank) на всех рангах, используя одно и то же значение для rank в каждом случае. Monitored Barrier# Начиная с версии v1.10, torch.distributed.monitored_barrier() существует как альтернатива torch.distributed.barrier(), которая при сбое выдает полезную информацию о том, какой ранг может быть неисправен, т.е. не все ранги вызывают torch.distributed.monitored_barrier() в течение заданного тайм-аута. torch.distributed.monitored_barrier() реализует барьер на стороне хоста, используя примитивы связи send/recv в процессе, похожем на подтверждения, позволяя рангу 0 сообщить, какой ранг(и) не смог подтвердить барьер вовремя. В качестве примера рассмотрим следующую функцию, где ранг 1 не вызывает torch.distributed.monitored_barrier() (на практике это может быть связано с ошибкой приложения или зависанием в предыдущем коллективе): import os from datetime import timedelta import torch import torch.distributed as dist import torch.multiprocessing as mp def worker(rank): dist.init_process_group("nccl", rank=rank, world_size=2) # monitored barrier требует группу процессов gloo для выполнения синхронизации на стороне хоста. 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" # установить DETAIL для логирования во время выполнения. 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(), если есть параметры, которые могут быть не использованы в прямом проходе, и начиная с v1.10, все выходные данные модели должны использоваться в вычислении потерь, поскольку torch.nn.parallel.DistributedDataParallel() не поддерживает неиспользуемые параметры в обратном проходе. Эти ограничения особенно сложны для больших моделей, поэтому при сбое с ошибкой torch.nn.parallel.DistributedDataParallel() будет логировать полное квалифицированное имя всех параметров, которые остались неиспользованными. Например, в приведенном выше приложении, если мы изменим потери на вычисление как loss = output[1], то TwoLinLayerNet.a не получит градиент в обратном проходе, и, следовательно, DDP завершится ошибкой. При сбое пользователю передается информация о параметрах, которые остались неиспользованными, что может быть сложно найти вручную для больших моделей: RuntimeError: Expected to have finished reduction in the prior iteration before starting a new one. This error indicates that your module has parameters that were not used in producing loss. You can enable unused parameter detection by passing the keyword argument find_unused_parameters=True to torch.nn.parallel.DistributedDataParallel, and by making sure all forward function outputs participate in calculating loss. If you already have done the above, then the distributed data parallel module wasn't able to locate the output tensors in the return value of your module's forward function. Please include the loss function and the structure of the return va lue of forward of your module when reporting this issue (e.g. list, dict, iterable). Parameters which did not receive grad for rank 0: a.weight Parameter indices which did not receive grad for rank 0: 0 Установка TORCH_DISTRIBUTED_DEBUG=DETAIL вызовет дополнительные проверки согласованности и синхронизации при каждом коллективном вызове, инициированном пользователем напрямую или косвенно (например, DDP allreduce). Это делается путем создания оберточной группы процессов, которая оборачивает все группы процессов, возвращаемые API torch.distributed.init_process_group() и torch.distributed.new_group(). В результате эти API вернут оберточную группу процессов, которую можно использовать точно так же, как обычную группу процессов, но она выполняет проверки согласованности перед отправкой коллектива в базовую группу процессов. В настоящее время эти проверки включают torch.distributed.monitored_barrier(), которая гарантирует, что все ранги завершат свои ожидающие коллективные вызовы, и сообщает о рангах, которые застряли. Затем сам коллектив проверяется на согласованность путем обеспечения соответствия всех коллективных функций и их вызова с согласованными формами тензоров. Если это не так, при сбое приложения включается подробный отчет об ошибке, а не зависание или неинформативное сообщение об ошибке. В качестве примера рассмотрим следующую функцию, которая имеет несовпадающие формы входных данных в torch.distributed.all_reduce(): import torch import torch.distributed as dist import torch.multiprocessing as mp def worker(rank): dist.init_process_group("nccl", rank=rank, world_size=2) torch.cuda.set_device(rank) tensor = torch.randn(10 if rank == 0 else 20).cuda() dist.all_reduce(tensor) torch.cuda.synchronize(device=rank) if name == "main": os.environ["MASTER_ADDR"] = "localhost" os.environ["MASTER_PORT"] = "29501" os.environ["TORCH_CPP_LOG_LEVEL"]="INFO" os.environ["TORCH_DISTRIBUTED_DEBUG"] = "DETAIL" mp.spawn(worker, nprocs=2, args=()) С бэкендом NCCL такое приложение, скорее всего, приведет к зависанию, которое может быть сложно диагностировать в нетривиальных сценариях. Если пользователь включит TORCH_DISTRIBUTED_DEBUG=DETAIL и перезапустит приложение, следующее сообщение об ошибке раскроет первопричину: work = default_pg.allreduce([tensor], opts) RuntimeError: Error when verifying shape tensors for collective ALLREDUCE on rank 0. This likely indicates that input shapes into the collective are mismatched across ranks. Got shapes: 10 [ torch.LongTensor{1} ] Примечание Для детального контроля уровня отладки во время выполнения также могут использоваться функции torch.distributed.set_debug_level(), torch.distributed.set_debug_level_from_env() и torch.distributed.get_debug_level(). Кроме того, TORCH_DISTRIBUTED_DEBUG=DETAIL может использоваться вместе с TORCH_SHOW_CPP_STACKTRACES=1 для логирования полного стека вызовов при обнаружении десинхронизации коллектива. Эти проверки десинхронизации коллектива будут работать для всех приложений, которые используют коллективные вызовы c10d, поддерживаемые группами процессов, созданными с помощью API torch.distributed.init_process_group() и torch.distributed.new_group(). Логирование# В дополнение к явной поддержке отладки через torch.distributed.monitored_barrier() и TORCH_DISTRIBUTED_DEBUG, базовая библиотека C++ torch.distributed также выводит сообщения журнала на различных уровнях. Эти сообщения могут быть полезны для понимания состояния выполнения распределенной обучающей задачи и устранения таких проблем, как сбои сетевого подключения. В следующей матрице показано, как можно настроить уровень логирования с помощью комбинации переменных окружения TORCH_CPP_LOG_LEVEL и TORCH_DISTRIBUTED_DEBUG. TORCH_CPP_LOG_LEVEL TORCH_DISTRIBUTED_DEBUG Эффективный уровень логирования ERROR игнорируется Error WARNING игнорируется Warning INFO игнорируется Info INFO DEBUG Info INFO DETAIL Trace (также известный как All) Распределенные компоненты вызывают пользовательские типы исключений, производные от RuntimeError: torch.distributed.DistError: Это базовый тип всех распределенных исключений. torch.distributed.DistBackendError: Это исключение выбрасывается, когда происходит ошибка, специфичная для бэкенда. Например, если используется бэкенд NCCL и пользователь пытается использовать GPU, недоступный для библиотеки NCCL. torch.distributed.DistNetworkError: Это исключение выбрасывается, когда сетевые библиотеки сталкиваются с ошибками (например, Connection reset by peer). torch.distributed.DistStoreError: Это исключение выбрасывается, когда Store сталкивается с ошибкой (например, TCPStore timeout). class torch.distributed.DistError# Исключение, возникающее при ошибке в распределенной библиотеке. class torch.distributed.DistBackendError# Исключение, возникающее при ошибке бэкенда в распределенной среде. class torch.distributed.DistNetworkError# Исключение, возникающее при сетевой ошибке в распределенной среде. class torch.distributed.DistStoreError# Исключение, возникающее при ошибке в распределенном хранилище. Если вы запускаете одноузловое обучение, может быть удобно интерактивно установить точку останова в вашем скрипте. Мы предлагаем способ удобной установки точки останова для одного ранга: torch.distributed.breakpoint(rank=0, skip=0, timeout_s=3600)[source]# Устанавливает точку останова, но только на одном ранге. Все остальные ранги будут ждать, пока вы закончите с точкой останова, прежде чем продолжить. Параметры rank (int) – Какой ранг остановить. По умолчанию: 0. skip (int) – Пропустить первые skip вызовов этой точки останова. По умолчанию: 0.
torch.distributed
Шаблон 3: Инициализация# Пакет необходимо инициализировать с помощью функции torch.distributed.init_process_group() или torch.distributed.device_mesh.init_device_mesh() перед вызовом любых других методов. Обе блокируются до тех пор, пока все процессы не присоединятся. Предупреждение Инициализация не является потокобезопасной. Создание группы процессов должно выполняться из одного потока, чтобы предотвратить непоследовательное назначение «UUID» между рангами и предотвратить состояния гонки во время инициализации, которые могут привести к зависаниям. torch.distributed.is_available()[source]# Возвращает True, если пакет распределенных вычислений доступен. В противном случае torch.distributed не предоставляет никаких других API. В настоящее время torch.distributed доступен на Linux, MacOS и Windows. Установите USE_DISTRIBUTED=1, чтобы включить его при сборке PyTorch из исходного кода. В настоящее время значение по умолчанию — USE_DISTRIBUTED=1 для Linux и Windows, USE_DISTRIBUTED=0 для MacOS. Тип возвращаемого значения bool torch.distributed.init_process_group(backend=None, init_method=None, timeout=None, world_size=-1, rank=-1, store=None, group_name='', pg_options=None, device_id=None)[source]# Инициализирует группу процессов по умолчанию. Это также инициализирует пакет распределенных вычислений. Есть 2 основных способа инициализации группы процессов: Явно указать store, rank и world_size. Указать init_method (строку URL), которая указывает, где/как обнаруживать пиров. Опционально указать rank и world_size или закодировать все необходимые параметры в URL и опустить их. Если не указано ни то, ни другое, init_method считается равным "env://". Параметры backend (str или Backend, опционально) – Используемый бэкенд. В зависимости от конфигурации сборки допустимые значения включают mpi, gloo, nccl, ucc, xccl или зарегистрированные сторонним плагином. Начиная с версии 2.6, если backend не указан, c10d будет использовать бэкенд, зарегистрированный для типа устройства, указанного в kwargs device_id (если он предоставлен). Известные регистрации по умолчанию на сегодня: nccl для cuda, gloo для cpu, xccl для xpu. Если не указаны ни backend, ни device_id, c10d обнаружит ускоритель на машине времени выполнения и использует бэкенд, зарегистрированный для этого обнаруженного ускорителя (или cpu). Это поле может быть задано в виде строки в нижнем регистре (например, "gloo"), к которой также можно получить доступ через атрибуты Backend (например, Backend.GLOO). При использовании нескольких процессов на машине с бэкендом nccl каждый процесс должен иметь эксклюзивный доступ к каждому используемому GPU, поскольку совместное использование GPU между процессами может привести к взаимоблокировке или недопустимому использованию NCCL. Бэкенд ucc является экспериментальным. Бэкенд по умолчанию для устройства можно запросить с помощью get_default_backend_for_device(). 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, опционально) – Тайм-аут для операций, выполняемых в группе процессов. Значение по умолчанию — 10 минут для NCCL и 30 минут для других бэкендов. Это продолжительность, после которой коллективы будут асинхронно прерваны, и процесс завершится с ошибкой. Это сделано потому, что выполнение CUDA является асинхронным, и больше небезопасно продолжать выполнение пользовательского кода, поскольку неудачные асинхронные операции NCCL могут привести к тому, что последующие операции CUDA будут выполняться с поврежденными данными. Когда установлен TORCH_NCCL_BLOCKING_WAIT, процесс будет блокироваться и ждать этого тайм-аута. group_name (str, опционально, устарело) – Имя группы. Этот аргумент игнорируется. pg_options (ProcessGroupOptions, опционально) – Параметры группы процессов, указывающие, какие дополнительные параметры необходимо передать при создании конкретных групп процессов. На данный момент единственная поддерживаемая опция — ProcessGroupNCCL.Options для бэкенда nccl; is_high_priority_stream может быть указан, чтобы бэкенд nccl мог использовать потоки CUDA с высоким приоритетом, когда ожидают вычислительные ядра. Другие доступные опции для настройки nccl см. https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/api/types.html#ncclconfig-t device_id (torch.device | int, опционально) – одно конкретное устройство, с которым будет работать этот процесс, что позволяет выполнять оптимизацию, специфичную для бэкенда. В настоящее время это имеет два эффекта, только под NCCL: коммуникатор формируется немедленно (вызов ncclCommInit сразу, а не обычный ленивый вызов), и подгруппы будут использовать ncclCommSplit, когда это возможно, чтобы избежать ненужных накладных расходов на создание группы. Если вы хотите узнать об ошибке инициализации NCCL на раннем этапе, вы также можете использовать это поле. Если указано int, API предполагает, что будет использоваться тип ускорителя во время компиляции. Примечание Чтобы включить backend == Backend.MPI, PyTorch должен быть собран из исходного кода в системе, поддерживающей MPI. Примечание Поддержка нескольких бэкендов является экспериментальной. В настоящее время, если бэкенд не указан, будут созданы бэкенды gloo и nccl. Бэкенд gloo будет использоваться для коллективов с тензорами CPU, а бэкенд nccl — для коллективов с тензорами CUDA. Пользовательский бэкенд может быть указан путем передачи строки в формате "<device_type>:<backend_name>,<device_type>:<backend_name>", например "cpu:gloo,cuda:custom_backend". torch.distributed.device_mesh.init_device_mesh(device_type, mesh_shape, , mesh_dim_names=None, backend_override=None)[source]# Инициализирует DeviceMesh на основе параметров device_type, mesh_shape и mesh_dim_names. Это создает DeviceMesh с n-мерным массивом, где n — длина mesh_shape. Если указан mesh_dim_names, каждое измерение помечается как mesh_dim_names[i]. Примечание init_device_mesh следует модели программирования SPMD, что означает, что одна и та же программа Python PyTorch выполняется на всех процессах/рангах в кластере. Убедитесь, что mesh_shape (размеры nD массива, описывающего расположение устройств) идентичен во всех рангах. Несогласованный mesh_shape может привести к зависанию. Примечание Если группа процессов не найдена, init_device_mesh инициализирует группу/группы процессов, необходимые для распределенных коммуникаций, «за кулисами». Параметры device_type (str) – Тип устройства сетки. В настоящее время поддерживается: "cpu", "cuda/cuda-like", "xpu". Передача типа устройства с индексом GPU, например "cuda:0", не допускается. mesh_shape (Tuple[int]) – Кортеж, определяющий размеры многомерного массива, описывающего расположение устройств. mesh_dim_names (Tuple[str], опционально) – Кортеж имен измерений сетки для назначения каждому измерению многомерного массива, описывающего расположение устройств. Его длина должна совпадать с длиной mesh_shape. Каждая строка в mesh_dim_names должна быть уникальной. backend_override (Dict[int | str, tuple[str, Options] | str | Options], опционально) – Переопределения для некоторых или всех ProcessGroups, которые будут созданы для каждого измерения сетки. Каждый ключ может быть либо индексом измерения, либо его именем (если указан mesh_dim_names). Каждое значение может быть кортежем, содержащим имя бэкенда и его параметры, или только одним из этих двух компонентов (в этом случае другой будет установлен в значение по умолчанию). Возвращает Объект DeviceMesh, представляющий расположение устройств. Тип возвращаемого значения DeviceMesh Пример: >>> from torch.distributed.device_mesh import init_device_mesh >>> >>> mesh_1d = init_device_mesh("cuda", mesh_shape=(8,)) >>> mesh_2d = init_device_mesh("cuda", mesh_shape=(2, 8), mesh_dim_names=("dp", "tp")) torch.distributed.is_initialized()[source]# Проверяет, была ли инициализирована группа процессов по умолчанию. Тип возвращаемого значения bool torch.distributed.is_mpi_available()[source]# Проверяет, доступен ли бэкенд MPI. Тип возвращаемого значения bool torch.distributed.is_nccl_available()[source]# Проверяет, доступен ли бэкенд NCCL. Тип возвращаемого значения bool torch.distributed.is_gloo_available()[source]# Проверяет, доступен ли бэкенд Gloo. Тип возвращаемого значения bool torch.distributed.distributed_c10d.is_xccl_available()[source]# Проверяет, доступен ли бэкенд XCCL. Тип возвращаемого значения bool torch.distributed.is_torchelastic_launched()[source]# Проверяет, был ли этот процесс запущен с помощью torch.distributed.elastic (также известного как torchelastic). Наличие переменной окружения TORCHELASTIC_RUN_ID используется в качестве прокси для определения того, был ли текущий процесс запущен с помощью torchelastic. Это разумный прокси, поскольку TORCHELASTIC_RUN_ID соответствует идентификатору рандеву, который всегда является ненулевым значением, указывающим идентификатор задания для целей обнаружения пиров. Тип возвращаемого значения bool torch.distributed.get_default_backend_for_device(device)[source]# Возвращает бэкенд по умолчанию для данного устройства. Параметры device (Union[str, torch.device]) – Устройство, для которого нужно получить бэкенд по умолчанию. Возвращает Бэкенд по умолчанию для данного устройства в виде строки в нижнем регистре. Тип возвращаемого значения str В настоящее время поддерживаются три метода инициализации: TCP инициализация# Есть два способа инициализации с использованием TCP, оба требуют сетевого адреса, доступного из всех процессов, и желаемого world_size. Первый способ требует указания адреса, принадлежащего процессу с рангом 0. Этот метод инициализации требует, чтобы все процессы вручную указали ранги. Обратите внимание, что многоадресный адрес больше не поддерживается в последней версии пакета распределенных вычислений. group_name также устарел. import torch.distributed as dist # Используйте адрес одной из машин dist.init_process_group(backend, init_method='tcp://10.1.1.20:23456', rank=args.rank, world_size=4) Инициализация через общую файловую систему# Другой метод инициализации использует файловую систему, которая является общей и видна со всех машин в группе, вместе с желаемым world_size. URL должен начинаться с file:// и содержать путь к несуществующему файлу (в существующем каталоге) в общей файловой системе. Файловая инициализация автоматически создаст этот файл, если он не существует, но не удалит его. Поэтому вы несете ответственность за то, чтобы файл был очищен перед следующим вызовом init_process_group() с тем же путем/именем файла. Обратите внимание, что автоматическое назначение рангов больше не поддерживается в последней версии пакета распределенных вычислений, и group_name также устарел. Предупреждение Этот метод предполагает, что файловая система поддерживает блокировку с помощью fcntl — большинство локальных систем и NFS поддерживают ее. Предупреждение Этот метод всегда будет создавать файл и приложит все усилия для его очистки и удаления в конце программы. Другими словами, каждая инициализация с помощью файлового метода init будет нуждаться в совершенно новом пустом файле для успешной инициализации. Если тот же файл, использованный в предыдущей инициализации (который не был очищен), используется снова, это неожиданное поведение и может часто вызывать взаимоблокировки и сбои. Поэтому, даже если этот метод приложит все усилия для очистки файла, если автоудаление окажется неудачным, вы несете ответственность за то, чтобы файл был удален в конце обучения, чтобы предотвратить повторное использование того же файла в следующий раз. Это особенно важно, если вы планируете вызывать init_process_group() несколько раз с одним и тем же именем файла. Другими словами, если файл не удален/не очищен, и вы снова вызываете init_process_group() для этого файла, ожидаются сбои. Эмпирическое правило здесь заключается в том, чтобы убедиться, что файл не существует или пуст каждый раз при вызове init_process_group(). import torch.distributed as dist # ранг всегда должен быть указан dist.init_process_group(backend, init_method='file:///mnt/nfs/sharedfile', world_size=4, rank=args.rank) Инициализация через переменные окружения# Этот метод будет читать конфигурацию из переменных окружения, что позволяет полностью настроить способ получения информации. Устанавливаемые переменные: MASTER_PORT — обязательно; должен быть свободным портом на машине с рангом 0 MASTER_ADDR — обязательно (кроме ранга 0); адрес узла с рангом 0 WORLD_SIZE — обязательно; может быть установлен здесь или в вызове функции init RANK — обязательно; может быть установлен здесь или в вызове функции init Машина с рангом 0 будет использоваться для установки всех соединений. Это метод по умолчанию, то есть init_method не нужно указывать (или он может быть env://). Улучшение времени инициализации# TORCH_GLOO_LAZY_INIT — устанавливает соединения по требованию, а не использует полную сетку, что может значительно улучшить время инициализации для операций, отличных от all2all.
torch.distributed.init_process_group()
Шаблон 4: Пример:
>>> from torch.distributed.device_mesh import init_device_mesh
>>>
>>> mesh_1d = init_device_mesh("cuda", mesh_shape=(8,))
>>> mesh_2d = init_device_mesh("cuda", mesh_shape=(2, 8), mesh_dim_names=("dp", "tp"))
Шаблон 5: Группы# По умолчанию коллективы работают с группой по умолчанию (также называемой миром) и требуют, чтобы все процессы вошли в вызов распределенной функции. Однако некоторые рабочие нагрузки могут выиграть от более детальной коммуникации. Здесь в игру вступают распределенные группы. Функция new_group() может использоваться для создания новых групп с произвольными подмножествами всех процессов. Она возвращает непрозрачный дескриптор группы, который может быть передан в качестве аргумента group всем коллективам (коллективы — это распределенные функции для обмена информацией в определенных хорошо известных шаблонах программирования). torch.distributed.new_group(ranks=None, timeout=None, backend=None, pg_options=None, use_local_synchronization=False, group_desc=None, device_id=None)[source]# Создает новую распределенную группу. Эта функция требует, чтобы все процессы в основной группе (т.е. все процессы, являющиеся частью распределенного задания) вошли в эту функцию, даже если они не будут членами группы. Кроме того, группы должны создаваться в одном и том же порядке во всех процессах. Предупреждение Безопасное параллельное использование: При использовании нескольких групп процессов с бэкендом NCCL пользователь должен обеспечить глобально согласованный порядок выполнения коллективов между рангами. Если несколько потоков в процессе выдают коллективы, необходима явная синхронизация для обеспечения согласованного порядка. При использовании асинхронных вариантов API коммуникации torch.distributed возвращается объект work, и ядро коммуникации помещается в отдельный поток CUDA, что позволяет перекрывать коммуникацию и вычисления. После того как одна или несколько асинхронных операций были выданы в одной группе процессов, они должны быть синхронизированы с другими потоками cuda путем вызова work.wait() перед использованием другой группы процессов. См. Using multiple NCCL communicators concurrently <https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/communicators.html#using-multiple-nccl-communicators-concurrently> для получения дополнительной информации. Параметры ranks (list[int]) – Список рангов членов группы. Если None, будет установлено на все ранги. По умолчанию None. timeout (timedelta, опционально) – см. init_process_group для подробностей и значения по умолчанию. backend (str или Backend, опционально) – Используемый бэкенд. В зависимости от конфигурации сборки допустимыми значениями являются gloo и nccl. По умолчанию используется тот же бэкенд, что и у глобальной группы. Это поле должно быть задано в виде строки в нижнем регистре (например, "gloo"), к которой также можно получить доступ через атрибуты Backend (например, Backend.GLOO). Если передано None, будет использоваться бэкенд, соответствующий группе процессов по умолчанию. По умолчанию None. pg_options (ProcessGroupOptions, опционально) – Параметры группы процессов, указывающие, какие дополнительные параметры необходимо передать при создании конкретных групп процессов. Например, для бэкенда nccl может быть указан is_high_priority_stream, чтобы группа процессов могла использовать потоки cuda с высоким приоритетом. Другие доступные опции для настройки nccl см. https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/api/types.html#ncclconfig-tuse_local_synchronization (bool, опционально): выполнить локальный барьер группы в конце создания группы процессов. Это отличается тем, что ранги, не являющиеся членами, не должны вызывать API и не присоединяются к барьеру. group_desc (str, опционально) – строка для описания группы процессов. device_id (torch.device, опционально) – одно конкретное устройство для «привязки» этого процесса. Вызов new_group попытается немедленно инициализировать бэкенд коммуникации для устройства, если это поле задано. Возвращает Дескриптор распределенной группы, который может быть передан коллективным вызовам или GroupMember.NON_GROUP_MEMBER, если ранг не входит в ranks. N.B. use_local_synchronization не работает с MPI. N.B. Хотя use_local_synchronization=True может быть значительно быстрее в больших кластерах и небольших группах процессов, следует соблюдать осторожность, поскольку это изменяет поведение кластера: ранги, не являющиеся членами, не присоединяются к group barrier(). N.B. use_local_synchronization=True может привести к взаимоблокировкам, когда каждый ранг создает несколько перекрывающихся групп процессов. Чтобы избежать этого, убедитесь, что все ранги следуют одному и тому же глобальному порядку создания. torch.distributed.get_group_rank(group, global_rank)[source]# Преобразует глобальный ранг в ранг группы. global_rank должен быть частью group, иначе возникает RuntimeError. Параметры group (ProcessGroup) – ProcessGroup, в которой нужно найти относительный ранг. global_rank (int) – Глобальный ранг для запроса. Возвращает Ранг группы для global_rank относительно group. Тип возвращаемого значения int N.B. вызов этой функции для группы процессов по умолчанию возвращает идентичность. torch.distributed.get_global_rank(group, group_rank)[source]# Преобразует ранг группы в глобальный ранг. group_rank должен быть частью group, иначе возникает RuntimeError. Параметры group (ProcessGroup) – ProcessGroup, из которой нужно получить глобальный ранг. group_rank (int) – Ранг группы для запроса. Возвращает Глобальный ранг group_rank относительно group. Тип возвращаемого значения int N.B. вызов этой функции для группы процессов по умолчанию возвращает идентичность. torch.distributed.get_process_group_ranks(group)[source]# Получить все ранги, связанные с group. Параметры group (Optional[ProcessGroup]) – ProcessGroup, из которой нужно получить все ранги. Если None, будет использоваться группа процессов по умолчанию. Возвращает Список глобальных рангов, упорядоченных по рангу группы. Тип возвращаемого значения list[int]
new_group()
Шаблон 6: Предупреждение Безопасное параллельное использование: При использовании нескольких групп процессов с бэкендом пользователь NCCL должен обеспечить глобально согласованный порядок выполнения коллективов между рангами. Если несколько потоков в процессе выдают коллективы, необходима явная синхронизация для обеспечения согласованного порядка. При использовании асинхронных вариантов API связи torch.distributed возвращает объект работы, а источник связи размещается в отдельном потоке CUDA, что позволяет перекрывать связь и вычисления. После того, как одна или несколько асинхронных операций были выполнены в одной группе процессов, они должны быть синхронизированы с другими потоками cuda с помощью вызова work.wait() перед использованием другой группы процессов. См. Одновременное использование нескольких коммуникаторов NCCL <https://docs.nvidia.com/deeplearning/nccl/user-guide/docs/usage/communicators.html#using-multiple-nccl-communicators-concurrently> для получения дополнительной информации.
NCCL
Шаблон 7: Примечание. Если вы используете DistributedDataParallel в сочетании с Distributed RPC Framework, вам всегда следует использовать torch.distributed.autograd.backward() для расчета градиентов и torch.distributed.optim.DistributedOptimizer для оптимизации параметров. Пример: >>> импортировать torch.distributed.autograd как dist_autograd >>> из torch.nn.parallel импортировать DistributedDataParallel как DDP >>> импортировать torch >>> из torch import optim >>> из torch.distributed.optim import DistributedOptimizer >>> импортировать torch.distributed.rpc как rpc >>> из torch.distributed.rpc import RRef >>> >>> t1 = torch.rand((3, 3), require_grad=True) >>> t2 = torch.rand((3, 3), require_grad=True) >>> rref = rpc.remote("worker1", torch.add, args=(t1, t2)) >>> ddp_model = DDP(my_model) >>> >>> # Настройка оптимизатора >>>оптимизатор_парамс = [rref] >>> для параметра в ddp_model.parameters(): >>>Optimer_params.append(RRef(param)) >>> >>> dist_optim = DistributedOptimizer( >>> optim.SGD, >>>Optimer_params, >>> lr=0.05, >>> ) >>> >>> с dist_autograd.context() как context_id: >>> pred = ddp_model(rref.to_here()) >>> loss = loss_func(pred, target) >>> dist_autograd.backward(context_id, [loss]) >>> dist_optim.step(context_id)
torch.distributed.autograd.backward()
Шаблон 8: static_graph (bool) – если установлено True, DDP знает, что обучаемый граф является статическим. Статический граф означает: 1) Набор дополнительных и неиспользуемых параметров не изменяется в течение всего цикла обучения; в этом случае значения не имеет, установил ли пользователь find_unused_parameters = True или нет. 2) Способность обучения графа не меняется в течение всего цикла обучения (то есть нет потока управления, обусловленного итерациями). Когда static_graph установлен в True, DDP будет поддерживать случаи, которые нельзя было удерживать в прошлом: 1) Рекуррентные обратные проходы. 2) Многократное использование контрольной точки активации. 3) Контрольная точка активации, когда модель имеет неиспользуемые параметры. 4) параметры модели, которые движутся вне функции вперед. 5) Потенциальное улучшение производительности при наличии неиспользуемых параметров, так как DDP не будет искать график на каждой итерации для обнаружения неиспользуемых параметров, когда static_graph установлен в True. Для проверки вы можете установить static_graph в True, один из — Способ проверки данных регистрации ddp в конце вашего запроса модели обработки. Если ddp_logging_data.get("can_set_static_graph") == True, скорее всего, вы также можете установить static_graph = True. Пример Пример:::>>>> model_DDP = torch.nn.parallel.DistributedDataParallel(model) >>> # Цикл обучения >>>... >>> ddp_logging_data = model_DDP._get_ddp_logging_data() >>> static_graph = ddp_logging_data.get("can_set_static_graph")
True
Справочные файлы
Этот навык включает в себя всестороннюю документацию в references/:
- other.md - Другая документация
Используйте view для чтения конкретных справочных документов, когда требуется дополнительная информация.
Работа с этим навыком
Для начинающих
Используйте материалы Getting_started или учебные пособия для изучения фундаментальных концепций.
Для определенных функций
Используйте соответствующие категории справочного файла (API, руководства и т.д.) для получения подробной информации.
Для примера кода
Раздел «Краткий справочник» выше содержит распространенные шаблоны, собранные из официальной документации.
Ресурсы
ссылки/
Организованная документация, полученная из официальных источников. Эти файлы содержат: - Подробные объяснения - Примеры кода с аннотациями языка - Ссылки на оригинальную документацию - Оглавление для внешней навигации
скрипты/
Добавьте сюда вспомогательные скрипты для автоматизации задач.
активы/
Добавьте сюда шаблоны, заготовки или готовые проекты.
Примечания
- Этот навык был автоматически сгенерирован из официальной документации.
- Справочные файлы, сохраняющие структуру и образцы исходных документов.
- Примеры кода включают определение языка для лучшей синтаксиса подсветки.
- Шаблоны быстрого доступа к ресурсам из примеров использования в документации.
Обновление
Чтобы обновить эти навыки с обновленной документацией: 1. Повторно запустите скребок с той же конфигурацией. 2. Навык будет перестроен с последней информацией