Функция distributed.reduce
Функция torch.distributed.reduce выполняет коллективную операцию редукции над тензором tensor на всех процессах в группе и сохраняет результат только на процессе с номером dst. Первым параметром передаётся тензор, который будет изменён на процессе-получателе. Вторым параметром указывается номер процесса-получателя (dst). Третьим параметром можно задать тип операции редукции (по умолчанию ReduceOp.SUM). Также доступны параметры group для указания пользовательской группы процессов и async_op для асинхронного выполнения.
Синтаксис
torch.distributed.reduce(tensor, dst, op=ReduceOp.SUM, group=None, async_op=False)
Пример
Давайте инициализируем группу процессов для двух узлов и выполним операцию суммирования на процессе с индексом 0:
import torch
import torch.distributed as dist
dist.init_process_group(backend="gloo")
rank = dist.get_rank()
t = torch.tensor([1, 2, 3, 4, 5])
if rank == 0:
print("Process 0 before reduce:", t)
dist.reduce(t, dst=0, op=dist.ReduceOp.SUM)
if rank == 0:
print("Process 0 after reduce:", t)
Результат выполнения кода на процессе 0:
"Process 0 before reduce: tensor([1, 2, 3, 4, 5])"
"Process 0 after reduce: tensor([2, 4, 6, 8, 10])"
Пример
Используем операцию умножения (ReduceOp.PRODUCT) для редукции на процессе с индексом 1:
import torch
import torch.distributed as dist
dist.init_process_group(backend="gloo")
rank = dist.get_rank()
t = torch.tensor([2, 3, 4])
if rank == 1:
print("Process 1 before reduce:", t)
dist.reduce(t, dst=1, op=dist.ReduceOp.PRODUCT)
if rank == 1:
print("Process 1 after reduce:", t)
Результат выполнения кода на процессе 1:
"Process 1 before reduce: tensor([2, 3, 4])"
"Process 1 after reduce: tensor([8, 27, 64])"
Пример
Выполним асинхронную редукцию с использованием параметра async_op и дождёмся её завершения:
import torch
import torch.distributed as dist
dist.init_process_group(backend="gloo")
rank = dist.get_rank()
t = torch.tensor([10, 20, 30])
work = dist.reduce(t, dst=0, op=dist.ReduceOp.SUM, async_op=True)
work.wait()
if rank == 0:
print("Result after async reduce:", t)
Результат выполнения кода на процессе 0:
"Result after async reduce: tensor([20, 40, 60])"
Смотрите также
-
функцию
all_reduce,
которая выполняет редукцию и рассылает результат всем процессам -
функцию
broadcast,
которая рассылает данные от одного процесса всем остальным -
перечисление
ReduceOp,
которое определяет возможные операции редукции (SUM, PRODUCT, MAX, MIN) -
функцию
init_process_group,
которая инициализирует группу процессов для коллективных операций