Распределенная система вызовов удаленных процедур (RPC)
Распределенная система RPC предоставляет механизмы для обучения моделей на нескольких машинах с помощью набора примитивов для удаленной связи и API более высокого уровня для автоматической дифференциации моделей, разделенных между несколькими машинами.
Предупреждение
API в пакете RPC стабильны. Проводятся работы по улучшению производительности и обработки ошибок, которые будут включены в будущие версии.
Предупреждение
Поддержка CUDA была добавлена в PyTorch 1.9 и до сих пор является бета-функцией. Не все функции пакета RPC совместимы с поддержкой CUDA, поэтому их использование не рекомендуется. К этим несовместимым функциям относятся: RRefs, совместимость JIT, dist autograd и dist optimizer, и профилирование. Эти недостатки будут устранены в будущих версиях.
Примечание
Обратитесь к Обзору распределённого обучения PyTorch для краткого ознакомления со всеми функциями, связанными с распределённым обучением.
Основы
Распределенная система RPC упрощает выполнение функций удаленно, поддерживает ссылки на удалённые объекты без копирования данных, и предоставляет API для autograd и оптимизатора, чтобы прозрачно запускать обратное распространение и обновлять параметры через границы RPC. Эти функции можно разделить на четыре набора API.
-
Вызов удалённой процедуры (RPC) поддерживает выполнение функции на указанном рабочем узле с заданными аргументами и получение значения результата обратно или создание ссылки на значение результата. Существует три основных API для RPC:
rpc_sync()(синхронный),rpc_async()(асинхронный) иremote()(асинхронный и возвращает ссылку на удалённый результат). Используйте синхронный API, если код пользователя не может продолжить без значения результата. В противном случае используйте асинхронный API, чтобы получить будущее значение и подождите его, когда значение результата потребуется на вызывающей стороне. APIremote()полезен, когда требуется создать что-то удалённо, но не нужно получать это значение на вызывающей стороне. Представьте случай, когда процесс-драйвер настраивает сервер параметров и трейнер. Драйвер может создать таблицу встраивания на сервере параметров и затем поделиться ссылкой на таблицу встраивания с трейнером, но сам никогда не будет использовать таблицу встраивания локально. В этом случаеrpc_sync()иrpc_async()больше не подходят, так как они всегда подразумевают, что значение результата будет возвращено вызывающей стороне немедленно или в будущем. -
Ссылка на удалённый объект (RRef) служит как распределённый указатель на локальный или удалённый объект. Он может быть разделён с другими рабочими узлами, и управление счётчиком ссылок будет выполняться прозрачно. Каждая RRef имеет только одного владельца, и объект существует только у этого владельца. Рабочие узлы, не являющиеся владельцами, хранящие RRefs, могут получить копии объекта от владельца, явно запросив их. Это полезно, когда рабочий узел нуждается в доступе к объекту данных, но сам не является ни создателем (вызывающим
remote()), ни владельцем объекта. Распределённый оптимизатор, как мы обсудим ниже, является одним из примеров таких случаев. - Распределённое вычисление градиентов (Autograd) связывает локальные механизмы autograd на всех рабочих узлах, участвующих в прямом проходе, и автоматически обращается к ним во время обратного прохода для вычисления градиентов. Это особенно полезно, если прямой проход должен охватывать несколько машин при проведении, например, распределённого обучения моделей параллельно, обучения на сервере параметров и т. д. С помощью этой функции код пользователя больше не должен беспокоиться о том, как отправлять градиенты через границы RPC и в каком порядке следует запускать локальные механизмы autograd, что может стать довольно сложным, когда в прямом проходе есть вложенные и взаимозависимые вызовы RPC.
-
Распределённый оптимизатор конструктор принимает
Optimizer()(например,SGD(),Adagrad()и т. д.) и список RRefs параметров, создаёт экземплярOptimizer()на каждом уникальном владельце RRef и обновляет параметры соответствующим образом при выполненииstep(). Когда у вас есть распределённые прямые и обратные проходы, параметры и градиенты будут распределены по нескольким рабочим узлам, и поэтому требуется оптимизатор на каждом из участвующих рабочих узлов. Распределённый оптимизатор объединяет все эти локальные оптимизаторы в один и предоставляет краткий конструктор иstep()API.
RPC
Перед использованием примитивов RPC и распределённого вычисления градиентов (Autograd) необходимо выполнить инициализацию. Для инициализации системы RPC необходимо использовать init_rpc(), что инициализирует систему RPC, систему RRef и распределённое вычисление градиентов.
-
torch.distributed.rpc.init_rpc(name, backend=None, rank=-1, world_size=None, rpc_backend_options=None)[source] -
Инициализирует примитивы RPC, такие как локальный агент RPC и распределённое вычисление градиентов (Autograd), что сразу подготавливает текущий процесс для отправки и получения RPC.
- Parameters
-
-
name (str) – глобально уникальное имя этого узла. (например,
Trainer3,ParameterServer2,Master,Worker1) Имя может содержать только цифры, буквы, знак подчёркивания, двоеточие и/или дефис, и должно быть короче 128 символов. -
backend (BackendType, optional) – Тип реализации бэкенда RPC. Поддерживаемые значения это
BackendType.TENSORPIPE(по умолчанию). Подробнее о бэкендах см. Бэкенды. - rank (int) – глобально уникальный идентификатор/ранг этого узла.
- world_size (int) – Количество рабочих узлов в группе.
-
rpc_backend_options (RpcBackendOptions, optional) – Параметры, передаваемые конструктору RpcAgent. Должен быть агенто-специфическим подклассом
RpcBackendOptionsи содержит конфигурации инициализации, специфичные для агента. По умолчанию для всех агентов устанавливается значение таймаута по умолчанию в 60 секунд и выполняется соединение с помощью группы процессов нижележащего уровня, инициализированной с использованиемinit_method = "env://", что означает, что переменные средыMASTER_ADDRиMASTER_PORTдолжны быть установлены правильно. Подробнее см. Бэкенды.
-
name (str) – глобально уникальное имя этого узла. (например,
Следующие API позволяют пользователям удаленно выполнять функции и создавать ссылки (RRefs) на удалённые объекты данных. В этих API, при передаче Tensor в качестве аргумента или возвращаемого значения, целевой рабочий узел попытается создать Tensor с теми же метаданными (т. е. форма, шаг и т. д.). Мы намеренно не позволяем передавать тензоры CUDA, так как это может привести к сбою, если списки устройств на исходном и целевом рабочих узлах не совпадают. В таких случаях приложения могут всегда явно перемещать входные тензоры на CPU на вызывающей стороне и перемещать их на нужные устройства на вызываемом узле при необходимости.
Предупреждение
Поддержка 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.
-
to (str или WorkerInfo или int) – имя/ранг/
- Возвращает
-
Возвращает результат выполнения
funcсargsиkwargs.
- Пример::
-
Убедитесь, что
MASTER_ADDRиMASTER_PORTправильно настроены на обоих узлах. Подробнее см. APIinit_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.
-
to (str или WorkerInfo или int) – имя/ранг/
- Возвращает
-
Возвращает объект
Future, на котором можно дождаться завершения. По завершении можно получить значение, возвращаемоеfuncнаargsиkwargs, из объектаFuture.
Предупреждение
Использование тензоров GPU в качестве аргументов или возвращаемых значений
funcне поддерживается, так как мы не поддерживаем передачу тензоров GPU по сети. Вам нужно явно скопировать тензоры GPU на CPU, прежде чем использовать их в качестве аргументов или возвращаемых значенийfunc.Предупреждение
API
rpc_asyncне копирует хранилища аргументных тензоров до отправки их по сети, что может выполняться другим потоком в зависимости от типа бэкенда RPC. Вызывающий код должен убедиться, что содержимое этих тензоров остается неизменным до завершения возвращаемогоFuture.- Пример::
-
Убедитесь, что
MASTER_ADDRиMASTER_PORTправильно настроены на обоих узлах. Подробнее см. APIinit_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.
-
to (str или WorkerInfo или int) – имя/ранг/
- Возвращает
-
Экземпляр пользовательского
RRefобъекта результата. Используйте блокирующую APItorch.distributed.rpc.RRef.to_here()для получения значения результата локально.
Предупреждение
API
remoteне копирует хранилища тензоров аргументов до их отправки по сети, что может выполняться другим потоком в зависимости от типа бэкенда RPC. Вызывающий код должен убедиться, что содержимое этих тензоров сохраняется до тех пор, пока возвращенный RRef не будет подтвержден владельцем, что можно проверить с помощью APItorch.distributed.rpc.RRef.confirmed_by_owner().Предупреждение
Ошибки, такие как таймауты для API
remoteобрабатываются на основе наилучших усилий. Это означает, что при сбое удаленных вызовов, инициированныхremote, например, с ошибкой таймаута, мы используем подход, основанный на наилучших усилиях, для обработки ошибок. Это означает, что ошибки обрабатываются и устанавливаются на результирующем RRef асинхронно. Если RRef не использовался приложением до обработки (например, при вызовеto_hereили вилки), будущие использования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, вернуть id текущего рабочего узла. (значение по умолчанию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установлены корректно на обоих рабочих узлах. Обратитесь к APIinit_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 также предоставляет декораторы, которые позволяют приложениям указывать, как должна обрабатываться данная функция на стороне вызываемого объекта.
END_OF_DOCUMENT_MARKER-
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 (Словарь[str, Словарь], необязательно) – Сопоставления расположения устройств от этого работника к вызывающему. Ключ — имя вызываемого работника, а значение — словарь (
Dictизint,str, илиtorch.device). Он сопоставляет устройства этого работника с устройствами вызываемого работника. (по умолчанию:None). -
devices (Список[int, str, или
torch.device], необязательно) – все локальные CUDA-устройства, используемые агентом RPC. По умолчанию он инициализируется всеми локальными устройствами из его собственногоdevice_mapsи соответствующими устройствами из устройств коллегdevice_maps. При обработке запросов CUDA RPC агент должным образом синхронизирует CUDA-потоки для всех устройств в этомList.
-
num_worker_threads (int, необязательно) – Количество потоков в пуле потоков, используемых
-
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. Эту функцию можно вызывать несколько раз, чтобы постепенно добавлять конфигурации расположения устройств.
- Параметры
-
- to (str) – Имя вызываемого.
- device_map (Словарь из 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 (Список из int, str, или torch.device) – локальные устройства, используемые агентом TensorPipe RPC.
Примечание
Фреймворк RPC не автоматически повторяет вызовы rpc_sync(), rpc_async() и remote(). Причина в том, что фреймворк RPC не может определить, является ли операция идемпотентной или нет, и безопасно ли её повторять. В результате, приложение должно обрабатывать ошибки и повторять попытки, если необходимо. Связь RPC основана на TCP, и в результате могут произойти ошибки из-за сбоев сети или прерывистых проблем с подключением к сети. В таких сценариях приложение должно повторить попытки с разумными задержками, чтобы обеспечить, что сеть не перегружена агрессивными повторами.
RRef
Предупреждение
RRef в настоящее время не поддерживаются при использовании CUDA тензоров
Ссылка на удалённую ссылку (Remote REFerence) — это ссылка на значение некоторого типа T (например, Tensor) на удалённом работнике. Этот маркер поддерживает вызываемое удалённое значение в жизни владельца, но нет никаких предположений, что значение будет перенесено на локального работника в будущем. RRefs могут использоваться в многомашинном обучении путём хранения ссылок на nn.Modules, которые существуют на других работниках, и вызова соответствующих функций для получения или изменения их параметров во время обучения. См. Протокол удалённой ссылки для получения дополнительной информации.
-
class torch.distributed.rpc.PyRRef(RRef) -
Класс, инкапсулирующий ссылку на значение некоторого типа на удалённом узле. Эта ссылка удерживает значение на удалённом узле живым. Ссылка
UserRRefбудет удалена, когда 1) к ней нет ссылок ни в коде приложения, ни в локальном контексте RRef, или 2) приложение вызвало плавную остановку. Вызов методов удалённой ссылки приведёт к неопределённому поведению. Реализация RRef предлагает только обнаружение ошибок с наилучшими усилиями, и приложения не должны использоватьUserRRefsпослеrpc.shutdown().Предупреждение
RRef можно сериализовать и десериализовать только с помощью модуля RPC. Сериализация и десериализация RRef без RPC (например, Python pickle, torch
save()/load(), JITsave()/load()и т.д.) приведёт к ошибкам.- Параметры
-
- value (object) – Значение, которое будет обернуто в эту RRef.
-
type_hint (Type, optional) – Тип Python, который должен быть передан компилятору
TorchScriptв качестве подсказки типа дляvalue.
- Пример::
-
Следующие примеры пропускают код инициализации и завершения RPC для простоты. Обратитесь к документации RPC для получения подробностей.
- Создание RRef с помощью rpc.remote
>>> import torch >>> import torch.distributed.rpc as rpc >>> rref = rpc.remote("worker1", torch.add, args=(torch.ones(2), 3)) >>> # get a copy of value from the RRef >>> x = rref.to_here()- Создание RRef из локального объекта
>>> import torch >>> from torch.distributed.rpc import RRef >>> x = torch.zeros(2, 2) >>> rref = RRef(x)
- Обмен RRef с другими узлами
>>> # On both worker0 and worker1: >>> def f(rref): >>> return rref.to_here() + 1
>>> # On worker0: >>> import torch >>> import torch.distributed.rpc as rpc >>> from torch.distributed.rpc import RRef >>> rref = RRef(torch.zeros(2, 2)) >>> # the following RPC shares the rref with worker1, reference >>> # count is automatically updated. >>> rpc.rpc_sync("worker1", f, args=(rref,))
-
backward(self: torch._C._distributed_rpc.PyRRef, dist_autograd_ctx_id: int = -1, retain_graph: bool = False) → None -
Выполняет обратное распространение ошибки, используя RRef в качестве корня обратного распространения ошибки. Если
dist_autograd_ctx_idпредоставлен, мы выполняем распределённое обратное распространение ошибки с использованием предоставленного ctx_id, начиная с владельца RRef. В этом случае,get_gradients()следует использовать для получения градиентов. Еслиdist_autograd_ctx_idNone, предполагается, что это локальный граф autograd, и мы выполним только локальное обратное распространение ошибки. В локальном случае узел, вызывающий этот API, должен быть владельцем RRef. Ожидается, что значение RRef будет скалярным тензором.- Параметры
-
- dist_autograd_ctx_id (int, optional) – Идентификатор распределённого контекста autograd, для которого мы должны получить градиенты (по умолчанию: -1).
-
retain_graph (bool, optional) – Если
False, граф, используемый для вычисления градиента, будет освобождён. Обратите внимание, что в большинстве случаев установление этого параметра вTrueне требуется и часто может быть обойдено более эффективным способом. Обычно вам нужно установить его вTrueдля выполнения обратного распространения ошибки несколько раз (по умолчанию: False).
- Пример::
-
>>> import torch.distributed.autograd as dist_autograd >>> with dist_autograd.context() as context_id: >>> rref.backward(context_id)
-
confirmed_by_owner(self: torch._C._distributed_rpc.PyRRef) → bool -
Возвращает, подтверждена ли эта ссылка
RRefвладельцем.OwnerRRefвсегда возвращает true, в то время какUserRRefвозвращает true только тогда, когда владелец знает об этойUserRRef.
-
is_owner(self: torch._C._distributed_rpc.PyRRef) → bool -
Возвращает, является ли текущий узел владельцем этой ссылки
RRef.
-
local_value(self: torch._C._distributed_rpc.PyRRef) → object -
Если текущий узел является владельцем, возвращает ссылку на локальное значение. В противном случае генерирует исключение.
-
owner(self: torch._C._distributed_rpc.PyRRef) → torch._C._distributed_rpc.WorkerInfo -
Возвращает информацию об узле, владеющем этой ссылкой
RRef.
-
owner_name(self: torch._C._distributed_rpc.PyRRef) → str -
Возвращает имя узла, владеющего этой ссылкой
RRef.
-
remote(self: torch._C._distributed_rpc.PyRRef, timeout: float = -1.0) → object -
Создаёт помощника-прокси для лёгкого запуска
remote, используя владельца RRef в качестве назначения для выполнения функций над объектом, на который ссылается эта RRef. Более конкретно,rref.remote().func_name(*args, **kwargs)эквивалентно следующему:>>> def run(rref, func_name, args, kwargs): >>> return getattr(rref.local_value(), func_name)(*args, **kwargs) >>> >>> rpc.remote(rref.owner(), run, args=(rref, func_name, args, kwargs))
- Параметры
-
timeout (float, optional) – Временной лимит для
rref.remote(). Если создание этой ссылкиRRefне завершится успешно в течение этого времени, при следующей попытке использования RRef (например,to_here) будет вызвано исключение по таймауту. Если не указано, будет использован стандартный таймаут RPC. См.rpc.remote()для подробностей о таймаутахRRef.
- Пример::
-
>>> from torch.distributed import rpc >>> rref = rpc.remote("worker1", torch.add, args=(torch.zeros(2, 2), 1)) >>> rref.remote().size().to_here() # returns torch.Size([2, 2]) >>> rref.remote().view(1, 4).to_here() # returns tensor([[1., 1., 1., 1.]])
-
rpc_async(self: torch._C._distributed_rpc.PyRRef, timeout: float = -1.0) → object -
Создаёт помощника-прокси для лёгкого запуска
rpc_async, используя владельца RRef в качестве назначения для выполнения функций над объектом, на который ссылается эта RRef. Более конкретно,rref.rpc_async().func_name(*args, **kwargs)эквивалентно следующему:>>> def run(rref, func_name, args, kwargs): >>> return getattr(rref.local_value(), func_name)(*args, **kwargs) >>> >>> rpc.rpc_async(rref.owner(), run, args=(rref, func_name, args, kwargs))
- Параметры
-
timeout (float, optional) – Временной лимит для
rref.rpc_async(). Если вызов не завершится в течение этого времени, будет возбуждено исключение. Если этот аргумент не указан, будет использован стандартный таймаут RPC.
- Пример::
-
>>> from torch.distributed import rpc >>> rref = rpc.remote("worker1", torch.add, args=(torch.zeros(2, 2), 1)) >>> rref.rpc_async().size().wait() # returns torch.Size([2, 2]) >>> rref.rpc_async().view(1, 4).wait() # returns tensor([[1., 1., 1., 1.]])
-
rpc_sync(self: torch._C._distributed_rpc.PyRRef, timeout: float = -1.0) → object -
Создаёт помощника-прокси для лёгкого запуска
rpc_sync, используя владельца RRef в качестве назначения для выполнения функций над объектом, на который ссылается эта RRef. Более конкретно,rref.rpc_sync().func_name(*args, **kwargs)эквивалентно следующему:>>> def run(rref, func_name, args, kwargs): >>> return getattr(rref.local_value(), func_name)(*args, **kwargs) >>> >>> rpc.rpc_sync(rref.owner(), run, args=(rref, func_name, args, kwargs))
- Параметры
-
timeout (float, optional) – Временной лимит для
rref.rpc_sync(). Если вызов не завершится в течение этого времени, будет возбуждено исключение. Если этот аргумент не указан, будет использован стандартный таймаут RPC.
- Пример::
-
>>> from torch.distributed import rpc >>> rref = rpc.remote("worker1", torch.add, args=(torch.zeros(2, 2), 1)) >>> rref.rpc_sync().size() # returns torch.Size([2, 2]) >>> rref.rpc_sync().view(1, 4) # returns tensor([[1., 1., 1., 1.]])
-
to_here(self: torch._C._distributed_rpc.PyRRef, timeout: float = -1.0) → object -
Блокирующий вызов, копирующий значение RRef с узла-владельца на локальный узел и возвращающий его. Если текущий узел является владельцем, возвращает ссылку на локальное значение.
- Параметры
-
timeout (float, optional) – Временной лимит для
to_here. Если вызов не завершится в течение этого времени, будет возбуждено исключение. Если этот аргумент не указан, будет использован стандартный таймаут RPC (60 сек).
Дополнительная информация о RRef
Модуль удалённого доступа
Предупреждение
В настоящее время 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выполняется на удалённом узле. Он обрабатывает запись автоградиента, чтобы обеспечить, что обратное распространение градиента вернётся к соответствующему удалённому модулю.Он генерирует два метода
forward_asyncиforwardна основе сигнатуры методаforwardмодуляmodule_cls.forward_asyncвыполняется асинхронно и возвращает будущее. Аргументы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) – аргументы, которые нужно передать в
module_cls. -
kwargs (Dict, optional) – ключевые слова, которые нужно передать в
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]]) параметров удалённого модуля. - Тип возвращаемого значения
Фреймворк распределённого автоградиента
Предупреждение
Распределённый автоградиент в настоящее время не поддерживается при использовании тензоров CUDA
Этот модуль предоставляет фреймворк распределённого автоградиента на основе RPC, который может быть использован для таких приложений, как обучение моделей параллельно. Короче говоря, приложения могут отправлять и получать тензоры записи градиента по RPC. В прямом проходе мы записываем, когда тензоры записи градиента отправляются по RPC, а в обратном проходе мы используем эту информацию для выполнения распределённого обратного прохода с помощью RPC. Более подробную информацию см. в проектировании распределённого автоградиента.
-
torch.distributed.autograd.backward(context_id: int, roots: List[Tensor], retain_graph=False) → None -
Запускает обратный проход распределённого автоградиента с использованием предоставленных корней. В настоящее время это реализует алгоритм алгоритм FAST, который предполагает, что все сообщения RPC, отправленные в одном контексте распределённого автоградиента на всех рабочих узлах, будут частью графа автоградиента во время обратного прохода.
Мы используем предоставленные корни для обнаружения графа автоградиента и вычисления соответствующих зависимостей. Этот метод блокируется до тех пор, пока не будет выполнено всё вычисление автоградиента.
Мы накапливаем градиенты в соответствующем
torch.distributed.autograd.contextна каждом из узлов. Контекст автоградиента, который нужно использовать, ищется поcontext_idпереданному во время вызоваtorch.distributed.autograd.backward(). Если нет допустимого контекста автоградиента, соответствующего данному идентификатору, выводится ошибка. Вы можете получить накопленные градиенты с помощью APIget_gradients().- Параметры
-
- context_id (int) – Идентификатор контекста автоградиента, для которого нужно получить градиенты.
- roots (список) – Тензоры, которые представляют корни вычисления автоградиента. Все тензоры должны быть скалярами.
- retain_graph (bool, optional) – Если 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] -
Объект контекста для обертывания прямых и обратных проходов при использовании распределённого автоградиента.
context_idсгенерированный в оператореwithнеобходим для уникальной идентификации распределённого обратного прохода на всех рабочих узлах. Каждый рабочий узел сохраняет метаданные, связанные с этимcontext_id, что необходимо для корректного выполнения распределённого прохода автоградиента.- Пример::
-
>>> 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в рамках обратного прохода распределённого автоградиента.- Параметры
-
context_id (int) – Идентификатор контекста автоградиента, для которого нужно получить градиенты.
- Возвращает
-
Отображение, где ключ – 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
Дистрибутивный оптимизатор
См. страницу torch.distributed.optim для документации по дистрибутивным оптимизаторам.
Примечания к проектированию
В примечаниях к проектированию дистрибутивного автодифференцирования описывается структура фреймворка дистрибутивного автодифференцирования на основе RPC, который полезен для таких применений, как обучение моделей параллельно.
В примечаниях к проектированию 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/2.1/rpc.html