Функция barrier
Функция distributed.barrier синхронизирует все процессы в распределенной группе, блокируя выполнение текущего процесса до тех пор, пока все процессы в группе не вызовут эту функцию. Это критически важно для координации распределенных операций, гарантируя, что все процессы достигли определенной точки выполнения перед продолжением работы.
Основные параметры функции:
-
group(необязательный) - группа процессов для синхронизации. По умолчанию используется группа по умолчанию (None). -
async_op(необязательный) - если установлен вTrue, функция возвращает объектWork, который позволяет выполнять асинхронную синхронизацию. По умолчаниюFalse. -
device_ids(необязательный) - список устройств для синхронизации, используется при работе с CUDA. По умолчаниюNone.
Синтаксис
torch.distributed.barrier(group=None, async_op=False, device_ids=None)
Пример
Рассмотрим простой пример синхронизации процессов с использованием barrier в распределенной среде:
import torch
import torch.distributed as dist
import os
# Инициализация распределенной среды
dist.init_process_group(backend='gloo', init_method='env://')
rank = dist.get_rank()
world_size = dist.get_world_size()
# Каждый процесс создает свой тензор
t = torch.tensor([rank * 10, rank * 10 + 1, rank * 10 + 2])
print(f"Process {rank}: tensor before sync - {t}")
# Синхронизация всех процессов
dist.barrier()
# После барьера все процессы продолжают работу
print(f"Process {rank}: all processes synced!")
dist.destroy_process_group()
Пример
Используем barrier для синхронизации перед выполнением операции all_reduce, чтобы гарантировать, что все процессы подготовили данные:
import torch
import torch.distributed as dist
import os
torch.manual_seed(0)
dist.init_process_group(backend='gloo', init_method='env://')
rank = dist.get_rank()
# Каждый процесс генерирует свои случайные данные
t = torch.randn(3)
print(f"Process {rank}: initial tensor - {t}")
# Синхронизация перед вычислениями
dist.barrier()
# Выполнение операции all_reduce после синхронизации
dist.all_reduce(t, op=dist.ReduceOp.SUM)
print(f"Process {rank}: reduced tensor - {t}")
dist.destroy_process_group()
Пример
Используем асинхронный режим barrier для неблокирующей синхронизации процессов:
import torch
import torch.distributed as dist
import os
import time
dist.init_process_group(backend='gloo', init_method='env://')
rank = dist.get_rank()
# Каждый процесс выполняет свою работу
if rank == 0:
print("Process 0: performing heavy computation...")
time.sleep(2)
else:
print(f"Process {rank}: waiting for process 0...")
# Асинхронный барьер
work = dist.barrier(async_op=True)
# Можно выполнять другие задачи во время ожидания
print(f"Process {rank}: doing other work while waiting...")
# Ожидание завершения барьера
work.wait()
print(f"Process {rank}: barrier complete, continuing work!")
dist.destroy_process_group()
Пример
Синхронизация через barrier при работе с CUDA устройствами для предотвращения конфликтов при доступе к GPU:
import torch
import torch.distributed as dist
import os
if torch.cuda.is_available():
device = torch.device(f'cuda:{dist.get_rank()}')
torch.cuda.set_device(device)
else:
device = torch.device('cpu')
dist.init_process_group(backend='gloo', init_method='env://')
rank = dist.get_rank()
# Каждый процесс создает тензор на своем устройстве
t = torch.tensor([rank, rank + 1, rank + 2]).to(device)
print(f"Process {rank} on {device}: tensor - {t}")
# Синхронизация с указанием устройств
device_ids = [device] if device.type == 'cuda' else None
dist.barrier(device_ids=device_ids)
print(f"Process {rank}: all devices synchronized!")
dist.destroy_process_group()
Смотрите также
-
функцию
init_process_group,
которая инициализирует распределенную среду для работы с процессами -
функцию
all_reduce,
которая выполняет редукцию данных между всеми процессами -
функцию
get_rank,
которая возвращает номер текущего процесса в группе -
функцию
destroy_process_group,
которая завершает работу распределенной группы