DistributedDataParallel
-
class torch.nn.parallel.DistributedDataParallel(module, device_ids=None, output_device=None, dim=0, broadcast_buffers=True, process_group=None, bucket_cap_mb=25, find_unused_parameters=False, check_reduction=False, gradient_as_bucket_view=False, static_graph=False, delay_all_reduce_named_params=None, param_to_hook_all_reduce=None, mixed_precision=None)[source] -
Реализует распределённое обучение с помощью метода данных, основанный на пакете
torch.distributedна уровне модуля.Этот контейнер обеспечивает параллельную обработку данных путём синхронизации градиентов между каждой копией модели. Устройства для синхронизации задаются входным параметром
process_group, который по умолчанию охватывает весь мир. Обратите внимание, чтоDistributedDataParallelне разбивает или иным образом не разделяет входные данные по участвующим графическим процессорам; пользователь отвечает за определение того, как это сделать, например, с помощьюDistributedSampler.См. также: Основы и Использование nn.parallel.DistributedDataParallel вместо multiprocessing или nn.DataParallel. Те же ограничения на вход, что и в
torch.nn.DataParallel, применяются.Для создания этого класса необходимо, чтобы
torch.distributedбыл уже инициализирован путём вызоваtorch.distributed.init_process_group().DistributedDataParallelдоказано, что значительно быстрее, чемtorch.nn.DataParallel, для обучения с использованием нескольких графических процессоров на одном узле.Для использования
DistributedDataParallelна узле с N графическими процессорами, необходимо запуститьNпроцессов, гарантируя, что каждый процесс работает исключительно с одним графическим процессором от 0 до N-1. Это можно сделать, установивCUDA_VISIBLE_DEVICESдля каждого процесса или вызвав:>>> torch.cuda.set_device(i)
где i изменяется от 0 до N-1. В каждом процессе вы должны выполнить следующие действия для создания этого модуля:
>>> torch.distributed.init_process_group( >>> backend='nccl', world_size=N, init_method='...' >>> ) >>> model = DistributedDataParallel(model, device_ids=[i], output_device=i)
Для запуска нескольких процессов на один узел можно использовать
torch.distributed.launchилиtorch.multiprocessing.spawn.Примечание
Для краткого введения во все функции, связанные с распределённым обучением, обратитесь к Обзору распределённого обучения PyTorch.
Примечание
DistributedDataParallelможет использоваться совместно сtorch.distributed.optim.ZeroRedundancyOptimizerдля уменьшения объёма памяти состояния оптимизатора на каждом процессе. Подробнее см. в рецепте ZeroRedundancyOptimizer.Примечание
ncclбекенд в настоящее время является самым быстрым и рекомендуемым бекендом при использовании графических процессоров. Это относится как к распределённому обучению на одном узле, так и к нескольким узлам.Примечание
Этот модуль также поддерживает распределённое обучение с промежуточной точностью. Это означает, что ваша модель может иметь различные типы параметров, такие как смешанные типы
fp16иfp32, а сокращение градиента на этих смешанных типах параметров будет работать нормально.Примечание
Если вы используете
torch.saveна одном процессе для сохранения модуля иtorch.loadна других процессах для его восстановления, убедитесь, чтоmap_locationправильно настроен для каждого процесса. Безmap_locationtorch.loadвосстановит модуль на устройствах, с которых он был сохранён.Примечание
Когда модель обучается на
Mузлах сbatch=N, градиент будетMраз меньше по сравнению с той же моделью, обученной на одном узле сbatch=M*N, если сумма потерь суммируется (а не усредняется, как обычно), по всем примерам в пакете (потому что градиенты между различными узлами усредняются). Вы должны учитывать это, когда хотите получить математически эквивалентный процесс обучения по сравнению с локальным вариантом обучения. Но в большинстве случаев вы можете рассматривать модель, обернутую DistributedDataParallel, модель, обернутую DataParallel, и обычную модель на одном графическом процессоре как одинаковые (например, использовать одну и ту же скорость обучения для эквивалентного размера пакета).Примечание
Параметры никогда не передаются между процессами. Модуль выполняет операцию all-reduce над градиентами и предполагает, что они будут изменены оптимизатором во всех процессах одинаковым образом. Буферы (например, статистические данные BatchNorm) передаются от модуля в процессе с номером ранга 0 ко всем другим репликам в системе в каждой итерации.
Примечание
Если вы используете DistributedDataParallel совместно с Распределённой RPC-системой, вы всегда должны использовать
torch.distributed.autograd.backward()для вычисления градиентов иtorch.distributed.optim.DistributedOptimizerдля оптимизации параметров.Пример:
>>> import torch.distributed.autograd as dist_autograd >>> from torch.nn.parallel import DistributedDataParallel as DDP >>> import torch >>> from torch import optim >>> from torch.distributed.optim import DistributedOptimizer >>> import torch.distributed.rpc as rpc >>> from torch.distributed.rpc import RRef >>> >>> t1 = torch.rand((3, 3), requires_grad=True) >>> t2 = torch.rand((3, 3), requires_grad=True) >>> rref = rpc.remote("worker1", torch.add, args=(t1, t2)) >>> ddp_model = DDP(my_model) >>> >>> # Setup optimizer >>> optimizer_params = [rref] >>> for param in ddp_model.parameters(): >>> optimizer_params.append(RRef(param)) >>> >>> dist_optim = DistributedOptimizer( >>> optim.SGD, >>> optimizer_params, >>> lr=0.05, >>> ) >>> >>> with dist_autograd.context() as context_id: >>> pred = ddp_model(rref.to_here()) >>> loss = loss_func(pred, target) >>> dist_autograd.backward(context_id, [loss]) >>> dist_optim.step(context_id)Примечание
DistributedDataParallel в настоящее время предлагает ограниченную поддержку отложенного вычисления градиентов с
torch.utils.checkpoint(). DDP будет работать как ожидается, если в модели нет неиспользуемых параметров и каждый слой проверяется не более одного раза (убедитесь, что вы не передаётеfind_unused_parameters=Trueв DDP). В настоящее время мы не поддерживаем случай, когда слой проверяется несколько раз или когда в проверенной модели есть неиспользуемые параметры.Примечание
Чтобы загрузить словарь состояния модели без DDP из модели DDP, необходимо применить
consume_prefix_in_state_dict_if_present()для удаления префикса «module.» в словаре состояния DDP перед загрузкой.Предупреждение
Конструктор, метод forward и дифференцирование выходных данных (или функции выходных данных этого модуля) являются точками синхронизации распределённых вычислений. Учтите это в случае, если разные процессы могут выполнять разные коды.
Предупреждение
Этот модуль предполагает, что все параметры зарегистрированы в модели к моменту её создания. После создания параметры не должны добавляться или удаляться.
Предупреждение
Этот модуль предполагает, что все параметры зарегистрированы в модели каждого процесса распределения находятся в одном и том же порядке. Сам модуль выполнит вычисление градиента
allreduceв обратном порядке зарегистрированных параметров модели. Другими словами, пользователи несут ответственность за обеспечение того, что у каждого распределённого процесса есть точно такая же модель, а значит и такой же порядок регистрации параметров.Предупреждение
Этот модуль позволяет параметрам с несмежными (non-rowmajor-contiguous) шагами. Например, ваша модель может содержать некоторые параметры, чей формат памяти
torch.memory_formatявляетсяtorch.contiguous_format, и другие, чей форматtorch.channels_last. Однако, соответствующие параметры в разных процессах должны иметь одинаковые шаги.Предупреждение
Этот модуль не работает с
torch.autograd.grad()(т.е. он будет работать только в том случае, если градиенты будут накапливаться в.gradатрибутах параметров).Предупреждение
Если вы планируете использовать этот модуль с
ncclбекендом илиglooбекендом (использующим Infiniband) совместно с DataLoader, использующим несколько рабочих потоков, измените метод запуска многопроцессорной обработки наforkserver(только Python 3) илиspawn. К сожалению, Gloo (использующий Infiniband) и NCCL2 не являются безопасными для форка, и у вас, скорее всего, возникнут тупики, если вы не измените это значение.Предупреждение
Никогда не пытайтесь изменить параметры вашей модели после обертывания вашей модели
DistributedDataParallel. Поскольку при обертывании вашей моделиDistributedDataParallel, конструкторDistributedDataParallelзарегистрирует дополнительные функции сокращения градиента для всех параметров самой модели в момент создания. Если вы измените параметры модели позже, функции сокращения градиента больше не будут соответствовать правильному набору параметров.Предупреждение
Использование
DistributedDataParallelсовместно с Распределённой RPC-системой экспериментально и может быть изменено.
- Параметры
-
- module (Модуль) – модуль, который нужно распараллелить
-
device_ids (список из целых чисел или torch.device) –
CUDA устройства. 1) Для модулей с одним устройством,
device_idsможет содержать ровно один идентификатор устройства, который представляет единственное CUDA устройство, где находится входной модуль, соответствующий этому процессу. В качестве альтернативы,device_idsтакже может бытьNone. 2) Для модулей с несколькими устройствами и модулей на CPU,device_idsдолжно бытьNone.Когда
device_idsравноNoneдля обоих случаев, как входные данные для прямого прохода, так и сам модуль должны быть размещены на правильном устройстве. (по умолчанию:None) -
output_device (целое число или torch.device) – Расположение устройства вывода для модулей CUDA с одним устройством. Для модулей с несколькими устройствами и модулей на CPU, оно должно быть
None, а сам модуль определяет местоположение вывода. (по умолчанию:device_ids[0]для модулей с одним устройством) -
broadcast_buffers (логическое значение) – Флаг, который включает синхронизацию (вещание) буферов модуля в начале функции
forward. (по умолчанию:True) -
process_group – Группа процессов, используемая для распределённого снижения данных. Если
None, будет использоваться стандартная группа процессов, которая создаётся функциейtorch.distributed.init_process_group(). (по умолчанию:None) -
bucket_cap_mb –
DistributedDataParallelбудут группировать параметры в несколько групп, чтобы снижение градиента каждой группы могло потенциально перекрываться с вычислением обратного прохода.bucket_cap_mbопределяет размер группы в мегабайтах (МБ). (по умолчанию: 25) -
find_unused_parameters (логическое значение) – Проходим по графу autograd от всех тензоров, содержащихся в возвращаемом значении функции обернутого модуля
forward. Параметры, которые не получают градиенты как часть этого графа, предварительно помечаются как готовые к снижению. Кроме того, параметры, которые могли быть использованы в функции обернутого модуляforward, но не входили в вычисление потерь, а следовательно, также не получали градиенты, предварительно помечаются как готовые к снижению. (по умолчанию:False) - check_reduction – Этот аргумент устарел.
-
gradient_as_bucket_view (логическое значение) – Если установлено в
True, градиенты будут представлениями, указывающими на разные смещения вallreduceгруппах коммуникации. Это может снизить максимальное использование памяти, где сохранённый размер памяти будет равен общему размеру градиентов. Кроме того, это исключает накладные расходы на копирование между градиентами иallreduceгруппами коммуникации. Когда градиенты являются представлениями,detach_()нельзя вызывать на градиентах. Если возникают такие ошибки, пожалуйста, исправьте их, обратившись к функцииzero_grad()вtorch/optim/optimizer.pyв качестве решения. Обратите внимание, что градиенты будут представлениями после первой итерации, поэтому максимальную экономию памяти следует проверять после первой итерации. -
static_graph (логическое значение) –
Если установлено в
True, DDP знает, что обученный граф является статическим. Статический граф означает, что 1) набор используемых и неиспользуемых параметров не изменится в течение всего цикла обучения; в этом случае неважно, установили ли пользователиfind_unused_parameters = Trueили нет. 2) Способ обучения графа не изменится в течение всего цикла обучения (то есть нет управления потоком, зависящего от итераций). Когда static_graph установлено вTrue, DDP будет поддерживать случаи, которые не поддерживались ранее: 1) Рекурсивный обратный проход. 2) Кэширование активаций несколько раз. 3) Кэширование активаций, когда модель имеет неиспользуемые параметры. 4) Параметры модели находятся вне функции forward. 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) >>> # Training loop >>> ... >>> ddp_logging_data = model_DDP._get_ddp_logging_data() >>> static_graph = ddp_logging_data.get("can_set_static_graph")
-
delay_all_reduce_named_params (список из кортежей из строки и torch.nn.Parameter) – список именованных параметров, снижение которых по всем процессам будет отложено, когда градиент параметра, указанного в
param_to_hook_all_reduce, будет готов. Другие аргументы DDP не применяются к именованным параметрам, указанным в этом аргументе, поскольку эти именованные параметры будут игнорироваться редуктором DDP. -
param_to_hook_all_reduce (torch.nn.Parameter) – параметр, который будет задействован для отложенного снижения по всем процессам параметров, указанных в
delay_all_reduce_named_params.
- Переменные
-
module (Модуль) – модуль, который нужно распараллелить.
Пример:
>>> torch.distributed.init_process_group(backend='nccl', world_size=4, init_method='...') >>> net = torch.nn.parallel.DistributedDataParallel(model)
-
join(divide_by_initial_world_size=True, enable=True, throw_on_early_termination=False)[source] -
Менеджер контекста, используемый совместно с экземпляром
torch.nn.parallel.DistributedDataParallel, чтобы иметь возможность обучаться с неравномерными входными данными на участвующих процессах.Этот менеджер контекста будет отслеживать уже подключенные процессы DDP и «затенять» прямые и обратные проходы, вставляя операции коллективной коммуникации, соответствующие операциям, созданным неподключенными процессами DDP. Это гарантирует, что каждый вызов коллективной операции имеет соответствующий вызов от уже подключенных процессов DDP, предотвращая зависания или ошибки, которые могут произойти при обучении с неравномерными входными данными на различных процессах. В качестве альтернативы, если флаг
throw_on_early_terminationзадан какTrue, все тренеры будут генерировать ошибку, как только один ранг исчерпает входные данные, что позволяет перехватить и обработать эти ошибки в соответствии с логикой приложения.После того, как все процессы DDP присоединятся, менеджер контекста разошлёт модель, соответствующую последнему присоединившемуся процессу, всем процессам, чтобы гарантировать, что модель одинаковая на всех процессах (что гарантируется DDP).
Чтобы использовать это для включения обучения с неравномерными входными данными на различных процессах, просто оберните этот менеджер контекста вокруг своего цикла обучения. Дополнительные изменения в модели или загрузке данных не требуются.
Предупреждение
Если модель или цикл обучения, вокруг которых обернут этот менеджер контекста, имеют дополнительные распределённые коллективные операции, такие как
SyncBatchNormв прямом проходе модели, то необходимо включить флагthrow_on_early_termination. Это связано с тем, что этот менеджер контекста не осведомлён о коллективной коммуникации, не относящейся к DDP. Этот флаг заставит все ранги выдать ошибку, когда любой ранг исчерпает входные данные, что позволит перехватить и восстановить эти ошибки на всех рангах.- Параметры
-
-
divide_by_initial_world_size (bool) – Если
True, градиенты будут разделены на начальноеworld_sizeDDP-обучения. ЕслиFalse, будет вычислена эффективная величина мира (количество рангов, которые еще не исчерпали свои входные данные), и градиенты будут разделены на неё во время allreduce. Установитеdivide_by_initial_world_size=True, чтобы гарантировать, что каждый образец входных данных, включая неравномерные входные данные, имеет одинаковый вес с точки зрения того, сколько они вносят вклад в глобальный градиент. Это достигается путём всегда деления градиента на начальнуюworld_sizeдаже при столкновении с неравномерными входными данными. Если вы установите это в значениеFalse, мы разделим градиент на оставшееся количество узлов. Это гарантирует соответствие обучению на меньшемworld_size, хотя это также означает, что неравномерные входные данные будут вносить больший вклад в глобальный градиент. Обычно вы захотите установить это в значениеTrueдля случаев, когда последние несколько входных данных вашей задачи обучения неравномерны. В крайних случаях, когда существует большая разница в количестве входных данных, установка этого значения вFalseможет дать лучшие результаты. -
enable (bool) – Включать или нет обнаружение неравномерных входных данных. Передайте
enable=Falseдля отключения в тех случаях, когда известно, что входные данные равномерны на участвующих процессах. По умолчаниюTrue. -
throw_on_early_termination (bool) – Бросить ошибку или продолжить обучение, когда по крайней мере один ранг исчерпал входные данные. Если
True, выбросит ошибку при достижении первой рангом конца данных. ЕслиFalse, продолжит обучение с меньшей эффективной величиной мира, пока все ранги не присоединятся. Обратите внимание, что если этот флаг задан, флагdivide_by_initial_world_sizeбудет проигнорирован. По умолчаниюFalse.
-
divide_by_initial_world_size (bool) – Если
Пример:
>>> import torch >>> import torch.distributed as dist >>> import os >>> import torch.multiprocessing as mp >>> import torch.nn as nn >>> # On each spawned worker >>> def worker(rank): >>> dist.init_process_group("nccl", rank=rank, world_size=2) >>> torch.cuda.set_device(rank) >>> model = nn.Linear(1, 1, bias=False).to(rank) >>> model = torch.nn.parallel.DistributedDataParallel( >>> model, device_ids=[rank], output_device=rank >>> ) >>> # Rank 1 gets one more input than rank 0. >>> inputs = [torch.tensor([1]).float() for _ in range(10 + rank)] >>> with model.join(): >>> for _ in range(5): >>> for inp in inputs: >>> loss = model(inp).sum() >>> loss.backward() >>> # Without the join() API, the below synchronization will hang >>> # blocking for rank 1's allreduce to complete. >>> torch.cuda.synchronize(device=rank)
-
join_hook(**kwargs)[source] -
Возвращает хук присоединения DDP, который позволяет обучать на неравномерных входных данных, затенением коллективных коммуникаций в прямых и обратных проходах.
- Параметры
-
kwargs (dict) – словарь, содержащий любые ключевые параметры для изменения поведения хука присоединения во время выполнения; все
Joinableэкземпляры, разделяющие один и тот же менеджер контекста присоединения, получают одинаковое значение дляkwargs.
- Хук поддерживает следующие ключевые параметры:
-
- divide_by_initial_world_size (bool, optional):
-
Если
True, то градиенты делятся на начальную величину мира, с которой был запущен DDP. ЕслиFalse, то градиенты делятся на эффективную величину мира (т. е. количество неподключенных процессов), что означает, что неравномерные входные данные вносят больший вклад в глобальный градиент. Обычно это должно быть установлено вTrueесли степень неравномерности невелика, но может быть установлено вFalseв крайних случаях для возможного улучшения результатов. По умолчаниюTrue.
-
no_sync()[source] -
Менеджер контекста для отключения синхронизации градиентов между процессами DDP. Внутри этого контекста градиенты будут накапливаться в переменных модуля, которые впоследствии будут синхронизированы в первом прямом-обратном проходе, выходящем из контекста.
Пример:
>>> ddp = torch.nn.parallel.DistributedDataParallel(model, pg) >>> with ddp.no_sync(): >>> for input in inputs: >>> ddp(input).backward() # no synchronization, accumulate grads >>> ddp(another_input).backward() # synchronize grads
Предупреждение
Прямой проход должен быть включён внутри менеджера контекста, иначе градиенты всё равно будут синхронизированы.
-
register_comm_hook(state, hook)[source] -
Регистрирует хук коммуникации, который является расширением, предоставляющим гибкий хук пользователям, где они могут указать, как DDP агрегирует градиенты на нескольких рабочих процессах.
Этот хук будет очень полезен для исследователей, чтобы опробовать новые идеи. Например, этот хук можно использовать для реализации нескольких алгоритмов, таких как GossipGrad и сжатие градиента, которые включают разные стратегии коммуникации для синхронизации параметров при выполнении обучения с распределёнными данными.
- Параметры
-
-
state (объект) –
Передаётся хуку для сохранения любой информации о состоянии во время процесса обучения. Примеры включают обратную связь об ошибках при сжатии градиентов, узлы для следующей коммуникации в GossipGrad и т. д.
Он хранится локально на каждом рабочем процессе и совместно используется всеми тензорами градиента на рабочем процессе.
-
hook (Callable) –
Вызываемый объект со следующей подписью:
hook(state: object, bucket: dist.GradBucket) -> torch.futures.Future[torch.Tensor]:Эта функция вызывается, как только ведро готово. Хук может выполнить любое необходимое обработку и вернуть будущее, указывающее на завершение любой асинхронной работы (например, allreduce). Если хук не выполняет никакой коммуникации, он всё равно должен вернуть завершённое будущее. Будущее должно содержать новое значение тензоров ведра градиентов. После того, как ведро готово, уменьшитель c10d вызовет этот хук и использует тензоры, возвращённые будущим, и скопирует градиенты в отдельные параметры. Обратите внимание, что тип возврата будущего должен быть единственным тензором.
Мы также предоставляем API под названием
get_futureдля получения будущего, связанного с завершениемc10d.ProcessGroup.Work.get_futureв настоящее время поддерживается для NCCL и также поддерживается для большинства операций на GLOO и MPI, за исключением операций «точка-точка» (send/recv).
-
Предупреждение
Тензоры ведра градиентов не будут предварительно разделены на world_size. Пользователь отвечает за деление на world_size в случае операций, таких как allreduce.
Предупреждение
Хук коммуникации DDP может быть зарегистрирован только один раз и должен быть зарегистрирован до вызова обратного прохода.
Предупреждение
Объект будущего, который возвращает хук, должен содержать единственный тензор, имеющий ту же форму, что и тензоры внутри ведра градиентов.
Предупреждение
get_futureAPI поддерживает бэкэнды NCCL и частично GLOO и MPI (нет поддержки операций «точка-точка», таких как send/recv), и вернётtorch.futures.Future.- Пример::
-
Ниже приведён пример хука «ничего не делать», который возвращает тот же тензор.
>>> def noop(state: object, bucket: dist.GradBucket) -> torch.futures.Future[torch.Tensor]: >>> fut = torch.futures.Future() >>> fut.set_result(bucket.buffer()) >>> return fut >>> ddp.register_comm_hook(state=None, hook=noop)
- Пример::
-
Ниже приведён пример алгоритма Parallel SGD, где градиенты кодируются до allreduce, а затем декодируются после allreduce.
>>> def encode_and_decode(state: object, bucket: dist.GradBucket) -> torch.futures.Future[torch.Tensor]: >>> encoded_tensor = encode(bucket.buffer()) # encode gradients >>> fut = torch.distributed.all_reduce(encoded_tensor).get_future() >>> # Define the then callback to decode. >>> def decode(fut): >>> decoded_tensor = decode(fut.value()[0]) # decode gradients >>> return decoded_tensor >>> return fut.then(decode) >>> ddp.register_comm_hook(state=None, hook=encode_and_decode)
-
© 2024, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://pytorch.org/docs/2.1/generated/torch.nn.parallel.DistributedDataParallel.html