Связки связи DDP
Связка связи DDP — это общий интерфейс для управления тем, как передавать градиенты между рабочими процессами, переопределяя стандартное allreduce в DistributedDataParallel. Предоставлено несколько встроенных связок связи, и пользователи могут легко применять любую из этих связок для оптимизации связи. Кроме того, интерфейс связки также может поддерживать пользовательские стратегии связи для более сложных случаев использования.
Как использовать связку связи?
Чтобы использовать связку связи, пользователю нужно просто позволить модели DDP зарегистрировать связку перед циклом обучения, как показано ниже.
torch.nn.parallel.DistributedDataParallel.register_comm_hook()
Что обрабатывает связка связи?
Связка связи предоставляет гибкий способ выполнить allreduce градиентов. Поэтому она в основном обрабатывает градиенты на каждой реплике перед allreduce, которые сгруппированы для увеличения перекрытия между связью и вычислениями. В частности, torch.distributed.GradBucket представляет собой пакет тензоров градиента, которые нужно allreduce.
-
class torch.distributed.GradBucket -
Этот класс в основном передаёт сплющеный тензор градиента (возвращаемый
buffer()) в связку связи DDP. Этот тензор можно далее разбить на список тензоров каждого параметра в этом пакете (возвращаемыйget_per_parameter_tensors()) для применения операций на уровне слоя.
-
torch.distributed.GradBucket.index(self: torch._C._distributed_c10d.GradBucket) → int -
Предупреждение
Поскольку пакеты перестраиваются после первой итерации, не следует полагаться на индексы в начале обучения.
- Возвращает:
-
Индекс пакета, хранящего градиенты нескольких смежных слоёв. Все градиенты сгруппированы.
-
torch.distributed.GradBucket.buffer(self: torch._C._distributed_c10d.GradBucket) → at::Tensor -
- Возвращает:
-
Сплющеный одномерный
torch.Tensorбуфер, который можно далее разбить на список тензоров каждого параметра в этом пакете.
-
torch.distributed.GradBucket.gradients(self: torch._C._distributed_c10d.GradBucket) → List[at::Tensor] -
- Возвращает:
-
Список
torch.Tensor. Каждый тензор в списке соответствует градиенту.
-
torch.distributed.GradBucket.is_last(self: torch._C._distributed_c10d.GradBucket) → bool -
- Возвращает:
-
Является ли этот пакет последним пакетом для allreduce в итерации. Это также означает, что этот пакет соответствует первым нескольким слоям в прямом проходе.
-
torch.distributed.GradBucket.set_buffer(self: torch._C._distributed_c10d.GradBucket, buffer: at::Tensor) → None -
Заменяет тензор в пакете на входной тензор-буфер.
-
torch.distributed.GradBucket.parameters(self: torch._C._distributed_c10d.GradBucket) → List[at::Tensor] -
- Возвращает:
-
Список
torch.Tensor. Каждый тензор в списке соответствует параметру модели.
Стандартные хуки связи
Стандартные хуки связи — это простые хуки без состояния, поэтому состояние ввода в register_comm_hook — это либо группа процессов, либо None. Вход bucket — это объект torch.distributed.GradBucket.
-
torch.distributed.algorithms.ddp_comm_hooks.default_hooks.allreduce_hook(process_group, bucket)[source] -
Этот хук связи DDP просто вызывает
allreduce, используя тензорыGradBucket. После агрегирования тензоров градиента по всем рабочим узлам его хук обратного вызоваthenвычисляет среднее значение и возвращает результат. Если пользователь зарегистрирует этот хук, результаты DDP должны быть такими же, как в случае, если хук не был зарегистрирован. Следовательно, это не изменит поведение DDP, и пользователь может использовать его в качестве ссылки или изменить этот хук для ведения полезной информации или любых других целей, не влияя на поведение DDP.- Пример::
-
>>> ddp_model.register_comm_hook(process_group, allreduce_hook)
-
torch.distributed.algorithms.ddp_comm_hooks.default_hooks.fp16_compress_hook(process_group, bucket)[source] -
Этот хук связи DDP реализует простой подход к сжатию градиента, который преобразует тензор
GradBucketв формат с плавающей запятой полуточной точности (torch.float16) и затем делит его на размер группы процессов. Он выполняет allreduce этих тензоров градиентаfloat16. После allreduce сжатых тензоров градиента, связанный хук обратного вызоваdecompressпреобразует его обратно в исходный тип данных (например,float32).- Пример::
-
>>> ddp_model.register_comm_hook(process_group, fp16_compress_hook)
-
torch.distributed.algorithms.ddp_comm_hooks.default_hooks.bf16_compress_hook(process_group, bucket)[source] -
Предупреждение: Этот API экспериментальный и требует версию NCCL, более позднюю чем 2.9.6.
Этот хук связи DDP реализует простой подход к сжатию градиента, который преобразует тензор
GradBucketв формат с плавающей запятой полуточной точности Brain floating point format (torch.bfloat16) и затем делит его на размер группы процессов. Он выполняет allreduce этих тензоров градиентаbfloat16. После allreduce сжатых тензоров градиента, связанный хук обратного вызоваdecompressпреобразует его обратно в исходный тип данных (например,float32).- Пример::
-
>>> ddp_model.register_comm_hook(process_group, bf16_compress_hook)
Кроме того, предоставляется обёртка для хука связи для поддержки fp16_compress_hook() или bf16_compress_hook() в качестве обёртки, которая может быть объединена с другими хуками связи.
-
torch.distributed.algorithms.ddp_comm_hooks.default_hooks.fp16_compress_wrapper(hook)[source] -
Эта обёртка преобразует входной тензор градиента заданного хука связи DDP в формат с плавающей запятой полуточной точности (
torch.float16), и преобразует полученный тензор заданного хука обратно в исходный тип данных, такой какfloat32.Следовательно,
fp16_compress_hookэквивалентноfp16_compress_wrapper(allreduce_hook).- Пример::
-
>>> state = PowerSGDState(process_group=process_group, matrix_approximation_rank=1, start_powerSGD_iter=10) >>> ddp_model.register_comm_hook(state, fp16_compress_wrapper(powerSGD_hook))
- Тип возвращаемого значения:
-
Callable[[Любой, GradBucket], Future[Тензор]]
-
torch.distributed.algorithms.ddp_comm_hooks.default_hooks.bf16_compress_wrapper(hook)[source] -
Предупреждение: Этот API экспериментальный и требует версию NCCL, более позднюю чем 2.9.6.
Эта обёртка преобразует входной тензор градиента заданного хука связи DDP в формат с плавающей запятой полуточной точности (
Brain floating point format <https://en.wikipedia.org/wiki/Bfloat16_floating-point_format> `_ (``torch.bfloat16`), и преобразует полученный тензор заданного хука обратно в исходный тип данных, такой какfloat32.Следовательно,
bf16_compress_hookэквивалентноbf16_compress_wrapper(allreduce_hook).- Пример::
-
>>> state = PowerSGDState(process_group=process_group, matrix_approximation_rank=1, start_powerSGD_iter=10) >>> ddp_model.register_comm_hook(state, bf16_compress_wrapper(powerSGD_hook))
- Тип возвращаемого значения:
-
Callable[[Любой, GradBucket], Future[Тензор]]
Связка PowerSGD для коммуникации
PowerSGD (Vogels et al., NeurIPS 2019) — алгоритм сжатия градиентов, который может обеспечить очень высокие скорости сжатия и ускорить распределённое обучение, ограниченное пропускной способностью. Этот алгоритм должен сохранять как некоторые гиперпараметры, так и внутреннее состояние. Поэтому связка PowerSGD для коммуникации — это состоятельная связка, и пользователь должен предоставить объект состояния, определённый ниже.
Состояние PowerSGD
-
class torch.distributed.algorithms.ddp_comm_hooks.powerSGD_hook.PowerSGDState(process_group, matrix_approximation_rank=1, start_powerSGD_iter=1000, min_compression_rate=2, use_error_feedback=True, warm_start=True, orthogonalization_epsilon=0, random_seed=0, compression_stats_logging_frequency=10000, batch_tensors_with_same_shape=False)[source] -
Хранит гиперпараметры алгоритма и внутреннее состояние для всех градиентов во время обучения. В частности,
matrix_approximation_rankиstart_powerSGD_iterявляются основными гиперпараметрами, которые должны быть настроены пользователем. Для повышения производительности рекомендуется удерживать бинарные гиперпараметрыuse_error_feedbackиwarm_startвключёнными.-
matrix_approximation_rankконтролирует размер сжатых тензоров низкого ранга, что определяет скорость сжатия. Чем ниже ранг, тем сильнее сжатие.1.1. Если
matrix_approximation_rankслишком низкий, для достижения качества полной модели потребуется больше шагов обучения или она вообще его не достигнет, что приведёт к потере точности.1.2. Увеличение
matrix_approximation_rankсущественно увеличивает вычислительные затраты на сжатие, и точность может больше не улучшиться за определённым порогомmatrix_approximation_rank.
Для настройки
matrix_approximation_rank, рекомендуется начинать с 1 и увеличивать с коэффициентом 2 (как экспоненциальный поиск по сетке, 1, 2, 4, …), пока не будет достигнута удовлетворительная точность. Как правило, используется небольшое значение 1-4. Для некоторых задач NLP (как показано в приложении D исходной статьи) это значение было увеличено до 32.-
start_powerSGD_iterоткладывает сжатие PowerSGD до шагаstart_powerSGD_iter, а обычное allreduce выполняется до шагаstart_powerSGD_iter. Эта гибридная схема обычного allreduce + PowerSGD может эффективно повысить точность, даже если используется относительно малое значениеmatrix_approximation_rank. Это связано с тем, что начальная фаза обучения обычно очень чувствительна к неточным градиентам, а сжатие градиентов слишком рано может заставить обучение быстро принять не оптимальную траекторию, что может привести к необратимому влиянию на точность.
Для настройки
start_powerSGD_iter, рекомендуется начинать с 10% от общего числа шагов обучения и увеличивать его до достижения удовлетворительной точности. Если в обучении есть этап разогрева,start_powerSGD_iterобычно не должен быть меньше числа шагов разогрева.-
min_compression_rate— минимальная требуемая скорость сжатия, когда слой сжимается. Из-за вычислительных накладных расходов на сжатие, тензор стоит сжимать только если можно получить существенную экономию пропускной способности, где(num_rows + num_cols) * matrix_approximation_rank * min_compression_rate < num_rows * num_cols. Если заданный пороговый уровень сжатия не может быть удовлетворён, тензор будет напрямую allreduce без сжатия.
Статистика сжатия регистрируется каждые
compression_stats_logging_frequencyитераций после начала сжатия PowerSGD.-
orthogonalization_epsilonможет быть очень малым значением (например, 1e-8), которое добавляется к каждому нормализованному столбцу матрицы на шаге ортогонализации, чтобы предотвратить ошибку деления на ноль, если какой-либо столбец имеет все нули. Если это уже можно предотвратить (например, с помощью нормализации по батчу), рекомендуется использовать эпсилон 0 для точности. -
batch_tensors_with_same_shapeуправляет тем, сжимать и распаковывать тензоры с одинаковой формой в пакетной операции, чтобы достичь большей параллельности. Обратите внимание, что вы также должны увеличить размер корзины (то есть аргументbucket_cap_mbв конструкторе DDP), чтобы больше тензоров с одинаковой формой появлялось в одной корзине, однако это может уменьшить перекрытие между вычислениями и коммуникацией и увеличить объём памяти из-за стекирования тензоров с одинаковой формой. Установите значениеTrueесли вычисления по сжатию/распаковке являются узким местом.
Предупреждение
Если включена обратная связь об ошибках или разогрев, минимальное разрешённое значение
start_powerSGD_iterв DDP равно 2. Это связано с тем, что существует другая внутренняя оптимизация, которая перестраивает корзины на итерации 1 в DDP, и это может вступать в конфликт с любым тензором, запомненным до процесса перестройки. -
Схемы PowerSGD
Предупреждение
PowerSGD обычно требует дополнительной памяти такого же размера, что и градиенты модели, для включения обратной связи об ошибках, что может компенсировать смещение сжатой связи и повысить точность.
Предупреждение
Схемы PowerSGD могут конфликтовать с пакетом автоматической смешанной точности Apex. Используйте вместо него родной пакет автоматической смешанной точности PyTorch.
-
torch.distributed.algorithms.ddp_comm_hooks.powerSGD_hook.powerSGD_hook(state, bucket)[source] -
Данная схема связи DDP реализует алгоритм сжатия градиентов PowerSGD, описанный в статье. После агрегирования тензоров градиентов на всех рабочих узлах, эта схема применяет сжатие следующим образом:
-
Представляет входной сглаженный тензор градиента размерности 1D как список тензоров на параметр и делит все тензоры на две группы:
1.1 Тензоры, которые должны быть сжаты перед allreduce, так как сжатие может обеспечить достаточную экономию полосы пропускания.
1.2 Остальные тензоры будут напрямую allreduce без сжатия, включая все векторные тензоры (для смещений).
-
Обрабатывает несжатые тензоры:
2.1 Выделяет непрерывную память для этих несжатых тензоров и выполняет allreduce всех несжатых тензоров как группы без сжатия;
2.2 Копирует отдельные несжатые тензоры из непрерывной памяти обратно в входной тензор.
-
Обрабатывает тензоры, которые должны быть сжаты с помощью сжатия PowerSGD:
3.1 Для каждого тензора M создаёт два тензора низкого ранга P и Q для разложения M, так что M = PQ^T, где Q инициализируется из стандартного нормального распределения и ортогонализируется;
3.2 Вычисляет каждый P в Ps, который равен MQ;
3.3 Выполняет allreduce Ps как группу;
3.4 Ортогонализирует каждый P в Ps;
3.5 Вычисляет каждый Q в Qs, который приблизительно равен M^TP;
3.6 Выполняет allreduce Qs как группу;
3.7 Вычисляет каждый M среди всех сжатых тензоров, который приблизительно равен PQ^T.
Обратите внимание, что эта схема связи применяет стандартный allreduce для первых
state.start_powerSGD_iterитераций. Это не только предоставляет пользователю больший контроль над компромиссом между ускорение и точность, но и помогает абстрагировать некоторую сложность внутренней оптимизации DDP для разработчиков будущих схем связи.- Параметры:
-
-
state (PowerSGDState) – Информация о состоянии для настройки коэффициента сжатия и поддержки обратной связи об ошибках, запуска с предыдущих данных и т. д. Для настройки параметров сжатия необходимо в основном настроить
matrix_approximation_rank,start_powerSGD_iterиmin_compression_rate. - bucket (dist.GradBucket) – Бакет, хранящий сглаженный тензор градиента 1D, который группирует несколько тензоров на переменную. Обратите внимание, что, поскольку схема связи DDP поддерживает только режим одного процесса и одного устройства, в этом бакете хранится ровно один тензор.
-
state (PowerSGDState) – Информация о состоянии для настройки коэффициента сжатия и поддержки обратной связи об ошибках, запуска с предыдущих данных и т. д. Для настройки параметров сжатия необходимо в основном настроить
- Возвращаемое значение:
-
Обработчик будущего связи, который обновляет градиенты на месте.
- Тип возвращаемого значения:
- Пример::
-
>>> state = PowerSGDState(process_group=process_group, matrix_approximation_rank=1, start_powerSGD_iter=10, min_compression_rate=0.5) >>> ddp_model.register_comm_hook(state, powerSGD_hook)
-
-
torch.distributed.algorithms.ddp_comm_hooks.powerSGD_hook.batched_powerSGD_hook(state, bucket)[source] -
Данная схема связи DDP реализует упрощённый алгоритм сжатия градиентов PowerSGD, описанный в статье. Этот вариант не сжимает градиенты по слоям, а вместо этого сжимает входной сглаженный тензор, который объединяет все градиенты. Поэтому он быстрее, чем
powerSGD_hook(), но обычно приводит к намного меньшей точности, еслиmatrix_approximation_rankравно 1.Предупреждение
Увеличение значения
matrix_approximation_rankздесь может не обязательно повысить точность, потому что объединение тензоров на параметр без выравнивания столбцов/строк может нарушить структуру низкого ранга. Поэтому пользователь должен всегда рассматриватьpowerSGD_hook()в первую очередь и рассматривать только этот вариант, когда удовлетворительная точность достигается, когдаmatrix_approximation_rankравно 1.После агрегирования тензоров градиентов на всех рабочих узлах, эта схема применяет сжатие следующим образом:
- Представляет входной сглаженный тензор градиента как квадратный тензор M с нулевыми заполнениями;
- Создаёт два тензора низкого ранга P и Q для разложения M, так что M = PQ^T, где Q инициализируется из стандартного нормального распределения и ортогонализируется;
- Вычисляет P, который равен MQ;
- Выполняет allreduce P;
- Ортогонализирует P;
- Вычисляет Q, который приблизительно равен M^TP;
- Выполняет allreduce Q;
- Вычисляет M, который приблизительно равен PQ^T.
- Обрезает входной тензор до исходной длины.
Обратите внимание, что эта схема связи применяет стандартный allreduce для первых
state.start_powerSGD_iterитераций. Это не только предоставляет пользователю больший контроль над компромиссом между ускорение и точность, но и помогает абстрагировать некоторую сложность внутренней оптимизации DDP для разработчиков будущих схем связи.- Параметры:
-
-
state (PowerSGDState) – Информация о состоянии для настройки коэффициента сжатия и поддержки обратной связи об ошибках, запуска с предыдущих данных и т. д. Для настройки параметров сжатия необходимо в основном настроить
matrix_approximation_rankиstart_powerSGD_iter. - bucket (dist.GradBucket) – Бакет, хранящий сглаженный тензор градиента 1D, который группирует несколько тензоров на переменную. Обратите внимание, что, поскольку схема связи DDP поддерживает только режим одного процесса и одного устройства, в этом бакете хранится ровно один тензор.
-
state (PowerSGDState) – Информация о состоянии для настройки коэффициента сжатия и поддержки обратной связи об ошибках, запуска с предыдущих данных и т. д. Для настройки параметров сжатия необходимо в основном настроить
- Возвращаемое значение:
-
Обработчик будущего связи, который обновляет градиенты на месте.
- Тип возвращаемого значения:
- Пример::
-
>>> state = PowerSGDState(process_group=process_group, matrix_approximation_rank=1) >>> ddp_model.register_comm_hook(state, batched_powerSGD_hook)
Отладка схем связи
Как следует из названия, отладка схем связи только используется для целей отладки и оптимизации производительности.
Предупреждение
Схемы отладки связи не обязательно выдают правильные результаты.
-
torch.distributed.algorithms.ddp_comm_hooks.debugging_hooks.noop_hook(_, bucket)[source] -
Данная схема связи DDP возвращает будущее, которое оборачивает входной объект, поэтому это операция бездействия, которая не требует расходов на связь.
Эта схема только должна использоваться для анализа возможностей allreduce оптимизации, а не для нормальной синхронизации градиентов. Например, если после регистрации этой схемы наблюдается ускорение обучения менее чем на 10%, это обычно означает, что allreduce не является узким местом для этого случая. Такая инструментация может быть особенно полезной, если трассировки GPU трудно получить или анализ трассировки усложняется такими факторами, как перекрытие allreduce и вычислений или десинхронизация между рангами.
- Пример::
-
>>> ddp_model.register_comm_hook(None, noop_hook)
Запись состояний коммуникационных хуков
Состоятельный коммуникационный хук может быть сохранен как часть сохранения модели, чтобы разрешить перезапуск трейнера. Чтобы сделать хук сериализуемым, должны быть определены __setstate__ и __getstate__.
Предупреждение
__getstate__ должен исключать несериализуемые атрибуты из возвращаемого словаря.
Предупреждение
__setstate__ должен должным образом инициализировать несериализуемые атрибуты, исключённые из предоставленного state.
PowerSGDState реализует __setstate__ и __getstate__ и может использоваться в качестве ссылки.
- classtorch.distributed.algorithms.ddp_comm_hooks.powerSGD_hook.PowerSGDState[source]
-
-
__getstate__()[source] -
Возвращает
Dict[str, Any], который будет сериализован и сохранён.process_groupне сериализуем и исключается из возвращаемого состояния.
-
__setstate__(state)[source] -
Принимает предоставленный
stateи извлекаетPowerSGDState.process_groupустанавливается по умолчанию.
-
Вот простой пример сохранения и загрузки состояния PowerSGD и хука.
import os
import sys
import tempfile
import torch
import torch.distributed as dist
import torch.nn as nn
import torch.optim as optim
from torch.distributed.algorithms.ddp_comm_hooks import powerSGD_hook as powerSGD
class SimpleModel(nn.Module):
def __init__(self):
super(SimpleModel, self).__init__()
self.fc1 = nn.Linear(24,24)
self.relu = nn.ReLU()
self.fc2 = nn.Linear(24,12)
def forward(self, x):
return self.fc2(self.relu(self.fc1(x)))
def setup(rank, world_size):
os.environ['MASTER_ADDR'] = 'localhost'
os.environ['MASTER_PORT'] = '12355'
# initialize the process group
dist.init_process_group("nccl", rank=rank, world_size=world_size)
def cleanup():
dist.destroy_process_group()
def run_demo(demo_fn, world_size):
mp.spawn(
demo_fn,
args=(world_size,),
nprocs=world_size,
join=True)
def demo_serialization(rank, world_size):
setup(rank, world_size)
CHECKPOINT = tempfile.gettempdir() + "/checkpoint.pt"
model = SimpleModel().to(rank)
ddp_model = DistributedDataParallel(model, device_ids=[rank])
powersgd_hook = powerSGD.powerSGD_hook
powersgd_state = powerSGD.PowerSGDState(process_group=None)
optimizer = optim.SGD(ddp_model.parameters(), lr=0.001)
ddp_model.register_comm_hook(powersgd_state, powersgd_hook)
state = {
'state_dict': ddp_model.state_dict(),
'comm_hook': hook,
'comm_hook_state': hook_state}
if rank == 0:
torch.save(state, CHECKPOINT)
dist.barrier()
map_location = {'cuda:%d' % 0: 'cuda:%d' % rank}
checkpoint = torch.load(CHECKPOINT, map_location=map_location)
ddp_model.load_state_dict(checkpoint['state_dict'])
powersgd_hook = checkpoint['comm_hook']
powersgd_state = checkpoint['comm_hook_state']
ddp_model.register_comm_hook(powersgd_state, powersgd_hook)
if rank == 0:
os.remove(CHECKPOINT)
cleanup()
if __name__ == "__main__":
n_gpus = torch.cuda.device_count()
assert n_gpus >= 2, f"Requires at least 2 GPUs to run, but got {n_gpus}"
world_size = n_gpus
run_demo(demo_serialization, world_size)
Благодарности
Большое спасибо автору статьи PowerSGD, Thijs Vogels, за рецензирование кода PowerSGD коммуникационного хука, а также за эксперименты сравнения, которые показывают, что производительность PowerSGD коммуникационного хука сопоставима с реализацией в оригинальной статье.
© 2024, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://pytorch.org/docs/1.13/ddp_comm_hooks.html