Распределенная 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 соединяет локальные движки автоградов на всех рабочих узлах, участвующих в прямом проходе, и автоматически обращается к ним во время обратного прохода для вычисления градиентов. Это особенно полезно, если прямой проход должен охватывать несколько машин при выполнении, например, распределенного обучения моделей, обучения с параметрическим сервером и т. д. С помощью этой функции код пользователя больше не должен беспокоиться о том, как отправлять градиенты через границы RPC и в каком порядке должны запускаться локальные движки автоградов, что может стать довольно сложным в случае вложенных и взаимозависимых вызовов RPC в прямом проходе.
-
Распределенный оптимизатор конструктор принимает
Optimizer()(например,SGD(),Adagrad()и т. д.) и список параметр RRefs, создает экземплярOptimizer()на каждом отдельном владельце RRef и обновляет параметры соответственно при выполненииstep(). Когда у вас есть распределенные прямой и обратный проходы, параметры и градиенты будут распределены по нескольким рабочим узлам, и поэтому требуется оптимизатор на каждом из вовлеченных рабочих узлов. Распределенный оптимизатор объединяет все эти локальные оптимизаторы в один и предоставляет краткий конструктор и APIstep().
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должны быть правильно установлены. Дополнительная информация и доступные параметры приведены в разделе Бэкенды.
-
name (str) – глобально уникальное имя этого узла. (например,
Следующие 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.
-
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или вызов 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правильно настроены на обоих узлах. Обратитесь к 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 также предоставляет декораторы, которые позволяют приложениям указать, как должна обрабатываться данная функция на стороне вызываемого.
-
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.
-
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. Эта функция может вызываться несколько раз для поэтапного добавления конфигураций размещения устройств.
- Параметры:
-
- 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
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]), указывающий на удаленный модуль.
-
remote_parameters(recurse=True) -
Возвращает список
RRef, указывающих на параметры удаленного модуля. Это обычно используется совместно сDistributedOptimizer.- Параметры:
-
recurse (bool) – если True, то возвращает параметры удаленного модуля и всех его подмодулей. В противном случае возвращает только параметры, которые являются непосредственными членами удаленного модуля.
- Возвращает:
-
Список
RRef(List[RRef[nn.Parameter]]) на параметры удаленного модуля. - Тип возвращаемого значения:
Фреймворк распределенного 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, соответствующего заданному идентификатору, мы выводим ошибку. Вы можете получить накопленные градиенты с помощью APIget_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
Распределённый оптимизатор
См. страницу torch.distributed.optim для документации по распределённым оптимизаторам.
Примечания по проектированию
Заметки по проектированию распределённого autograd охватывают проектирование фреймворка распределённого autograd на основе 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/1.13/rpc.html