Функция distributed.send
Функция distributed.send предназначена для отправки тензора из текущего процесса указанному получателю в группе процессов.
Первый параметр - tensor, который нужно отправить.
Второй параметр - dst, целочисленный идентификатор ранга процесса-получателя.
Третий параметр - group (опционально), задаёт группу процессов для коммуникации.
Функция возвращает объект DistributedRequest, который поддерживает метод wait для синхронизации.
Синтаксис
torch.distributed.send(tensor, dst, [group])
Пример
Отправим тензор из процесса с рангом 0 в процесс с рангом 1:
import torch
import torch.distributed as dist
dist.init_process_group(backend='gloo', init_method='tcp://localhost:23456', rank=0, world_size=2)
t = torch.tensor([1, 2, 3, 4, 5])
req = dist.send(t, dst=1)
req.wait()
print("Tensor sent")
dist.destroy_process_group()
Результат выполнения кода:
"Tensor sent"
Пример
Пример работы функции в паре с приёмом на стороне получателя. Отправка тензора из процесса 0 в процесс 1:
import torch
import torch.distributed as dist
dist.init_process_group(backend='gloo', init_method='tcp://localhost:23456', rank=0, world_size=2)
t = torch.tensor([10, 20, 30, 40, 50])
req = dist.send(t, dst=1)
req.wait()
print(f"Process {dist.get_rank()} sent tensor: {t}")
dist.destroy_process_group()
Результат выполнения кода:
"Process 0 sent tensor: tensor([10, 20, 30, 40, 50])"
Пример
Использование с пользовательской группой процессов. Создадим подгруппу и отправим тензор внутри неё:
import torch
import torch.distributed as dist
dist.init_process_group(backend='gloo', init_method='tcp://localhost:23456', rank=0, world_size=4)
new_group = dist.new_group([0, 2, 3])
if dist.get_rank() == 0:
t = torch.tensor([100, 200, 300])
req = dist.send(t, dst=2, group=new_group)
req.wait()
print("Tensor sent to rank 2 in custom group")
dist.destroy_process_group()
Результат выполнения кода:
"Tensor sent to rank 2 in custom group"
Смотрите также
-
функцию
recv,
которая принимает тензор от другого процесса -
функцию
broadcast,
которая отправляет один и тот же тензор всем процессам в группе -
функцию
all_reduce,
которая применяет операцию редукции ко всем процессам -
функцию
barrier,
которая синхронизирует все процессы в группе