Пакет multiprocessing - torch.multiprocessing
torch.multiprocessing является оболочкой над родным модулем multiprocessing. Он регистрирует пользовательские редукторы, которые используют общую память для предоставления общих представлений об одних и тех же данных в разных процессах. После перемещения тензора/хранилища в shared_memory (см. share_memory_()), его можно будет отправлять в другие процессы без создания копий.
API полностью совместим с оригинальным модулем — достаточно изменить import multiprocessing на import torch.multiprocessing для перемещения всех тензоров через очереди или совместного использования через другие механизмы в общую память.
Из-за сходства API мы не документируем большую часть содержимого этого пакета и рекомендуем обратиться к отличной документации оригинального модуля.
Предупреждение
Если основной процесс неожиданно завершается (например, из-за входящего сигнала), Python’s multiprocessing иногда не удается очистить своих дочерние процессы. Это известный недостаток, поэтому если вы видите утечки ресурсов после прерывания интерпретатора, это, вероятно, означает, что это произошло с вами.
Управление стратегиями
-
torch.multiprocessing.get_all_sharing_strategies()[source] -
Возвращает набор стратегий совместного использования, поддерживаемых на текущей системе.
-
torch.multiprocessing.get_sharing_strategy()[source] -
Возвращает текущую стратегию совместного использования тензоров CPU.
-
torch.multiprocessing.set_sharing_strategy(new_strategy)[source] -
Устанавливает стратегию для совместного использования тензоров CPU.
- Parameters
-
new_strategy (str) – Название выбранной стратегии. Должно быть одним из значений, возвращаемых
get_all_sharing_strategies().
Детали совместного использования тензоров CUDA
В отличие от тензоров CPU, процесс отправки должен хранить оригинальный тензор до тех пор, пока процесс приема сохраняет копию тензора. Управление ссылками реализуется внутри, но пользователи должны следовать лучшим практикам.
Предупреждение
Если процесс-получатель неожиданно завершается из-за фатального сигнала, общий тензор может быть навсегда сохранен в памяти до тех пор, пока процесс отправки работает.
- Освободите память как можно скорее в процессе-получателе.
## 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)
2. Поддерживайте работу процесса-отправителя до тех пор, пока все процессы-получатели не завершат работу. Это предотвратит ситуацию, когда процесс-отправитель освобождает память, которая все еще используется процессом-получателем.
## producer # send tensors, do something event.wait()
## consumer # receive tensors and use them event.set()
- Не передавайте полученные тензоры.
# 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(fn, args=(), nprocs=1, join=True, daemon=False, start_method='spawn')[source] -
Запускает
nprocsпроцессы, выполняющиеfnсargs.Если один из процессов завершается с ненулевым кодом завершения, оставшиеся процессы убиваются, и возникает исключение с причиной завершения. В случае, если исключение было перехвачено в дочернем процессе, оно пересылается, а его стек отслеживания включен в исключение, поднятое в родительском процессе.
- Parameters
-
-
fn (function) –
Функция используется в качестве точки входа запущенного процесса. Эта функция должна быть определена на верхнем уровне модуля, чтобы ее можно было сериализовать и запустить. Это требование, наложенное модулем multiprocessing.
Функция вызывается как
fn(i, *args), гдеi— индекс процесса, аargs— переданная кортеж аргументов. -
args (tuple) – Аргументы, передаваемые в
fn. - nprocs (int) – Количество процессов для запуска.
- join (bool) – Выполнить блокирующее присоединение ко всем процессам.
- daemon (bool) – Флаг демонического процесса для запущенных процессов. Если установлено в True, создаются демонические процессы.
-
start_method (str) – (устаревшее) этот метод всегда будет использовать
spawnв качестве метода запуска. Чтобы использовать другой метод запуска, используйтеstart_processes().
-
- Returns
-
None, если
join—True,ProcessContextеслиjoin—False
-
class torch.multiprocessing.SpawnContext[source] -
Возвращается функцией
spawn()при вызове сjoin=False.-
join(timeout=None) -
Пытается объединить один или несколько процессов в этом контексте запуска. Если один из них завершился с ненулевым кодом выхода, эта функция убивает оставшиеся процессы и вызывает исключение с причиной завершения первого процесса.
Возвращает
Trueесли все процессы успешно объединены,Falseесли есть дополнительные процессы, которые необходимо объединить.- Параметры
-
timeout (float) – Ждать столько времени, прежде чем отказаться от ожидания.
-
© 2024, PyTorch Contributors
PyTorch has a BSD-style license, as found in the LICENSE file.
https://pytorch.org/docs/2.1/multiprocessing.html