Функция new_group
Функция new_group создаёт новую подгруппу
(communication group) в рамках распределённой среды,
инициализированной с помощью init_process_group.
Группа определяет набор процессов, которые могут обмениваться данными
с помощью коллективных операций, таких как all_reduce
или broadcast.
Первым параметром функция принимает список рангов процессов,
которые войдут в подгруппу. Вторым параметром можно указать
бэкенд для этой группы (по умолчанию используется бэкенд
основной группы).
Синтаксис
torch.distributed.new_group(
ranks,
backend=None,
timeout=None,
pg_options=None
)
Пример
Давайте создадим группу, состоящую из процессов с рангами 0 и 1:
import torch
import torch.distributed as dist
dist.init_process_group(backend='gloo')
group = dist.new_group(ranks=[0, 1])
print(group)
Результат выполнения кода (объект группы):
<torch.distributed.distributed_c10d.ProcessGroupGloo object at 0x7f8a1c2d4a90>
Пример
Создадим группу из всех доступных процессов с явным указанием бэкенда:
import torch
import torch.distributed as dist
dist.init_process_group(backend='gloo')
world_size = dist.get_world_size()
ranks = list(range(world_size))
group = dist.new_group(ranks, backend='gloo')
print(f"Group created with {group.size()} processes")
Результат выполнения кода:
"Group created with 4 processes"
Пример
Используем созданную группу для операции all_reduce
только внутри подгруппы:
import torch
import torch.distributed as dist
dist.init_process_group(backend='gloo')
group = dist.new_group(ranks=[0, 2])
t = torch.tensor([1.0, 2.0])
if dist.get_rank() in [0, 2]:
dist.all_reduce(t, op=dist.ReduceOp.SUM, group=group)
print(f"Rank {dist.get_rank()}: {t}")
Результат выполнения кода (для рангов 0 и 2):
"Rank 0: tensor([2., 4.])"
"Rank 2: tensor([2., 4.])"
Смотрите также
-
функцию
init_process_group,
которая инициализирует распределённую среду -
функцию
all_reduce,
которая выполняет редукцию данных внутри группы -
функцию
get_rank,
которая возвращает ранг текущего процесса -
функцию
barrier,
которая синхронизирует процессы в группе