Spec-Zone.ru › PyTorch 2.14

Пакет многопроцессорной обработки — torch.multiprocessing

Создано: 23 дек. 2016 г. | Последнее обновление: 13 апр. 2026 г.

torch.multiprocessing — это оболочка над встроенным модулем multiprocessing.

Он регистрирует пользовательские редукторы, использующие разделяемую память для предоставления общих представлений одних и тех же данных в разных процессах. После перемещения тензора/хранилища в shared_memory (см. share_memory_()) его можно будет отправлять другим процессам без копирования.

API на 100% совместим с исходным модулем — достаточно заменить import multiprocessing на import torch.multiprocessing, чтобы все тензоры, отправленные через очереди или переданные с помощью других механизмов, перемещались в разделяемую память.

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

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

Если основной процесс завершается аварийно (например, из-за полученного сигнала), Python иногда не удается очистить его дочерние процессы с помощью multiprocessing. Это известная проблема, поэтому, если после прерывания интерпретатора вы заметили утечку ресурсов, вероятно, произошло именно это.

Управление стратегиями

torch.multiprocessing.get_all_sharing_strategies() [исходный код]

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

torch.multiprocessing.get_sharing_strategy() [исходный код]

Возвращает текущую стратегию совместного использования тензоров CPU.

torch.multiprocessing.set_sharing_strategy(new_strategy) [исходный код]

Задает стратегию совместного использования тензоров CPU.

Параметры:

new_strategy (str) – Название выбранной стратегии. Должно совпадать с одним из значений, возвращаемых функцией get_all_sharing_strategies().

Совместное использование тензоров CUDA

Совместное использование тензоров CUDA между процессами поддерживается только в Python 3 при использовании методов запуска spawn или forkserver.

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

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

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

  1. Как можно скорее освобождайте память в процессе-потребителе.
## Good
x = queue.get()
# do somethings with x
del x
## Bad
x = queue.get()
# do somethings with x
# do everything else (producer have to keep x in memory)
  1. Не завершайте процесс-производитель, пока не выйдут все процессы-потребители. Это предотвратит ситуацию, когда процесс-производитель освобождает память, которая все еще используется процессом-потребителем.
## producer
# send tensors, do something
event.wait()
## consumer
# receive tensors and use them
event.set()
  1. Не передавайте полученные тензоры дальше.
# not going to work
x = queue.get()
queue_2.put(x)
# you need to create a process-local copy
x = queue.get()
x_clone = x.clone()
queue_2.put(x_clone)
# putting and getting from the same queue in the same process will likely end up with segfault
queue.put(tensor)
x = queue.get()

Стратегии совместного использования

В этом разделе кратко описано, как работают различные стратегии совместного использования. Обратите внимание, что они применяются только к тензорам CPU: для тензоров CUDA всегда используется API CUDA, поскольку это единственный способ совместно использовать их.

Дескриптор файла — file_descriptor

Примечание

Это стратегия по умолчанию (кроме macOS и OS X, где она не поддерживается).

В этой стратегии в качестве дескрипторов разделяемой памяти используются дескрипторы файлов. Когда хранилище перемещается в разделяемую память, дескриптор файла, полученный из shm_open, кэшируется вместе с объектом; при отправке объекта другим процессам этот дескриптор передается им (например, через сокеты UNIX). Получатель также кэширует дескриптор файла и mmap его, чтобы получить общее представление данных хранилища.

Обратите внимание: если совместно используется много тензоров, эта стратегия будет большую часть времени удерживать открытыми многочисленные дескрипторы файлов. Если в вашей системе установлены низкие ограничения на число открытых дескрипторов файлов и вы не можете их повысить, следует использовать стратегию file_system.

Файловая система — file_system

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

Для устранения утечек файлов разделяемой памяти torch.multiprocessing запускает демон с именем torch_shm_manager, который изолируется от текущей группы процессов и отслеживает все выделения разделяемой памяти. Когда все подключенные к нему процессы завершаются, он немного ждет, чтобы убедиться в отсутствии новых подключений, а затем перебирает все файлы разделяемой памяти, выделенные группой. Если какие-либо из них все еще существуют, демон освобождает их. Мы протестировали этот метод и убедились, что он устойчив к различным сбоям. Тем не менее, если в вашей системе достаточно высокие ограничения и поддерживается стратегия file_descriptor, мы не рекомендуем переключаться на эту стратегию.

Запуск дочерних процессов

Примечание

Доступно для Python >= 3.4.

Работа зависит от метода запуска spawn в пакете multiprocessing Python.

Для выполнения некоторой функции в нескольких дочерних процессах можно создать экземпляры Process и вызвать join, чтобы дождаться их завершения. Этот подход хорошо работает с одним дочерним процессом, но при работе с несколькими процессами могут возникнуть проблемы.

В частности, последовательное присоединение процессов означает, что они будут завершаться последовательно. Если это не так и первый процесс не завершится, завершение остальных процессов останется незамеченным. Кроме того, встроенных средств для передачи ошибок нет.

Функция spawn, приведенная ниже, решает эти проблемы: она обеспечивает передачу ошибок, обрабатывает завершение процессов не по порядку и активно завершает процессы при обнаружении ошибки в одном из них.

torch.multiprocessing.spawn.spawn(fn, args=(), nprocs=1, join=True, daemon=False, start_method='spawn') [исходный код]

Запускает nprocs процессов, выполняющих fn с аргументами args.

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

Параметры:
  • fn (функция) –

    Функция вызывается как точка входа запущенного процесса. Она должна быть определена на верхнем уровне модуля, чтобы ее можно было сериализовать с помощью pickle и запустить. Это требование, накладываемое модулем multiprocessing.

    Функция вызывается как fn(i, *args), где i — индекс процесса, а args — переданный кортеж аргументов.

  • args (tuple) – Аргументы, передаваемые в fn.
  • nprocs (int) – Число запускаемых процессов.
  • join (bool) – Блокирующее ожидание завершения всех процессов.
  • daemon (bool) – Флаг демона для запускаемых процессов. Если установлено значение True, будут созданы процессы-демоны.
  • start_method (str) – (устарел) этот метод всегда использует spawn в качестве метода запуска. Чтобы использовать другой метод запуска, воспользуйтесь start_processes().
Возвращает:

None, если join имеет значение True, ProcessContext, если join имеет значение False

class torch.multiprocessing.SpawnContext [исходный код]

Возвращается функцией spawn(), если она вызвана с join=False.

join(timeout=None, grace_period=None) [исходный код]

Ожидает завершения одного или нескольких процессов в контексте запуска.

Пытается дождаться завершения одного или нескольких процессов в этом контексте запуска. Если один из них завершился с ненулевым кодом возврата, функция принудительно завершает остальные процессы (при необходимости предоставляя им время на завершение) и возбуждает исключение с причиной завершения первого процесса.

Возвращает True, если все процессы успешно завершили работу, и False, если остались процессы, завершения которых нужно дождаться.

Параметры:
  • timeout (float) – Время ожидания в секундах, по истечении которого ожидание прекращается.
  • grace_period (float) – Если какой-либо процесс завершится с ошибкой, столько секунд другие процессы будут ожидать перед завершением, чтобы успеть завершиться корректно. Если они все еще не завершились, перед их принудительным завершением выжидается еще один льготный период.

torch.multiprocessing.pool

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

Spec-Zone.ru

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