Функция isend
Функция distributed.isend выполняет асинхронную отправку тензора
указанному процессу в распределенной группе. В отличие от синхронной
функции send, данная функция не блокирует выполнение программы
до завершения отправки, а возвращает объект Work, который
можно использовать для проверки статуса операции.
Первым параметром передается тензор для отправки, вторым - ранг
процесса-получателя, третьим - группа процессов (по умолчанию
используется группа по умолчанию), четвертым - тег для идентификации
сообщения.
Синтаксис
torch.distributed.isend(tensor, dst, group=None, tag=0)
Пример
Давайте создадим простую распределенную среду с двумя процессами и выполним асинхронную отправку тензора из процесса с рангом 0 в процесс с рангом 1:
import torch
import torch.distributed as dist
# Инициализация процесса
dist.init_process_group(backend='gloo')
rank = dist.get_rank()
if rank == 0:
# Процесс-отправитель
t = torch.tensor([1, 2, 3, 4, 5])
work = dist.isend(t, dst=1, tag=0)
work.wait() # Ожидаем завершения отправки
print("Тензор отправлен")
elif rank == 1:
# Процесс-получатель
t = torch.zeros(5)
dist.recv(t, src=0, tag=0)
print(t)
dist.destroy_process_group()
Результат выполнения кода в процессе с рангом 1:
tensor([1., 2., 3., 4., 5.])
Пример
Асинхронная отправка позволяет выполнять другие операции до завершения передачи данных. В этом примере отправитель выполняет вычисления, пока данные передаются:
import torch
import torch.distributed as dist
import time
dist.init_process_group(backend='gloo')
rank = dist.get_rank()
if rank == 0:
# Создаем большой тензор для отправки
t = torch.ones(1000000)
work = dist.isend(t, dst=1, tag=0)
# Выполняем другие вычисления во время отправки
print("Начинаем вычисления...")
res = torch.sum(t) * 2
print(f"Результат вычислений: {res}")
# Ждем завершения отправки
work.wait()
print("Отправка завершена")
elif rank == 1:
t = torch.zeros(1000000)
dist.recv(t, src=0, tag=0)
print(f"Получен тензор размером {t.shape}")
dist.destroy_process_group()
Результат выполнения кода в процессе с рангом 0:
"Начинаем вычисления..."
"Результат вычислений: 2000000.0"
"Отправка завершена"
Пример
Использование пользовательской группы процессов для отправки с тегом для идентификации сообщения:
import torch
import torch.distributed as dist
dist.init_process_group(backend='gloo')
rank = dist.get_rank()
world_size = dist.get_world_size()
# Создаем подгруппу из первых двух процессов
group = dist.new_group([0, 1])
if rank == 0:
t = torch.tensor([10, 20, 30, 40, 50])
# Отправляем с тегом 100
work = dist.isend(t, dst=1, group=group, tag=100)
work.wait()
print("Тензор отправлен через подгруппу")
elif rank == 1:
t = torch.zeros(5)
dist.recv(t, src=0, group=group, tag=100)
print(t)
dist.destroy_process_group()
Результат выполнения кода в процессе с рангом 1:
tensor([10., 20., 30., 40., 50.])
Смотрите также
-
функцию
send,
которая выполняет синхронную отправку тензора -
функцию
irecv,
которая выполняет асинхронный прием тензора -
функцию
recv,
которая выполняет синхронный прием тензора -
функцию
init_process_group,
которая инициализирует распределенную группу процессов