Функция irecv
Функция distributed.irecv выполняет асинхронный (неблокирующий) приём данных
от другого процесса в распределённой среде PyTorch. Первым параметром функция
принимает тензор tensor, в который будут записаны полученные данные.
Вторым параметром src передаётся ранг процесса-отправителя.
Третьим параметром можно указать группу процессов group, по умолчанию
используется группа по умолчанию. Функция возвращает объект Work,
который можно использовать для проверки завершения операции.
Синтаксис
torch.distributed.irecv(tensor, src, [group])
Пример
Базовый пример использования асинхронного приёма данных:
import torch
import torch.distributed as dist
torch.manual_seed(0)
dist.init_process_group("gloo", rank=0, world_size=2)
recv_tensor = torch.zeros(5)
work = dist.irecv(recv_tensor, src=1)
work.wait()
print(recv_tensor)
Результат выполнения кода (при условии, что процесс с рангом 1 отправил данные):
tensor([1., 2., 3., 4., 5.])
Пример
Пример с проверкой завершения операции без блокировки:
import torch
import torch.distributed as dist
dist.init_process_group("gloo", rank=0, world_size=2)
recv_tensor = torch.ones(3)
work = dist.irecv(recv_tensor, src=1)
while not work.is_completed():
pass
print(recv_tensor)
Результат выполнения кода (если процесс 1 отправил тензор [7, 8, 9]):
tensor([7., 8., 9.])
Пример
Использование irecv с пользовательской группой процессов:
import torch
import torch.distributed as dist
dist.init_process_group("gloo", rank=0, world_size=4)
ranks = [0, 1]
group = dist.new_group(ranks)
recv_tensor = torch.zeros(4)
work = dist.irecv(recv_tensor, src=1, group=group)
work.wait()
print(recv_tensor)
Результат выполнения кода:
tensor([1., 1., 1., 1.])
Смотрите также
-
функцию
isend,
которая выполняет асинхронную отправку данных -
функцию
send,
которая выполняет синхронную отправку данных -
функцию
recv,
которая выполняет синхронный приём данных -
функцию
init_process_group,
которая инициализирует распределённую среду