Spec-Zone.ru › PyTorch 1

Распределенная RPC-фреймворк

Распределенный RPC-фреймворк предоставляет механизмы для обучения моделей на нескольких машинах с помощью набора примитивов для удаленной связи и API более высокого уровня для автоматической дифференциации моделей, разделенных между несколькими машинами.

Предупреждение

API в пакете RPC стабильны. Ведется несколько работ по улучшению производительности и обработки ошибок, которые будут реализованы в будущих версиях.

Предупреждение

Поддержка CUDA была добавлена в PyTorch 1.9 и до сих пор является бета-функцией. Не все функции пакета RPC совместимы с поддержкой CUDA, и поэтому их использование не рекомендуется. К этим несовместимым функциям относятся: RRefs, совместимость JIT, dist autograd и dist optimizer, а также профилирование. Эти недостатки будут устранены в будущих выпусках.

Примечание

Для краткого введения во все функции, связанные с распределенным обучением, обратитесь к Обзору распределенного обучения PyTorch.

Основы

Распределенный RPC-фреймворк упрощает выполнение функций удаленно, поддерживает ссылки на удаленные объекты без копирования реальных данных и предоставляет API autograd и оптимизатора для прозрачного выполнения обратного прохода и обновления параметров через границы RPC. Эти функции можно разделить на четыре набора API.

  1. Удаленный вызов процедуры (RPC) поддерживает выполнение функции на указанном целевом рабочем узле с заданными аргументами и получение возвращаемого значения или создание ссылки на возвращаемое значение. Существует три основных API RPC: rpc_sync() (синхронный), rpc_async() (асинхронный) и remote() (асинхронный и возвращает ссылку на удаленное возвращаемое значение). Используйте синхронный API, если код пользователя не может продолжить без возвращаемого значения. В противном случае используйте асинхронный API, чтобы получить будущее, и подождите будущего, когда возвращаемое значение нужно на вызывающей стороне. API remote() полезен, когда требуется создать что-то удаленно, но никогда не нужно получать это на вызывающей стороне. Представьте случай, когда процесс драйвера настраивает сервер параметров и трейнер. Драйвер может создать таблицу встраивания на сервере параметров, а затем поделиться ссылкой на таблицу встраивания с трейнером, но сам никогда не будет использовать таблицу встраивания локально. В этом случае rpc_sync() и rpc_async() больше не подходят, так как они всегда подразумевают, что возвращаемое значение будет возвращено вызывающей стороне немедленно или в будущем.
  2. Удаленная ссылка (RRef) служит распределенным общим указателем на локальный или удаленный объект. Его можно обмениваться с другими рабочими узлами, и обработка подсчета ссылок будет выполняться прозрачно. У каждой RRef есть только один владелец, и объект существует только у этого владельца. Рабочие узлы, не являющиеся владельцами, имеющие RRefs, могут получить копии объекта от владельца, явно запросив их. Это полезно, когда рабочий узел нуждается в доступе к некоторому объекту данных, но сам по себе не является создателем (вызывающей стороной remote()) или владельцем объекта. Распределенный оптимизатор, как мы обсудим ниже, является одним из примеров таких случаев.
  3. Распределенный Autograd соединяет локальные движки автоградов на всех рабочих узлах, участвующих в прямом проходе, и автоматически обращается к ним во время обратного прохода для вычисления градиентов. Это особенно полезно, если прямой проход должен охватывать несколько машин при выполнении, например, распределенного обучения моделей, обучения с параметрическим сервером и т. д. С помощью этой функции код пользователя больше не должен беспокоиться о том, как отправлять градиенты через границы RPC и в каком порядке должны запускаться локальные движки автоградов, что может стать довольно сложным в случае вложенных и взаимозависимых вызовов RPC в прямом проходе.
  4. Распределенный оптимизатор конструктор принимает Optimizer() (например, SGD(), Adagrad() и т. д.) и список параметр RRefs, создает экземпляр Optimizer() на каждом отдельном владельце RRef и обновляет параметры соответственно при выполнении step(). Когда у вас есть распределенные прямой и обратный проходы, параметры и градиенты будут распределены по нескольким рабочим узлам, и поэтому требуется оптимизатор на каждом из вовлеченных рабочих узлов. Распределенный оптимизатор объединяет все эти локальные оптимизаторы в один и предоставляет краткий конструктор и API step().

RPC

Перед использованием RPC и распределенных примитивов autograd, необходимо произвести инициализацию. Для инициализации фреймворка RPC необходимо использовать init_rpc(), что инициализирует фреймворк RPC, фреймворк RRef и распределенный autograd.

torch.distributed.rpc.init_rpc(name, backend=None, rank=- 1, world_size=None, rpc_backend_options=None) [source]

Инициализирует RPC примитивы, такие как локальный агент RPC и распределенный autograd, что сразу делает текущий процесс готовым для отправки и получения RPC.

Параметры:
  • name (str) – глобально уникальное имя этого узла. (например, Trainer3, ParameterServer2, Master, Worker1). Имя может содержать только цифры, буквы, символы подчеркивания, двоеточия и/или дефисы, и должно быть короче 128 символов.
  • backend (BackendType, необязательно) – тип реализации бэкенда RPC. Поддерживаемые значения – BackendType.TENSORPIPE (по умолчанию). Подробнее см. Бэкенды.
  • rank (int) – глобально уникальный идентификатор/ранг этого узла.
  • world_size (int) – количество рабочих процессов в группе.
  • rpc_backend_options (RpcBackendOptions, необязательно) – параметры, передаваемые конструктору RpcAgent. Должно быть подклассом RpcBackendOptions, специфичным для агента, и содержит конфигурации инициализации, специфичные для агента. По умолчанию для всех агентов устанавливается значение таймаута по умолчанию в 60 секунд и выполняется установление связи с базовой группой процессов, инициализированной с помощью init_method = "env://", что означает, что переменные среды MASTER_ADDR и MASTER_PORT должны быть правильно установлены. Дополнительная информация и доступные параметры приведены в разделе Бэкенды.

Следующие API позволяют удаленно выполнять функции, а также создавать ссылки (RRef) на удаленные объекты данных. В этих API, при передаче Tensor в качестве аргумента или значения возврата, целевой рабочий процесс попытается создать Tensor с теми же метаданными (т.е. формой, шагом и т. д.). Мы намеренно запрещаем передачу тензоров CUDA, так как это может привести к сбою, если списки устройств на источнике и целевом рабочих процессах не совпадают. В таких случаях приложения всегда могут явно перемещать входные тензоры на ЦП вызывающего и перемещать их на нужные устройства на вызываемом, если необходимо.

Предупреждение

Поддержка TorchScript в RPC является прототипом и может быть изменена. Начиная с версии v1.5.0, torch.distributed.rpc поддерживает вызов функций TorchScript в качестве целевых функций RPC, что поможет улучшить параллельность на стороне вызываемого процесса, так как выполнение функций TorchScript не требует блокировки GIL.

torch.distributed.rpc.rpc_sync(to, func, args=None, kwargs=None, timeout=- 1.0) [source]

Выполняет блокирующий вызов RPC для выполнения функции func на рабочем процессе to. Сообщения RPC отправляются и принимаются параллельно с выполнением кода Python. Этот метод потокобезопасен.

Параметры:
  • to (str или WorkerInfo или int) – имя/ранг/WorkerInfo целевого рабочего процесса.
  • func (Callable) – вызываемая функция, такая как Python-функции, встроенные операторы (например, add()) и аннотированные функции TorchScript.
  • args (tuple) – кортеж аргументов для вызова func.
  • kwargs (dict) – словарь ключевых аргументов для вызова func.
  • timeout (float, необязательно) – таймаут в секундах для этого RPC. Если RPC не завершается в течение этого времени, будет поднято исключение, указывающее, что таймаут истек. Значение 0 указывает на бесконечный таймаут, т.е. ошибка таймаута никогда не будет поднята. Если не указано, используется значение по умолчанию, установленное при инициализации или с помощью _set_rpc_timeout.
Возвращает:

Возвращает результат выполнения func с args и kwargs.

Пример::

Убедитесь, что MASTER_ADDR и MASTER_PORT правильно установлены на обоих рабочих процессах. Обратитесь к API init_process_group() для получения дополнительной информации. Например,

export MASTER_ADDR=localhost export MASTER_PORT=5678

Затем выполните следующий код в двух разных процессах:

>>> # On worker 0:
>>> import torch
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker0", rank=0, world_size=2)
>>> ret = rpc.rpc_sync("worker1", torch.add, args=(torch.ones(2), 3))
>>> rpc.shutdown()
>>> # On worker 1:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker1", rank=1, world_size=2)
>>> rpc.shutdown()

Ниже приведен пример запуска функции TorchScript с помощью RPC.

>>> # On both workers:
>>> @torch.jit.script
>>> def my_script_add(t1, t2):
>>>    return torch.add(t1, t2)
>>> # On worker 0:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker0", rank=0, world_size=2)
>>> ret = rpc.rpc_sync("worker1", my_script_add, args=(torch.ones(2), 3))
>>> rpc.shutdown()
>>> # On worker 1:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker1", rank=1, world_size=2)
>>> rpc.shutdown()
torch.distributed.rpc.rpc_async(to, func, args=None, kwargs=None, timeout=- 1.0) [source]

Выполните асинхронный вызов RPC для выполнения функции func на работнике to. Сообщения RPC отправляются и принимаются параллельно выполнению Python-кода. Этот метод потокобезопасен. Этот метод сразу же вернёт Future, на котором можно дождаться завершения.

Параметры:
  • to (str или WorkerInfo или int) – имя/ранг/WorkerInfo целевого работника.
  • func (Callable) – вызываемая функция, например, Python-функции, встроенные операторы (например, add()) и аннотированные TorchScript-функции.
  • args (tuple) – кортеж аргументов для вызова func.
  • kwargs (dict) – словарь ключевых аргументов для вызова func.
  • timeout (float, необязательно) – таймаут в секундах для этого RPC. Если RPC не завершится в течение этого времени, будет поднято исключение, указывающее, что таймаут истек. Значение 0 означает бесконечный таймаут, т.е. исключение по таймауту никогда не будет поднято. Если не указано, используется значение по умолчанию, установленное во время инициализации или с помощью _set_rpc_timeout.
Возвращаемое значение:

Возвращает объект Future, на котором можно дождаться завершения. По завершении, возвращаемое значение func на args и kwargs можно получить из объекта Future.

Предупреждение

Использование тензоров GPU в качестве аргументов или возвращаемых значений func не поддерживается, так как мы не поддерживаем отправку тензоров GPU по сети. Вам необходимо явно скопировать тензоры GPU на CPU, прежде чем использовать их в качестве аргументов или возвращаемых значений func.

Предупреждение

API rpc_async не копирует хранилища аргументных тензоров до отправки их по сети, что может выполняться другим потоком в зависимости от типа бэкенда RPC. Вызывающий метод должен убедиться, что содержимое этих тензоров остаётся неизменным до завершения возвращаемого Future.

Пример::

Убедитесь, что MASTER_ADDR и MASTER_PORT правильно установлены на обоих работниках. Обратитесь к API init_process_group() для получения дополнительной информации. Например,

export MASTER_ADDR=localhost export MASTER_PORT=5678

Затем выполните следующий код в двух разных процессах:

>>> # On worker 0:
>>> import torch
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker0", rank=0, world_size=2)
>>> fut1 = rpc.rpc_async("worker1", torch.add, args=(torch.ones(2), 3))
>>> fut2 = rpc.rpc_async("worker1", min, args=(1, 2))
>>> result = fut1.wait() + fut2.wait()
>>> rpc.shutdown()
>>> # On worker 1:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker1", rank=1, world_size=2)
>>> rpc.shutdown()

Ниже приведен пример выполнения TorchScript-функции с использованием RPC.

>>> # On both workers:
>>> @torch.jit.script
>>> def my_script_add(t1, t2):
>>>    return torch.add(t1, t2)
>>> # On worker 0:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker0", rank=0, world_size=2)
>>> fut = rpc.rpc_async("worker1", my_script_add, args=(torch.ones(2), 3))
>>> ret = fut.wait()
>>> rpc.shutdown()
>>> # On worker 1:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker1", rank=1, world_size=2)
>>> rpc.shutdown()
torch.distributed.rpc.remote(to, func, args=None, kwargs=None, timeout=- 1.0) [source]

Выполните удалённый вызов для выполнения func на работнике to и сразу же верните RRef к результату. Работник to будет владельцем возвращённого RRef, а вызывающий remote – пользователем. Владелец управляет глобальным счётчиком ссылок своего RRef, и владелец RRef уничтожается только тогда, когда глобально нет живых ссылок на него.

Параметры:
  • to (str или WorkerInfo или int) – имя/ранг/WorkerInfo целевого работника.
  • func (Callable) – вызываемая функция, например, Python-функции, встроенные операторы (например, add()) и аннотированные TorchScript-функции.
  • args (tuple) – кортеж аргументов для вызова func.
  • kwargs (dict) – словарь ключевых аргументов для вызова func.
  • timeout (float, необязательно) – таймаут в секундах для этого удалённого вызова. Если создание этого RRef на работнике to не успешно обработано на этом работнике в течение этого таймаута, то при следующей попытке использования RRef (например, to_here()) будет поднят таймаут, указывающий на эту ошибку. Значение 0 означает бесконечный таймаут, т.е. исключение по таймауту никогда не будет поднято. Если не указано, используется значение по умолчанию, установленное во время инициализации или с помощью _set_rpc_timeout.
Возвращаемое значение:

Экземпляр пользовательского RRef к результату. Используйте блокирующую API torch.distributed.rpc.RRef.to_here() для получения значения результата локально.

Предупреждение

API remote не копирует хранилища аргументных тензоров до отправки их по сети, что может выполняться другим потоком в зависимости от типа бэкенда RPC. Вызывающий метод должен убедиться, что содержимое этих тензоров остаётся неизменным до подтверждения возвращаемого RRef владельцем, что можно проверить с помощью API torch.distributed.rpc.RRef.confirmed_by_owner().

Предупреждение

Обработка ошибок, таких как таймауты для API remote выполняется на основе наилучших усилий. Это означает, что когда удалённые вызовы, инициированные remote, терпят неудачу, например, из-за ошибки таймаута, мы используем подход обработки ошибок на основе наилучших усилий. Это означает, что ошибки обрабатываются и устанавливаются на возвращаемом RRef асинхронно. Если RRef не использовался приложением до этой обработки (например, to_here или вызов fork), то будущие использования RRef будут должным образом поднимать ошибки. Однако возможно, что пользовательское приложение будет использовать RRef до обработки ошибок. В этом случае ошибки могут не быть подняты, так как они ещё не были обработаны.

Пример:

Make sure that ``MASTER_ADDR`` and ``MASTER_PORT`` are set properly
on both workers. Refer to :meth:`~torch.distributed.init_process_group`
API for more details. For example,

export MASTER_ADDR=localhost
export MASTER_PORT=5678

Then run the following code in two different processes:

>>> # On worker 0:
>>> import torch
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker0", rank=0, world_size=2)
>>> rref1 = rpc.remote("worker1", torch.add, args=(torch.ones(2), 3))
>>> rref2 = rpc.remote("worker1", torch.add, args=(torch.ones(2), 1))
>>> x = rref1.to_here() + rref2.to_here()
>>> rpc.shutdown()

>>> # On worker 1:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker1", rank=1, world_size=2)
>>> rpc.shutdown()

Below is an example of running a TorchScript function using RPC.

>>> # On both workers:
>>> @torch.jit.script
>>> def my_script_add(t1, t2):
>>>    return torch.add(t1, t2)

>>> # On worker 0:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker0", rank=0, world_size=2)
>>> rref = rpc.remote("worker1", my_script_add, args=(torch.ones(2), 3))
>>> rref.to_here()
>>> rpc.shutdown()

>>> # On worker 1:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker1", rank=1, world_size=2)
>>> rpc.shutdown()
torch.distributed.rpc.get_worker_info(worker_name=None) [source]

Получение WorkerInfo данного работника по имени. Используйте этот WorkerInfo для избежания передачи дорогостоящей строки при каждом вызове.

Параметры:

worker_name (str) – строковое имя работника. Если None, возвращает идентификатор текущего работника. (по умолчанию None)

Возвращаемое значение:

Экземпляр WorkerInfo для данного worker_name или WorkerInfo текущего работника, если worker_name является None.

torch.distributed.rpc.shutdown(graceful=True, timeout=0) [source]

Выполните завершение агента RPC и уничтожьте агент RPC. Это останавливает локальный агент от принятия ожидающих запросов и завершает фреймворк RPC, завершая все потоки RPC. Если graceful=True, это заблокирует выполнение, пока все локальные и удалённые процессы RPC не достигнут этой функции и не подождут завершения всех ожидающих задач. В противном случае, если graceful=False, это локальное завершение, и оно не ожидает, пока другие процессы RPC достигнут этой функции.

Предупреждение

Для объектов Future, возвращаемых функцией rpc_async(), future.wait() не должно вызываться после shutdown().

Параметры:

graceful (bool) – Выполнять ли плавное завершение. Если True, это 1) ожидает, пока не появятся незавершенные системные сообщения для UserRRefs и удаляет их; 2) блокирует выполнение, пока все локальные и удалённые процессы RPC не достигнут этой функции и не подождут завершения всех ожидающих задач.

Пример::

Убедитесь, что MASTER_ADDR и MASTER_PORT правильно настроены на обоих узлах. Обратитесь к API init_process_group() для получения более подробной информации. Например,

export MASTER_ADDR=localhost export MASTER_PORT=5678

Затем выполните следующий код в двух разных процессах:

>>> # On worker 0:
>>> import torch
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker0", rank=0, world_size=2)
>>> # do some work
>>> result = rpc.rpc_sync("worker1", torch.add, args=(torch.ones(1), 1))
>>> # ready to shutdown
>>> rpc.shutdown()
>>> # On worker 1:
>>> import torch.distributed.rpc as rpc
>>> rpc.init_rpc("worker1", rank=1, world_size=2)
>>> # wait for worker 0 to finish work, and then shutdown.
>>> rpc.shutdown()
class torch.distributed.rpc.WorkerInfo

Структура, которая содержит информацию о рабочем узле в системе. Содержит имя и идентификатор рабочего узла. Этот класс не предназначен для непосредственного создания, а экземпляр можно получить с помощью get_worker_info(), а результат можно передавать в такие функции, как rpc_sync(), rpc_async(), remote() для избежания копирования строки при каждом вызове.

property id

Уникальный идентификатор для идентификации рабочего узла.

property name

Имя рабочего узла.

Пакет RPC также предоставляет декораторы, которые позволяют приложениям указать, как должна обрабатываться данная функция на стороне вызываемого.

torch.distributed.rpc.functions.async_execution(fn) [source]

Декоратор для функции, указывающий, что значение возврата функции гарантированно является объектом Future, и эта функция может выполняться асинхронно на вызываемом узле RPC. Более конкретно, вызываемый узел извлекает объект Future, возвращённый обернутой функцией, и устанавливает последующие шаги обработки в качестве обратного вызова для этого объекта Future. Установленный обратный вызов будет считывать значение из объекта Future при завершении и отправлять значение обратно в качестве ответа RPC. Это также означает, что возвращённый объект Future существует только на стороне вызываемого узла и никогда не передаётся через RPC. Этот декоратор полезен, когда выполнение обернутой функции (fn) необходимо приостановить и возобновить из-за, например, содержащихся в ней вызовов rpc_async() или ожидания других сигналов.

Примечание

Для включения асинхронного выполнения приложения должны передавать объект функции, возвращаемый этим декоратором, в API RPC. Если RPC обнаружит атрибуты, установленные этим декоратором, он узнает, что эта функция возвращает объект Future и будет обрабатывать это соответствующим образом. Однако это не означает, что этот декоратор должен быть внешним при определении функции. Например, при объединении с @staticmethod или @classmethod, @rpc.functions.async_execution должен быть внутренним декоратором, чтобы позволить целевой функции распознаваться как статическая или классовая функция. Эта целевая функция всё равно может выполняться асинхронно, потому что при доступе статический или классовый метод сохраняет атрибуты, установленные декоратором @rpc.functions.async_execution.

Пример::

Возвращаемый объект Future может поступать из rpc_async(), then() или конструктора Future. Приведённый ниже пример показывает непосредственное использование объекта Future, возвращённого функцией then().

>>> from torch.distributed import rpc
>>>
>>> # omitting setup and shutdown RPC
>>>
>>> # On all workers
>>> @rpc.functions.async_execution
>>> def async_add_chained(to, x, y, z):
>>>     # This function runs on "worker1" and returns immediately when
>>>     # the callback is installed through the `then(cb)` API. In the
>>>     # mean time, the `rpc_async` to "worker2" can run concurrently.
>>>     # When the return value of that `rpc_async` arrives at
>>>     # "worker1", "worker1" will run the lambda function accordingly
>>>     # and set the value for the previously returned `Future`, which
>>>     # will then trigger RPC to send the result back to "worker0".
>>>     return rpc.rpc_async(to, torch.add, args=(x, y)).then(
>>>         lambda fut: fut.wait() + z
>>>     )
>>>
>>> # On worker0
>>> ret = rpc.rpc_sync(
>>>     "worker1",
>>>     async_add_chained,
>>>     args=("worker2", torch.ones(2), 1, 1)
>>> )
>>> print(ret)  # prints tensor([3., 3.])

При объединении с декораторами TorchScript этот декоратор должен быть внешним.

>>> from torch import Tensor
>>> from torch.futures import Future
>>> from torch.distributed import rpc
>>>
>>> # omitting setup and shutdown RPC
>>>
>>> # On all workers
>>> @torch.jit.script
>>> def script_add(x: Tensor, y: Tensor) -> Tensor:
>>>     return x + y
>>>
>>> @rpc.functions.async_execution
>>> @torch.jit.script
>>> def async_add(to: str, x: Tensor, y: Tensor) -> Future[Tensor]:
>>>     return rpc.rpc_async(to, script_add, (x, y))
>>>
>>> # On worker0
>>> ret = rpc.rpc_sync(
>>>     "worker1",
>>>     async_add,
>>>     args=("worker2", torch.ones(2), 1)
>>> )
>>> print(ret)  # prints tensor([2., 2.])

При объединении со статическими или классовыми методами этот декоратор должен быть внутренним.

>>> from torch.distributed import rpc
>>>
>>> # omitting setup and shutdown RPC
>>>
>>> # On all workers
>>> class AsyncExecutionClass:
>>>
>>>     @staticmethod
>>>     @rpc.functions.async_execution
>>>     def static_async_add(to, x, y, z):
>>>         return rpc.rpc_async(to, torch.add, args=(x, y)).then(
>>>             lambda fut: fut.wait() + z
>>>         )
>>>
>>>     @classmethod
>>>     @rpc.functions.async_execution
>>>     def class_async_add(cls, to, x, y, z):
>>>         ret_fut = torch.futures.Future()
>>>         rpc.rpc_async(to, torch.add, args=(x, y)).then(
>>>             lambda fut: ret_fut.set_result(fut.wait() + z)
>>>         )
>>>         return ret_fut
>>>
>>>     @rpc.functions.async_execution
>>>     def bound_async_add(self, to, x, y, z):
>>>         return rpc.rpc_async(to, torch.add, args=(x, y)).then(
>>>             lambda fut: fut.wait() + z
>>>         )
>>>
>>> # On worker0
>>> ret = rpc.rpc_sync(
>>>     "worker1",
>>>     AsyncExecutionClass.static_async_add,
>>>     args=("worker2", torch.ones(2), 1, 2)
>>> )
>>> print(ret)  # prints tensor([4., 4.])
>>>
>>> ret = rpc.rpc_sync(
>>>     "worker1",
>>>     AsyncExecutionClass.class_async_add,
>>>     args=("worker2", torch.ones(2), 1, 2)
>>> )
>>> print(ret)  # prints tensor([4., 4.])

Этот декоратор также работает с помощниками RRef, т. е. torch.distributed.rpc.RRef.rpc_sync(), torch.distributed.rpc.RRef.rpc_async(), и torch.distributed.rpc.RRef.remote().

>>> from torch.distributed import rpc
>>>
>>> # reuse the AsyncExecutionClass class above
>>> rref = rpc.remote("worker1", AsyncExecutionClass)
>>> ret = rref.rpc_sync().static_async_add("worker2", torch.ones(2), 1, 2)
>>> print(ret)  # prints tensor([4., 4.])
>>>
>>> rref = rpc.remote("worker1", AsyncExecutionClass)
>>> ret = rref.rpc_async().static_async_add("worker2", torch.ones(2), 1, 2).wait()
>>> print(ret)  # prints tensor([4., 4.])
>>>
>>> rref = rpc.remote("worker1", AsyncExecutionClass)
>>> ret = rref.remote().static_async_add("worker2", torch.ones(2), 1, 2).to_here()
>>> print(ret)  # prints tensor([4., 4.])

Бэкэнды

Модуль RPC может использовать различные бэкэнды для выполнения связи между узлами. Бэкенд, который следует использовать, может быть указан в функции init_rpc(), передав определённое значение перечисления BackendType. Независимо от используемого бэкенда, остальная часть API RPC не изменится. Каждый бэкенд также определяет свой собственный подкласс класса RpcBackendOptions, экземпляр которого также может быть передан в функцию init_rpc() для настройки поведения бэкенда.

class torch.distributed.rpc.BackendType(value)

Перечисление доступных бэкэндов.

PyTorch поставляется с встроенным BackendType.TENSORPIPE бэкэндом. Дополнительные бэкэнды можно зарегистрировать, используя функцию register_backend().

class torch.distributed.rpc.RpcBackendOptions

Абстрактная структура, содержащая параметры, передаваемые бэкенду RPC. Экземпляр этого класса можно передавать в функцию init_rpc() для инициализации RPC со специфическими конфигурациями, такими как таймаут RPC и используемый init_method.

property init_method

URL, указывающий, как инициализировать группу процессов. По умолчанию env://

property rpc_timeout

Число с плавающей точкой, указывающее таймаут для всех вызовов RPC. Если вызов RPC не завершается в течение этого времени, он завершается с исключением, указывающим, что он истек.

Бэкенд TensorPipe

Агент TensorPipe, являющийся по умолчанию, использует библиотеку TensorPipe, которая предоставляет примитив для точечной коммуникации, специально предназначенный для машинного обучения и фундаментально устраняющий некоторые ограничения Gloo. По сравнению с Gloo, он обладает преимуществом асинхронности, что позволяет большому количеству операций передачи происходить одновременно, каждая со своей скоростью, без блокировки друг друга. Он будет открывать каналы только между парами узлов при необходимости, по требованию, и при сбое одного узла будут закрыты только его инцидентные каналы, в то время как все остальные будут продолжать работать нормально. Кроме того, он способен поддерживать несколько различных транспортов (TCP, конечно, но также и общую память, NVLink, InfiniBand, …) и может автоматически определять их доступность и договариваться о наилучшем транспорте для каждого канала.

Бэкенд TensorPipe был представлен в PyTorch v1.6 и активно развивается. В настоящее время он поддерживает только тензоры CPU, поддержка GPU появится в ближайшее время. Он поставляется с транспортным слоем на основе TCP, как и Gloo. Он также способен автоматически разбивать и мультиплексировать большие тензоры по нескольким сокетам и потокам, чтобы достичь очень высокой пропускной способности. Агент сможет самостоятельно выбрать лучший транспорт без какого-либо вмешательства.

Пример:

>>> import os
>>> from torch.distributed import rpc
>>> os.environ['MASTER_ADDR'] = 'localhost'
>>> os.environ['MASTER_PORT'] = '29500'
>>>
>>> rpc.init_rpc(
>>>     "worker1",
>>>     rank=0,
>>>     world_size=2,
>>>     rpc_backend_options=rpc.TensorPipeRpcBackendOptions(
>>>         num_worker_threads=8,
>>>         rpc_timeout=20 # 20 second timeout
>>>     )
>>> )
>>>
>>> # omitting init_rpc invocation on worker2
class torch.distributed.rpc.TensorPipeRpcBackendOptions(*, num_worker_threads=16, rpc_timeout=60.0, init_method='env://', device_maps=None, devices=None, _transports=None, _channels=None) [source]

Параметры бэкенда для TensorPipeAgent, полученные из RpcBackendOptions.

Параметры:
  • num_worker_threads (int, необязательно) – Количество потоков в пуле потоков, используемых TensorPipeAgent для выполнения запросов (по умолчанию: 16).
  • rpc_timeout (float, необязательно) – Время ожидания по умолчанию в секундах для запросов RPC (по умолчанию: 60 секунд). Если RPC не завершится в течение этого времени, будет возбуждено исключение, указывающее на это. Вызывающие стороны могут переопределить это время ожидания для отдельных RPC в rpc_sync() и rpc_async(), если необходимо.
  • init_method (str, необязательно) – URL для инициализации распределённого хранилища, используемого для встречи. Он принимает любое значение, допускаемое для того же аргумента функции init_process_group() (по умолчанию: env://).
  • device_maps (Dict[str, Dict], необязательно) – Сопоставления размещения устройств с этого узла до узла вызываемого. Ключ — имя узла вызываемого, а значение — словарь (Dict из int, str, или torch.device) который сопоставляет устройства этого узла с устройствами узла вызываемого. (по умолчанию: None)
  • devices (Список из int, str или torch.device, необязательно) – все локальные устройства CUDA, используемые агентом RPC. По умолчанию он будет инициализирован всеми локальными устройствами с его собственного device_maps и соответствующими устройствами с устройств коллег device_maps. При обработке запросов CUDA RPC агент правильно синхронизирует потоки CUDA для всех устройств в этом List.
property device_maps

Расположения отображения устройств.

property devices

Все устройства, используемые локальным агентом.

property init_method

URL, определяющий, как инициализировать группу процессов. По умолчанию env://

property num_worker_threads

Количество потоков в пуле потоков, используемых TensorPipeAgent для выполнения запросов.

property rpc_timeout

Значение с плавающей точкой, указывающее используемое время ожидания для всех RPC. Если RPC не завершается в течение этого интервала времени, оно завершается с исключением, указывающим на истечение времени ожидания.

set_device_map(to, device_map) [source]

Установите отображение устройств между каждой парой вызывающего и вызываемого RPC. Эта функция может вызываться несколько раз для поэтапного добавления конфигураций размещения устройств.

Параметры:
  • worker_name (str) – Имя вызываемого.
  • device_map (Словарь python:int, str, или torch.device) – Отображения размещения устройств с этого узла до узла вызываемого. Это отображение должно быть обратимым.
Пример::
>>> # both workers
>>> def add(x, y):
>>>     print(x)  # tensor([1., 1.], device='cuda:1')
>>>     return x + y, (x + y).to(2)
>>>
>>> # on worker 0
>>> options = TensorPipeRpcBackendOptions(
>>>     num_worker_threads=8,
>>>     device_maps={"worker1": {0: 1}}
>>>     # maps worker0's cuda:0 to worker1's cuda:1
>>> )
>>> options.set_device_map("worker1", {1: 2})
>>> # maps worker0's cuda:1 to worker1's cuda:2
>>>
>>> rpc.init_rpc(
>>>     "worker0",
>>>     rank=0,
>>>     world_size=2,
>>>     backend=rpc.BackendType.TENSORPIPE,
>>>     rpc_backend_options=options
>>> )
>>>
>>> x = torch.ones(2)
>>> rets = rpc.rpc_sync("worker1", add, args=(x.to(0), 1))
>>> # The first argument will be moved to cuda:1 on worker1. When
>>> # sending the return value back, it will follow the invert of
>>> # the device map, and hence will be moved back to cuda:0 and
>>> # cuda:1 on worker0
>>> print(rets[0])  # tensor([2., 2.], device='cuda:0')
>>> print(rets[1])  # tensor([2., 2.], device='cuda:1')
set_devices(devices) [source]

Установите локальные устройства, используемые агентом TensorPipe RPC. При обработке запросов CUDA RPC агент TensorPipe RPC правильно синхронизирует потоки CUDA для всех устройств в этом List.

Параметры:

devices (Список из python:int, str, или torch.device) – локальные устройства, используемые агентом TensorPipe RPC.

Примечание

Фреймворк RPC не автоматически повторно выполняет вызовы rpc_sync(), rpc_async() и remote(). Причина в том, что фреймворк RPC не может определить, является ли операция идемпотентной или нет, и безопасно ли её повторно выполнять. В результате, приложение должно обрабатывать ошибки и повторно выполнять операции при необходимости. Связь RPC основана на TCP, и в результате могут произойти сбои из-за сбоев сети или временных проблем с подключением к сети. В таких ситуациях приложение должно повторно выполнять операции с разумными задержками, чтобы не перегрузить сеть агрессивными повторениями.

RRef

Предупреждение

В настоящее время RRef не поддерживаются при использовании тензоров CUDA

RRef (Удаленная ссылка) — это ссылка на значение некоторого типа T (например, Tensor ) на удаленном обработчике. Этот обработчик поддерживает живучесть удаленного значения у владельца, но не подразумевает, что значение будет передано на локальный обработчик в будущем. RRef можно использовать в обучении на нескольких машинах, сохраняя ссылки на nn.Modules, которые существуют на других обработчиках, и вызывая соответствующие функции для получения или изменения их параметров во время обучения. Более подробную информацию см. в разделе Протокол удаленных ссылок.

class torch.distributed.rpc.RRef [source]

Дополнительная информация о RRef

  • Протокол удаленных ссылок
    • Обзор
    • Предположения
    • Жизненный цикл RRef
      • Обоснование проектирования
      • Реализация
    • Сценарии протокола
      • Пользовательский обмен RRef с владельцем в качестве возвращаемого значения
      • Пользовательский обмен RRef с владельцем в качестве аргумента
      • Владелец обменивает RRef с пользователем
      • Пользователь обменивается RRef с пользователем

RemoteModule

Предупреждение

RemoteModule в настоящее время не поддерживается при использовании тензоров CUDA

RemoteModule — это простой способ создать nn.Module удаленно на другом процессе. Фактический модуль находится на удаленном узле, но локальный узел имеет обработчик этого модуля и может вызывать этот модуль аналогично обычному nn.Module. Однако вызов влечет за собой вызовы RPC к удаленному узлу и может выполняться асинхронно, если это необходимо, с помощью дополнительных API, поддерживаемых RemoteModule.

class torch.distributed.nn.api.remote_module.RemoteModule(*args, **kwargs) [source]

Экземпляр RemoteModule может быть создан только после инициализации RPC. Он создает указанный пользователем модуль на указанном удаленном узле. Он ведет себя как обычный nn.Module , за исключением того, что метод forward выполняется на удаленном узле. Он обрабатывает запись autograd, чтобы обеспечить распространение градиентов обратного прохода обратно в соответствующий удаленный модуль.

Он генерирует два метода forward_async и forward на основе сигнатуры метода forward модуля module_cls. forward_async выполняется асинхронно и возвращает Future. Аргументы forward_async и forward такие же, как и аргументы метода forward модуля, возвращаемого методом module_cls.

Например, если module_cls возвращает экземпляр nn.Linear, имеющий метод forward с сигнатурой def forward(input: Tensor) -> Tensor:, сгенерированный RemoteModule будет иметь 2 метода с сигнатурами:

Параметры:
  • remote_device (str) – Устройство на целевом обработчике, на котором мы хотим разместить этот модуль. Формат должен быть «<имя_обработчика>/<устройство>», где поле устройства может быть проанализировано как тип torch.device. Например, «trainer0/cpu», «trainer0», «ps0/cuda:0». Кроме того, поле устройства может быть необязательным, и значение по умолчанию — «cpu».
  • module_cls (nn.Module) –

    Класс для удаленно создаваемого модуля. Например,

    >>> class MyModule(nn.Module):
    >>>     def forward(input):
    >>>         return input + 1
    >>>
    >>> module_cls = MyModule
    
  • args (Sequence, optional) – args, которые нужно передать в module_cls.
  • kwargs (Dict, optional) – kwargs, которые нужно передать в module_cls.
Возвращает:

Экземпляр удаленного модуля, который оборачивает Module , созданный пользователем с помощью module_cls. У него есть метод forward для блокирующего вызова и метод forward_async для асинхронного вызова, возвращающий будущие значения вызова forward на удаленной стороне.

Пример::

Запустите следующий код в двух разных процессах:

>>> # On worker 0:
>>> import torch
>>> import torch.distributed.rpc as rpc
>>> from torch import nn, Tensor
>>> from torch.distributed.nn.api.remote_module import RemoteModule
>>>
>>> rpc.init_rpc("worker0", rank=0, world_size=2)
>>> remote_linear_module = RemoteModule(
>>>     "worker1/cpu", nn.Linear, args=(20, 30),
>>> )
>>> input = torch.randn(128, 20)
>>> ret_fut = remote_linear_module.forward_async(input)
>>> ret = ret_fut.wait()
>>> rpc.shutdown()
>>> # On worker 1:
>>> import torch
>>> import torch.distributed.rpc as rpc
>>>
>>> rpc.init_rpc("worker1", rank=1, world_size=2)
>>> rpc.shutdown()

Кроме того, более практичный пример, который сочетается с DistributedDataParallel (DDP), можно найти в этом учебнике.

get_module_rref()

Возвращает RRef (RRef[nn.Module] ), указывающий на удаленный модуль.

Тип возвращаемого значения:

RRef[Модуль]

remote_parameters(recurse=True)

Возвращает список RRef, указывающих на параметры удаленного модуля. Это обычно используется совместно с DistributedOptimizer.

Параметры:

recurse (bool) – если True, то возвращает параметры удаленного модуля и всех его подмодулей. В противном случае возвращает только параметры, которые являются непосредственными членами удаленного модуля.

Возвращает:

Список RRef (List[RRef[nn.Parameter]]) на параметры удаленного модуля.

Тип возвращаемого значения:

Список[RRef[Параметр]]

Фреймворк распределенного Autograd

Предупреждение

Распределённый autograd в настоящее время не поддерживается при использовании тензоров CUDA

Этот модуль предоставляет фреймворк распределённого autograd на основе RPC, который можно использовать для таких задач, как обучение модели параллельно. Вкратце, приложения могут отправлять и получать тензоры записи градиента через RPC. В прямом проходе мы записываем, когда тензоры записи градиента отправляются через RPC, а в обратном проходе мы используем эту информацию для выполнения распределённого обратного прохода с помощью RPC. Более подробную информацию см. в проектирование распределённого Autograd.

torch.distributed.autograd.backward(context_id: int, roots: List[Tensor], retain_graph=False) → None

Инициализирует распределённый обратный проход с использованием предоставленных корней. В настоящее время реализуется алгоритм алгоритм FAST, который предполагает, что все сообщения RPC, отправленные в одном контексте распределённого autograd на всех узлах, будут частью графа autograd во время обратного прохода.

Мы используем предоставленные корни для обнаружения графа autograd и вычисления соответствующих зависимостей. Этот метод блокируется до завершения всего вычисления autograd.

Мы накапливаем градиенты в соответствующем torch.distributed.autograd.context на каждом из узлов. Контекст autograd для использования определяется по context_id, который передаётся при вызове torch.distributed.autograd.backward(). Если нет допустимого контекста autograd, соответствующего заданному идентификатору, мы выводим ошибку. Вы можете получить накопленные градиенты с помощью API get_gradients().

Параметры:
  • context_id (int) – Идентификатор контекста autograd, для которого необходимо получить градиенты.
  • roots (список) – Тензоры, которые представляют корни вычисления autograd. Все тензоры должны быть скалярами.
  • retain_graph (bool, необязательно) – Если False, граф, используемый для вычисления grad, будет освобождён. Обратите внимание, что в почти всех случаях установка этого параметра в True не нужна и часто может быть обойдена гораздо более эффективным способом. Обычно вам нужно установить его в True для многократного запуска обратного прохода.
Пример::
>>> import torch.distributed.autograd as dist_autograd
>>> with dist_autograd.context() as context_id:
>>>     pred = model.forward()
>>>     loss = loss_func(pred, loss)
>>>     dist_autograd.backward(context_id, loss)
class torch.distributed.autograd.context [source]

Объект контекста для обертывания прямых и обратных проходов при использовании распределённого autograd. context_id, сгенерированный в инструкции with, требуется для уникальной идентификации распределённого обратного прохода на всех узлах. Каждый узел сохраняет метаданные, связанные с этим context_id, что необходимо для правильного выполнения распределённого прохода autograd.

Пример::
>>> import torch.distributed.autograd as dist_autograd
>>> with dist_autograd.context() as context_id:
>>>   t1 = torch.rand((3, 3), requires_grad=True)
>>>   t2 = torch.rand((3, 3), requires_grad=True)
>>>   loss = rpc.rpc_sync("worker1", torch.add, args=(t1, t2)).sum()
>>>   dist_autograd.backward(context_id, [loss])
torch.distributed.autograd.get_gradients(context_id: int) → Dict[Tensor, Tensor]

Возвращает отображение из Tensor к соответствующему градиенту для этого Tensor, накопленному в предоставленном контексте, соответствующем данному context_id в рамках распределённого обратного прохода autograd.

Параметры:

context_id (int) – Идентификатор контекста autograd, для которого необходимо получить градиенты.

Возвращает:

Отображение, где ключ — это Tensor, а значение — соответствующий градиент для этого Tensor.

Пример::
>>> import torch.distributed.autograd as dist_autograd
>>> with dist_autograd.context() as context_id:
>>>     t1 = torch.rand((3, 3), requires_grad=True)
>>>     t2 = torch.rand((3, 3), requires_grad=True)
>>>     loss = t1 + t2
>>>     dist_autograd.backward(context_id, [loss.sum()])
>>>     grads = dist_autograd.get_gradients(context_id)
>>>     print(grads[t1])
>>>     print(grads[t2])

Дополнительная информация о RPC Autograd

  • Проектирование распределённого Autograd
    • Предыстория
    • Запись autograd во время прямого прохода
    • Контекст распределённого Autograd
    • Распределённый обратный проход
      • Вычисление зависимостей
      • Алгоритм FAST
      • Алгоритм SMART
    • Распределённый оптимизатор
    • Простой пример от начала до конца

Распределённый оптимизатор

См. страницу torch.distributed.optim для документации по распределённым оптимизаторам.

Примечания по проектированию

Заметки по проектированию распределённого autograd охватывают проектирование фреймворка распределённого autograd на основе RPC, полезного для таких приложений, как параллельное обучение модели.

  • Проектирование распределённого Autograd

Заметки по проектированию RRef охватывают проектирование протокола RRef (Удалённая ссылка), используемого для ссылки на значения на удалённых узлах в рамках фреймворка.

  • Протокол удалённой ссылки

Учебники

Учебники по RPC знакомят пользователей с фреймворком RPC, предоставляют несколько примеров приложений с использованием API torch.distributed.rpc и демонстрируют, как использовать профилировщик для профилирования RPC-задач.

  • Начало работы с распределённым фреймворком RPC
  • Реализация сервера параметров с использованием распределённого фреймворка RPC
  • Объединение Distributed DataParallel с распределённым фреймворком RPC (также охватывает RemoteModule)
  • Профилирование RPC-задач
  • Реализация пакетной обработки RPC
  • Распределённое параллельное конвейерное обучение

© 2024, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://pytorch.org/docs/1.13/rpc.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API