multiprocessing — Параллелизм на основе процессов
Исходный код: Lib/multiprocessing/
Доступность: не поддерживается на Android, iOS и WASI.
Этот модуль не поддерживается на мобильных платформах и платформах WebAssembly.
Введение
multiprocessing — это пакет, который поддерживает создание процессов с помощью API, аналогичного модулю threading. Пакет multiprocessing обеспечивает как локальную, так и удалённую параллельную работу, эффективно обходя глобальную блокировку интерпретатора за счёт использования дочерних процессов вместо потоков. Благодаря этому модуль multiprocessing позволяет программисту в полной мере использовать несколько процессоров на компьютере. Он работает как в POSIX, так и в Windows.
Модуль multiprocessing также представляет объект Pool, предоставляющий удобный способ распараллелить выполнение функции для нескольких входных значений, распределяя входные данные между процессами (параллелизм данных). В следующем примере показана распространённая практика определения таких функций в модуле, чтобы дочерние процессы могли успешно импортировать этот модуль. Этот простой пример параллелизма данных с помощью Pool,
from multiprocessing import Pool
def f(x):
return x*x
if __name__ == '__main__':
with Pool(5) as p:
print(p.map(f, [1, 2, 3]))
выведет в стандартный вывод
[1, 4, 9]
Модуль multiprocessing также предоставляет API, у которого нет аналогов в модуле threading, например возможность terminate, interrupt или kill работающий процесс.
См. также
concurrent.futures.ProcessPoolExecutor предоставляет интерфейс более высокого уровня для отправки задач фоновому процессу без блокировки выполнения вызывающего процесса. По сравнению с непосредственным использованием интерфейса Pool, API concurrent.futures позволяет проще отделить отправку задач в пул процессов от ожидания результатов.
Класс Process
В multiprocessing процессы создаются путём создания объекта Process и последующего вызова его метода start(). Process соответствует API класса threading.Thread. Простейшая программа с несколькими процессами выглядит так:
from multiprocessing import Process
def f(name):
print('hello', name)
if __name__ == '__main__':
p = Process(target=f, args=('bob',))
p.start()
p.join()
Вот расширенный пример, показывающий идентификаторы отдельных процессов:
from multiprocessing import Process
import os
def info(title):
print(title)
print('module name:', __name__)
print('parent process:', os.getppid())
print('process id:', os.getpid())
def f(name):
info('function f')
print('hello', name)
if __name__ == '__main__':
info('main line')
p = Process(target=f, args=('bob',))
p.start()
p.join()
Объяснение необходимости части if __name__ == '__main__' см. в разделе Рекомендации по программированию.
Аргументы для Process обычно должны поддерживать сериализацию с помощью pickle, чтобы их можно было передать дочернему процессу. Если попробовать ввести приведённый выше пример непосредственно в REPL, в дочернем процессе может возникнуть AttributeError при попытке найти функцию f в модуле __main__.
Контексты и методы запуска
В зависимости от платформы multiprocessing поддерживает три способа запуска процесса. Эти методы запуска:
- spawn
-
Родительский процесс запускает новый процесс интерпретатора Python. Дочерний процесс унаследует только те ресурсы, которые необходимы для выполнения метода
run()объекта процесса. В частности, ненужные файловые дескрипторы и дескрипторы из родительского процесса унаследованы не будут. Запуск процесса этим способом довольно медленный по сравнению с использованием fork или forkserver.Доступен на платформах POSIX и Windows. Используется по умолчанию в Windows и macOS.
- fork
-
Родительский процесс использует
os.fork()для создания копии интерпретатора Python. После запуска дочерний процесс практически идентичен родительскому. Дочерний процесс наследует все ресурсы родительского процесса. Обратите внимание, что безопасно создавать копию многопоточного процесса проблематично.Доступен в системах POSIX.
Изменено в версии 3.14: Этот метод больше не используется по умолчанию ни на одной платформе. В коде, которому необходим fork, его нужно явно указывать с помощью
get_context()илиset_start_method().Изменено в версии 3.12: Если Python может определить, что ваш процесс является многопоточным, функция
os.fork(), вызываемая этим методом запуска, выдаст предупреждениеDeprecationWarning. Используйте другой метод запуска. Дополнительные пояснения см. в документацииos.fork().
- forkserver
-
При запуске программы с выбором метода forkserver создаётся серверный процесс. После этого, когда требуется создать новый процесс, родительский процесс подключается к серверу и запрашивает создание нового процесса путём копирования. Серверный процесс создания копий является однопоточным, если только системные библиотеки или предварительно загруженные импорты не создают потоки в качестве побочного эффекта, поэтому обычно он может безопасно использовать
os.fork(). Ненужные ресурсы не наследуются.Доступен на платформах POSIX, поддерживающих передачу файловых дескрипторов через каналы Unix, например в Linux. На таких платформах используется по умолчанию.
Изменено в версии 3.14: Этот метод стал методом запуска по умолчанию на платформах POSIX.
Изменено в версии 3.4: Метод spawn добавлен на все платформы POSIX, а forkserver — на некоторые платформы POSIX. Дочерние процессы в Windows больше не наследуют все дескрипторы родительского процесса, помеченные как наследуемые.
Изменено в версии 3.8: В macOS метод запуска spawn теперь используется по умолчанию. Метод fork следует считать небезопасным, поскольку он может привести к сбоям дочернего процесса: системные библиотеки macOS могут запускать потоки. См. bpo-33725.
Изменено в версии 3.14: На платформах POSIX метод запуска по умолчанию изменён с fork на forkserver, чтобы сохранить производительность и избежать распространённых проблем несовместимости с многопоточными процессами. См. gh-84559.
При использовании методов запуска spawn или forkserver в POSIX также запускается процесс отслеживания ресурсов, который отслеживает именованные системные ресурсы, для которых была отменена связь (например, именованные семафоры или объекты SharedMemory), созданные процессами программы. После завершения всех процессов отслеживатель ресурсов удаляет все оставшиеся отслеживаемые объекты. Обычно таких объектов быть не должно, но если процесс был завершён сигналом, некоторые ресурсы могут «утечь». (Ни утёкшие семафоры, ни сегменты разделяемой памяти не будут автоматически удалены до следующей перезагрузки. Это проблема для обоих типов объектов, поскольку система допускает лишь ограниченное количество именованных семафоров, а сегменты разделяемой памяти занимают место в основной памяти.)
Чтобы выбрать метод запуска, используйте set_start_method() в блоке if __name__ == '__main__' главного модуля. Например:
import multiprocessing as mp
def foo(q):
q.put('hello')
if __name__ == '__main__':
mp.set_start_method('spawn')
q = mp.Queue()
p = mp.Process(target=foo, args=(q,))
p.start()
print(q.get())
p.join()
set_start_method() не следует использовать в программе более одного раза.
В качестве альтернативы можно использовать get_context(), чтобы получить объект контекста. Объекты контекста имеют тот же API, что и модуль multiprocessing, и позволяют использовать в одной программе несколько методов запуска.
import multiprocessing as mp
def foo(q):
q.put('hello')
if __name__ == '__main__':
ctx = mp.get_context('spawn')
q = ctx.Queue()
p = ctx.Process(target=foo, args=(q,))
p.start()
print(q.get())
p.join()
Обратите внимание, что объекты, связанные с одним контекстом, могут быть несовместимы с процессами другого контекста. В частности, блокировки, созданные с помощью контекста fork, нельзя передавать процессам, запущенным с помощью методов spawn или forkserver.
Библиотеки, использующие multiprocessing или ProcessPoolExecutor, следует проектировать так, чтобы пользователи могли передавать собственный контекст multiprocessing. Использование в библиотеке конкретного контекста может привести к несовместимости с остальной частью приложения её пользователя. Всегда документируйте, если библиотеке требуется определённый метод запуска.
Предупреждение
Методы запуска 'spawn' и 'forkserver' обычно нельзя использовать в «замороженных» исполняемых файлах (то есть в бинарных файлах, созданных такими пакетами, как PyInstaller и cx_Freeze) в системах POSIX. Метод запуска 'fork' может работать, если код не использует потоки.
Обмен объектами между процессами
multiprocessing поддерживает два типа каналов связи между процессами:
Очереди
Класс Queue почти полностью повторяет queue.Queue. Например:
from multiprocessing import Process, Queue
def f(q):
q.put([42, None, 'hello'])
if __name__ == '__main__':
q = Queue()
p = Process(target=f, args=(q,))
p.start()
print(q.get()) # prints "[42, None, 'hello']"
p.join()
Очереди безопасны для потоков и процессов. Любой объект, помещённый в очередь multiprocessing, будет сериализован.
Каналы
Функция Pipe() возвращает пару объектов соединения, связанных каналом, который по умолчанию является двунаправленным. Например:
from multiprocessing import Process, Pipe
def f(conn):
conn.send([42, None, 'hello'])
conn.close()
if __name__ == '__main__':
parent_conn, child_conn = Pipe()
p = Process(target=f, args=(child_conn,))
p.start()
print(parent_conn.recv()) # prints "[42, None, 'hello']"
p.join()
Два объекта соединения, возвращаемые функцией Pipe(), представляют два конца канала. У каждого объекта соединения есть методы send() и recv() (а также другие). Обратите внимание, что данные в канале могут быть повреждены, если два процесса (или потока) одновременно попытаются читать из одного и того же конца канала или записывать в него. Конечно, если процессы одновременно используют разные концы канала, риск повреждения данных отсутствует.
Метод send() сериализует объект, а recv() воссоздаёт его.
Синхронизация между процессами
multiprocessing содержит аналоги всех примитивов синхронизации из threading. Например, можно использовать блокировку, чтобы гарантировать, что в каждый момент времени только один процесс выводит данные в стандартный вывод:
from multiprocessing import Process, Lock
def f(l, i):
l.acquire()
try:
print('hello world', i)
finally:
l.release()
if __name__ == '__main__':
lock = Lock()
for num in range(10):
Process(target=f, args=(lock, num)).start()
Без блокировки вывод разных процессов может смешиваться.
Использование пула рабочих процессов
Класс Pool представляет пул рабочих процессов. Он предоставляет методы, позволяющие несколькими способами передавать задачи рабочим процессам для выполнения.
Например:
from multiprocessing import Pool, TimeoutError
import time
import os
def f(x):
return x*x
if __name__ == '__main__':
# start 4 worker processes
with Pool(processes=4) as pool:
# print "[0, 1, 4,..., 81]"
print(pool.map(f, range(10)))
# print same numbers in arbitrary order
for i in pool.imap_unordered(f, range(10)):
print(i)
# evaluate "f(20)" asynchronously
res = pool.apply_async(f, (20,)) # runs in *only* one process
print(res.get(timeout=1)) # prints "400"
# evaluate "os.getpid()" asynchronously
res = pool.apply_async(os.getpid, ()) # runs in *only* one process
print(res.get(timeout=1)) # prints the PID of that process
# launching multiple evaluations asynchronously *may* use more processes
multiple_results = [pool.apply_async(os.getpid, ()) for i in range(4)]
print([res.get(timeout=1) for res in multiple_results])
# make a single worker sleep for 10 seconds
res = pool.apply_async(time.sleep, (10,))
try:
print(res.get(timeout=1))
except TimeoutError:
print("We lacked patience and got a multiprocessing.TimeoutError")
print("For the moment, the pool remains available for more work")
# exiting the 'with'-block has stopped the pool
print("Now the pool is closed and no longer available")
Обратите внимание, что методы пула должен использовать только создавший его процесс.
Примечание
Для работы функциональности этого пакета дочерние процессы должны иметь возможность импортировать модуль __main__. Это рассматривается в разделе Рекомендации по программированию, однако на этом стоит остановиться отдельно. Это означает, что некоторые примеры, например примеры с multiprocessing.pool.Pool, не будут работать в интерактивном интерпретаторе. Например:
>>> from multiprocessing import Pool >>> p = Pool(5) >>> def f(x): ... return x*x ... >>> with p: ... p.map(f, [1,2,3]) Process PoolWorker-1: Process PoolWorker-2: Process PoolWorker-3: Traceback (most recent call last): Traceback (most recent call last): Traceback (most recent call last): AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)> AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)> AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)>
(Если вы попробуете выполнить этот пример, он выведет три полные трассировки стека, перемежающиеся в полуслучайном порядке, после чего, возможно, придётся каким-либо образом остановить родительский процесс.)
Справочник
Пакет multiprocessing в основном воспроизводит API модуля threading.
Глобальный метод запуска
Python поддерживает несколько способов создания и инициализации процесса. Глобальный метод запуска задаёт механизм создания процесса по умолчанию.
Некоторые функции и методы multiprocessing, которые также могут создавать определённые объекты, неявно устанавливают глобальный метод запуска, принятый в системе по умолчанию, если он ещё не был установлен. Глобальный метод запуска можно установить только один раз. Если нужно изменить метод запуска, принятый в системе по умолчанию, следует явно задать глобальный метод запуска до вызова функций или методов либо создания этих объектов.
Process и исключения
-
class multiprocessing.Process(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None) -
Объекты Process представляют деятельность, выполняемую в отдельном процессе. Класс
Processсодержит аналоги всех методовthreading.Thread.Конструктор всегда следует вызывать с именованными аргументами. Значение group всегда должно быть
None; оно существует исключительно для совместимости сthreading.Thread. target — вызываемый объект, который будет вызван методомrun(). По умолчанию этоNone, то есть ничего не вызывается. name — имя процесса (подробности см. вname). args — кортеж аргументов для вызова целевого объекта. kwargs — словарь именованных аргументов для вызова целевого объекта. Если задан необязательный именованный аргумент daemon, он устанавливает флагdaemonпроцесса вTrueилиFalse. Если указаноNone(значение по умолчанию), этот флаг наследуется от создающего процесса.По умолчанию целевому объекту target не передаются аргументы. Аргумент args, значение которого по умолчанию —
(), позволяет задать список или кортеж аргументов, передаваемых в target.Если подкласс переопределяет конструктор, он должен вызвать конструктор базового класса (
super().__init__()), прежде чем выполнять какие-либо другие действия с процессом.Примечание
Как правило, все аргументы, передаваемые в
Process, должны поддерживать pickling. Это часто обнаруживается при попытке создатьProcessили использоватьconcurrent.futures.ProcessPoolExecutorв REPL с локально определённой функцией target.Передача вызываемого объекта, определённого в текущем сеансе REPL, приводит к аварийному завершению дочернего процесса из-за необработанного исключения
AttributeErrorпри запуске, поскольку target должен быть определён в импортируемом модуле, чтобы его можно было загрузить при распаковке.Пример этой необрабатываемой ошибки в дочернем процессе:
>>> import multiprocessing as mp >>> def knigit(): ... print("Ni!") ... >>> process = mp.Process(target=knigit) >>> process.start() >>> Traceback (most recent call last): File ".../multiprocessing/spawn.py", line ..., in spawn_main File ".../multiprocessing/spawn.py", line ..., in _main AttributeError: module '__main__' has no attribute 'knigit' >>> process <SpawnProcess name='SpawnProcess-1' pid=379473 parent=378707 stopped exitcode=1>См. Методы запуска spawn и forkserver. Это ограничение не действует при использовании метода запуска
"fork", однако начиная с Python3.14он больше не является методом по умолчанию ни на одной платформе. См. Контексты и методы запуска. См. также gh-132898.Изменено в версии 3.3: Добавлен параметр daemon.
-
run() -
Метод, представляющий деятельность процесса.
Этот метод можно переопределить в подклассе. Стандартный метод
run()вызывает вызываемый объект, переданный конструктору объекта в качестве целевого аргумента, если он задан, с позиционными и именованными аргументами, взятыми соответственно из аргументов args и kwargs.Передача списка или кортежа в качестве аргумента args для
Processдаёт тот же результат.Пример:
>>> from multiprocessing import Process >>> p = Process(target=print, args=[1]) >>> p.run() 1 >>> p = Process(target=print, args=(1,)) >>> p.run() 1
-
start() -
Запускает выполнение процесса.
Для каждого объекта процесса этот метод можно вызвать не более одного раза. Он обеспечивает вызов метода объекта
run()в отдельном процессе.
-
join([timeout]) -
Если необязательный аргумент timeout равен
None(значение по умолчанию), метод блокирует выполнение до завершения процесса, для которого был вызван методjoin(). Если timeout — положительное число, блокировка длится не более указанного количества секунд. Обратите внимание: метод возвращаетNone, если процесс завершился или истекло время ожидания. Чтобы определить, завершился ли процесс, проверьте егоexitcode.К процессу можно присоединяться многократно.
Процесс не может присоединиться к самому себе, поскольку это приведёт к взаимной блокировке. Попытка присоединиться к процессу до его запуска является ошибкой.
-
name -
Имя процесса. Имя — это строка, используемая только для идентификации. Оно не имеет семантического значения. Нескольким процессам можно присвоить одинаковое имя.
Начальное имя задаётся конструктором. Если конструктору не передано явное имя, формируется имя вида «Process-N1:N2:…:Nk», где каждый Nk — N-й дочерний процесс своего родителя.
-
is_alive() -
Возвращает признак того, что процесс активен.
Процесс считается активным с момента возврата метода
start()и до завершения дочернего процесса.
-
daemon -
Флаг демона процесса, логическое значение. Его необходимо установить до вызова
start().Начальное значение наследуется от создающего процесса.
При завершении процесса он пытается завершить все свои дочерние процессы-демоны.
Обратите внимание: процесс-демон не может создавать дочерние процессы. Иначе, если родительский процесс завершится, дочерние процессы-демоны останутся без родителя. Кроме того, это не демоны или службы Unix, а обычные процессы, которые будут завершены (без присоединения к ним), если все процессы, не являющиеся демонами, завершились.
Помимо API
threading.Thread, объектыProcessтакже поддерживают следующие атрибуты и методы:-
pid -
Возвращает идентификатор процесса. До создания процесса возвращается
None.
-
exitcode -
Код завершения дочернего процесса. Если процесс ещё не завершился, возвращается
None.Если метод
run()дочернего процесса завершился штатно, код завершения будет равен 0. Если процесс завершился черезsys.exit()с целочисленным аргументом N, код завершения будет равен N.Если дочерний процесс завершился из-за исключения, не перехваченного внутри
run(), код завершения будет равен 1. Если процесс был завершён сигналом N, код завершения будет равен отрицательному значению -N.
-
authkey -
Ключ аутентификации процесса (байтовая строка).
При инициализации
multiprocessingглавному процессу назначается случайная строка с помощьюos.urandom().При создании объекта
Processон наследует ключ аутентификации родительского процесса, однако его можно изменить, присвоив атрибутуauthkeyдругую байтовую строку.См. Ключи аутентификации.
-
sentinel -
Числовой дескриптор системного объекта, который перейдёт в состояние «готов», когда процесс завершится.
Это значение можно использовать, если нужно ожидать сразу несколько событий с помощью
multiprocessing.connection.wait(). В противном случае проще вызватьjoin().В Windows это дескриптор ОС, пригодный для использования с семейством вызовов API
WaitForSingleObjectиWaitForMultipleObjects. В POSIX это файловый дескриптор, пригодный для использования с примитивами модуляselect.Добавлено в версии 3.3.
-
interrupt() -
Прерывает процесс. В POSIX используется сигнал
SIGINT. Поведение в Windows не определено.По умолчанию дочерний процесс завершается путём возбуждения исключения
KeyboardInterrupt. Это поведение можно изменить, задав соответствующий обработчик сигнала в дочернем процессе с помощьюsignal.signal()для сигналаSIGINT.Примечание: если дочерний процесс перехватывает и игнорирует
KeyboardInterrupt, процесс не будет завершён.Примечание: поведение по умолчанию также устанавливает
exitcodeв1, как если бы в дочернем процессе было возбуждено необработанное исключение. Чтобы получить другое значениеexitcode, достаточно перехватитьKeyboardInterruptи вызватьexit(your_code).Добавлено в версии 3.14.
-
terminate() -
Завершает процесс. В POSIX для этого используется сигнал
SIGTERM; в Windows используетсяTerminateProcess(). Обратите внимание: обработчики выхода, блоки finally и т. п. выполняться не будут.Обратите внимание: дочерние процессы данного процесса не будут завершены — они просто останутся без родителя.
Предупреждение
Если этот метод используется, когда связанный процесс работает с каналом или очередью, канал или очередь могут быть повреждены и стать непригодными для других процессов. Аналогично, если процесс захватил блокировку, семафор и т. п., его завершение может привести к взаимной блокировке других процессов.
-
kill() -
То же, что и
terminate(), но в POSIX используется сигналSIGKILL.Добавлено в версии 3.7.
-
close() -
Закрывает объект
Processи освобождает все связанные с ним ресурсы. Если базовый процесс всё ещё выполняется, возбуждается исключениеValueError. После успешного возвратаclose()большинство других методов и атрибутов объектаProcessбудут вызывать исключениеValueError.Добавлено в версии 3.7.
Обратите внимание: методы
start(),join(),is_alive(),terminate()иexitcodeдолжен вызывать только процесс, создавший объект процесса.Пример использования некоторых методов
Process:>>> import multiprocessing, time, signal >>> mp_context = multiprocessing.get_context('spawn') >>> p = mp_context.Process(target=time.sleep, args=(1000,)) >>> print(p, p.is_alive()) <...Process ... initial> False >>> p.start() >>> print(p, p.is_alive()) <...Process ... started> True >>> p.terminate() >>> time.sleep(0.1) >>> print(p, p.is_alive()) <...Process ... stopped exitcode=-SIGTERM> False >>> p.exitcode == -signal.SIGTERM True -
-
exception multiprocessing.ProcessError -
Базовый класс всех исключений
multiprocessing.
-
exception multiprocessing.BufferTooShort -
Исключение, которое возбуждает
Connection.recv_bytes_into(), если переданный объект буфера слишком мал для прочитанного сообщения.Если
eявляется экземпляромBufferTooShort, тоe.args[0]вернёт сообщение в виде байтовой строки.
-
exception multiprocessing.AuthenticationError -
Возникает при ошибке аутентификации.
-
exception multiprocessing.TimeoutError -
Возникает в методах с тайм-аутом по истечении времени ожидания.
Трубы и очереди
При использовании нескольких процессов для взаимодействия между ними обычно применяют передачу сообщений, избегая необходимости использовать примитивы синхронизации, такие как блокировки.
Для передачи сообщений можно использовать Pipe() (для соединения двух процессов) или очередь (которая допускает нескольких производителей и потребителей).
Типы Queue, SimpleQueue и JoinableQueue представляют собой очереди FIFO с несколькими производителями и потребителями, созданные по образцу класса queue.Queue из стандартной библиотеки. Они отличаются тем, что в Queue отсутствуют методы task_done() и join(), добавленные в класс queue.Queue в Python 2.5.
Если вы используете JoinableQueue, то обязательно вызывайте JoinableQueue.task_done() для каждой задачи, извлечённой из очереди, иначе семафор, используемый для подсчёта количества незавершённых задач, может переполниться, что приведёт к возникновению исключения.
Одно из отличий от других реализаций очередей в Python состоит в том, что очереди multiprocessing сериализуют все помещаемые в них объекты с помощью pickle. Объект, возвращаемый методом get, является заново созданным объектом, который не использует общую память с исходным объектом.
Обратите внимание, что общую очередь также можно создать с помощью объекта-менеджера — см. раздел Менеджеры.
Примечание
multiprocessing использует стандартные исключения queue.Empty и queue.Full для сигнализации о тайм-ауте. Они недоступны в пространстве имён multiprocessing, поэтому их необходимо импортировать из queue.
Примечание
Когда объект помещается в очередь, он сериализуется с помощью pickle, а фоновый поток позднее передаёт сериализованные данные в нижележащий канал. Это имеет некоторые, немного неожиданные последствия, однако они не должны вызывать практических затруднений. Если же они действительно вас беспокоят, можно вместо этого использовать очередь, созданную с помощью менеджера.
- После помещения объекта в пустую очередь может пройти ничтожно малое время, прежде чем метод очереди
empty()вернётFalse, а вызовget_nowait()сможет завершиться без исключенияqueue.Empty. - Если несколько процессов помещают объекты в очередь, то на другой стороне они могут быть получены не в том порядке. Однако объекты, добавленные одним и тем же процессом, всегда будут расположены в ожидаемом порядке относительно друг друга.
Предупреждение
Если процесс завершён с помощью Process.terminate() или os.kill(), пока он пытается использовать Queue, данные в очереди, вероятно, будут повреждены. Впоследствии это может привести к исключению при попытке любого другого процесса использовать эту очередь.
Предупреждение
Как упоминалось выше, если дочерний процесс поместил элементы в очередь (и не использовал JoinableQueue.cancel_join_thread), он не завершится, пока все элементы из буфера не будут переданы в канал.
Это означает, что при попытке присоединиться к такому процессу может возникнуть взаимная блокировка, если вы не уверены, что все помещённые в очередь элементы были обработаны. Аналогично, если дочерний процесс не является демоническим, родительский процесс может зависнуть при завершении, когда попытается присоединиться ко всем своим недемоническим дочерним процессам.
Обратите внимание, что у очереди, созданной с помощью менеджера, такой проблемы нет. См. раздел Рекомендации по программированию.
Пример использования очередей для межпроцессного взаимодействия см. в разделе Примеры.
-
multiprocessing.Pipe(duplex=True) -
Возвращает пару
(conn1, conn2)объектовConnection, представляющих концы канала.Если duplex равно
True(значение по умолчанию), канал является двунаправленным. Если duplex равноFalse, канал является однонаправленным:conn1можно использовать только для получения сообщений, аconn2— только для их отправки.Метод
send()сериализует объект с помощьюpickle, аrecv()воссоздаёт объект.
-
class multiprocessing.Queue([maxsize]) -
Возвращает очередь, совместно используемую процессами и реализованную с помощью канала и нескольких блокировок/семафоров. Когда процесс впервые помещает элемент в очередь, запускается поток-поставщик, который передаёт объекты из буфера в канал.
Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. раздел Глобальный метод запуска.
Для сигнализации о тайм-ауте вызываются стандартные исключения
queue.Emptyиqueue.Fullиз модуляqueueстандартной библиотеки.Queueреализует все методыqueue.Queue, кромеtask_done(),join()иshutdown().-
qsize() -
Возвращает приблизительный размер очереди. Из-за особенностей многопоточности и многопроцессности это число ненадёжно.
Обратите внимание, что на таких платформах, как macOS, где
sem_getvalue()не реализован, этот метод может вызвать исключениеNotImplementedError.
-
empty() -
Возвращает
True, если очередь пуста, иFalseв противном случае. Из-за особенностей многопоточности и многопроцессности этот результат ненадёжен.Для закрытых очередей может вызвать исключение
OSError(не гарантируется).
-
full() -
Возвращает
True, если очередь заполнена, иFalseв противном случае. Из-за особенностей многопоточности и многопроцессности этот результат ненадёжен.
-
put(obj[, block[, timeout]]) -
Помещает obj в очередь. Если необязательный аргумент block равен
True(значение по умолчанию), а timeout равенNone(значение по умолчанию), метод при необходимости блокируется до появления свободного места. Если timeout — положительное число, метод блокируется не более timeout секунд и вызывает исключениеqueue.Full, если за это время свободное место не появилось. В противном случае (block равенFalse) элемент помещается в очередь, если свободное место доступно немедленно; иначе вызывается исключениеqueue.Full(timeout в этом случае игнорируется).Изменено в версии 3.8: Если очередь закрыта, вместо
AssertionErrorвызываетсяValueError.
-
put_nowait(obj) -
Эквивалентно
put(obj, False).
-
get([block[, timeout]]) -
Извлекает и возвращает элемент из очереди. Если необязательный аргумент block равен
True(значение по умолчанию), а timeout равенNone(значение по умолчанию), метод при необходимости блокируется до появления элемента. Если timeout — положительное число, метод блокируется не более timeout секунд и вызывает исключениеqueue.Empty, если за это время элемент не появился. В противном случае (block равенFalse) метод возвращает элемент, если он доступен немедленно; иначе вызывается исключениеqueue.Empty(timeout в этом случае игнорируется).Изменено в версии 3.8: Если очередь закрыта, вместо
OSErrorвызываетсяValueError.
-
get_nowait() -
Эквивалентно
get(False).
У
multiprocessing.Queueесть несколько дополнительных методов, отсутствующих вqueue.Queue. Обычно эти методы не нужны в большинстве программ:-
close() -
Закрывает очередь и освобождает внутренние ресурсы.
После закрытия очередь больше нельзя использовать. Например, больше нельзя вызывать методы
get(),put()иempty().Фоновый поток завершит работу после передачи всех данных из буфера в канал. Этот метод вызывается автоматически при сборке мусора для очереди.
-
join_thread() -
Ожидает завершения фонового потока. Этот метод можно использовать только после вызова
close(). Он блокируется до завершения фонового потока, гарантируя, что все данные из буфера переданы в канал.По умолчанию при завершении процесса, который не создал очередь, процесс попытается дождаться завершения фонового потока очереди. Процесс может вызвать
cancel_join_thread(), чтобыjoin_thread()ничего не делал.
-
cancel_join_thread() -
Предотвращает блокировку при вызове
join_thread(). В частности, это предотвращает автоматическое ожидание завершения фонового потока при выходе из процесса — см.join_thread().Более подходящим названием для этого метода могло бы быть
allow_exit_without_flush(). При его использовании данные в очереди, вероятно, будут потеряны, и вам почти наверняка не понадобится его вызывать. Он предназначен лишь для случаев, когда текущему процессу нужно немедленно завершиться, не дожидаясь передачи данных из очереди в нижележащий канал, и потеря данных не имеет значения.
Примечание
Для работы этого класса необходима действующая реализация семафоров совместного использования в операционной системе хоста. Если такая реализация отсутствует, функциональность класса будет отключена, а попытка создать экземпляр
Queueприведёт к исключениюImportError. Дополнительные сведения см. в bpo-3770. То же относится ко всем перечисленным ниже специализированным типам очередей. -
-
class multiprocessing.SimpleQueue -
Это упрощённый тип
Queue, очень близкий к каналуPipeс блокировкой.Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. раздел Глобальный метод запуска.
-
close() -
Закрывает очередь и освобождает внутренние ресурсы.
После закрытия очередь больше нельзя использовать. Например, больше нельзя вызывать методы
get(),put()иempty().Добавлено в версии 3.9.
-
empty() -
Возвращает
True, если очередь пуста, иFalseв противном случае.Если SimpleQueue закрыта, всегда вызывает исключение
OSError.
-
get() -
Извлекает и возвращает элемент из очереди.
-
put(item) -
Помещает item в очередь.
-
-
class multiprocessing.JoinableQueue([maxsize]) -
JoinableQueue— подклассQueue, представляющий собой очередь, которая дополнительно имеет методыtask_done()иjoin().Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. раздел Глобальный метод запуска.
-
task_done() -
Указывает, что ранее помещённая в очередь задача выполнена. Используется потребителями очереди. Для каждого вызова
get(), использованного для получения задачи, последующий вызовtask_done()сообщает очереди, что обработка задачи завершена.Если вызов
join()в данный момент заблокирован, он возобновит работу, когда все элементы будут обработаны (то есть когда для каждого элемента, помещённого в очередь с помощьюput(), будет получен вызовtask_done()).Вызывает исключение
ValueError, если вызвано больше раз, чем элементов было помещено в очередь.
-
join() -
Блокируется до тех пор, пока все элементы очереди не будут получены и обработаны.
Счётчик незавершённых задач увеличивается при добавлении элемента в очередь. Он уменьшается, когда потребитель вызывает
task_done(), сообщая, что элемент получен и вся работа с ним завершена. Когда счётчик незавершённых задач достигает нуля, вызовjoin()разблокируется.
-
Разное
-
multiprocessing.active_children() -
Возвращает список всех работающих дочерних процессов текущего процесса.
Вызов этой функции имеет побочный эффект: она «присоединяет» все процессы, которые уже завершились.
-
multiprocessing.cpu_count() -
Возвращает количество процессоров в системе.
Это число не эквивалентно количеству процессоров, которые может использовать текущий процесс. Количество доступных процессоров можно получить с помощью
os.process_cpu_count()(илиlen(os.sched_getaffinity(0))).Если количество процессоров не удаётся определить, возникает исключение
NotImplementedError.См. также
Изменено в версии 3.13: Возвращаемое значение также можно переопределить с помощью флага
-X cpu_countилиPYTHON_CPU_COUNT, поскольку эта функция лишь обёртывает API подсчёта процессоровos.
-
multiprocessing.current_process() -
Возвращает объект
Process, соответствующий текущему процессу.Аналог
threading.current_thread().
-
multiprocessing.parent_process() -
Возвращает объект
Process, соответствующий родительскому процессуcurrent_process(). Для главного процессаparent_processбудетNone.Добавлено в версии 3.8.
-
multiprocessing.freeze_support() -
Добавляет поддержку программ, использующих
multiprocessing, которые были заморожены для создания исполняемого файла. (Проверено с py2exe, PyInstaller и cx_Freeze.)Эту функцию необходимо вызвать сразу после строки
if __name__ == '__main__'главного модуля. Например:from multiprocessing import Process, freeze_support def f(): print('hello world!') if __name__ == '__main__': freeze_support() Process(target=f).start()Если строка
freeze_support()опущена, попытка запустить замороженный исполняемый файл вызовет исключениеRuntimeError.Вызов
freeze_support()не оказывает влияния, если метод запуска не равен spawn. Кроме того, если модуль запускается обычным образом интерпретатором Python (программа не была заморожена),freeze_support()не оказывает влияния.
-
multiprocessing.get_all_start_methods() -
Возвращает список поддерживаемых методов запуска; первый из них используется по умолчанию. Возможные методы запуска:
'fork','spawn'и'forkserver'. Не все платформы поддерживают все методы. См. раздел Контексты и методы запуска.Добавлено в версии 3.4.
-
multiprocessing.get_context(method=None) -
Возвращает объект контекста, обладающий теми же атрибутами, что и модуль
multiprocessing.Если method равен
None, возвращается контекст по умолчанию. Обратите внимание: если глобальный метод запуска не был задан, эта функция задаст для него системное значение по умолчанию. Подробнее см. раздел Глобальный метод запуска. В противном случае method должен быть равен'fork','spawn'или'forkserver'. Если указанный метод запуска недоступен, возникает исключениеValueError. См. раздел Контексты и методы запуска.Добавлено в версии 3.4.
-
multiprocessing.get_start_method(allow_none=False) -
Возвращает название метода запуска, используемого для запуска процессов.
Если глобальный метод запуска не задан и allow_none равен
False, для глобального метода запуска устанавливается значение по умолчанию и возвращается его название. Подробнее см. раздел Глобальный метод запуска.Возвращаемым значением может быть
'fork','spawn','forkserver'илиNone. См. раздел Контексты и методы запуска.Добавлено в версии 3.4.
Изменено в версии 3.8: В macOS метод запуска spawn теперь используется по умолчанию. Метод запуска fork следует считать небезопасным, так как он может привести к аварийному завершению подпроцесса. См. bpo-33725.
-
multiprocessing.set_executable(executable) -
Задаёт путь к интерпретатору Python, который будет использоваться при запуске дочернего процесса. (По умолчанию используется
sys.executable.) Для встраивания интерпретатора, вероятно, потребуется выполнить что-то вроде следующего:set_executable(os.path.join(sys.exec_prefix, 'pythonw.exe'))
прежде чем создавать дочерние процессы.
Изменено в версии 3.4: Теперь поддерживается в POSIX при использовании метода запуска
'spawn'.Изменено в версии 3.11: Принимает объект, подобный пути.
-
multiprocessing.set_forkserver_preload(module_names) -
Задаёт список имён модулей, которые основной процесс forkserver должен попытаться импортировать, чтобы процессы, созданные с помощью fork, унаследовали состояние уже импортированных модулей. Все возникающие при этом исключения
ImportErrorигнорируются без уведомления. Это можно использовать для повышения производительности, чтобы избежать повторения одних и тех же действий в каждом процессе.Чтобы это работало, функцию необходимо вызвать до запуска процесса forkserver (до создания
Poolили запускаProcess).Имеет смысл только при использовании метода запуска
'forkserver'. См. раздел Контексты и методы запуска.Добавлено в версии 3.4.
-
multiprocessing.set_start_method(method, force=False) -
Задаёт метод запуска дочерних процессов. Аргумент method может быть равен
'fork','spawn'или'forkserver'. Если метод запуска уже задан, а force не равенTrue, возникает исключениеRuntimeError. Если method равенNoneи force равенTrue, метод запуска устанавливается вNone. Если method равенNoneи force равенFalse, для контекста устанавливается значение по умолчанию.Обратите внимание: эту функцию следует вызывать не более одного раза, и её вызов должен быть защищён условием
if __name__ == '__main__'главного модуля.См. раздел Контексты и методы запуска.
Добавлено в версии 3.4.
Примечание
В multiprocessing нет аналогов для threading.active_count(), threading.enumerate(), threading.settrace(), threading.setprofile(), threading.Timer или threading.local.
Объекты соединения
Объекты соединения позволяют отправлять и получать объекты, допускающие упаковку с помощью pickle, или строки. Их можно рассматривать как соединённые сокеты, ориентированные на сообщения.
Объекты соединения обычно создаются с помощью Pipe — см. также раздел Слушатели и клиенты.
-
class multiprocessing.connection.Connection -
-
send(obj) -
Отправляет объект на другой конец соединения; его следует прочитать с помощью
recv().Объект должен допускать упаковку с помощью pickle. Для очень больших объектов pickle (примерно от 32 МиБ, хотя это зависит от ОС) может возникнуть исключение
ValueError.
-
recv() -
Возвращает объект, отправленный с другого конца соединения с помощью
send(). Блокирует выполнение до получения данных. Если получать больше нечего и другой конец соединения закрыт, возникает исключениеEOFError.
-
fileno() -
Возвращает файловый дескриптор или дескриптор, используемый соединением.
-
close() -
Закрывает соединение.
Этот метод вызывается автоматически при сборке мусора для соединения.
-
poll([timeout]) -
Возвращает значение, указывающее, доступны ли данные для чтения.
Если timeout не указан, функция возвращает результат немедленно. Если timeout — число, оно задаёт максимальное время блокировки в секундах. Если timeout равен
None, используется бесконечное время ожидания.Обратите внимание: одновременно опрашивать несколько объектов соединения можно с помощью
multiprocessing.connection.wait().
-
send_bytes(buf[, offset[, size]]) -
Отправляет байтовые данные из объекта, подобного bytes, в виде целого сообщения.
Если задан offset, данные считываются из этой позиции в buf. Если задан size, из buf считывается указанное количество байтов. Для очень больших буферов (примерно от 32 МиБ, хотя это зависит от ОС) может возникнуть исключение
ValueError.
-
recv_bytes([maxlength]) -
Возвращает в виде строки полное сообщение с байтовыми данными, отправленное с другого конца соединения. Блокирует выполнение до получения данных. Если получать больше нечего и другой конец соединения закрыт, возникает исключение
EOFError.Если задан maxlength и длина сообщения превышает maxlength, возникает исключение
OSError, после чего соединение больше нельзя будет читать.
-
recv_bytes_into(buf[, offset]) -
Считывает в buf полное сообщение с байтовыми данными, отправленное с другого конца соединения, и возвращает количество байтов в сообщении. Блокирует выполнение до получения данных. Если получать больше нечего и другой конец соединения закрыт, возникает исключение
EOFError.buf должен быть доступным для записи объектом, подобным bytes. Если задан offset, сообщение записывается в буфер начиная с этой позиции. Смещение должно быть неотрицательным целым числом, меньшим длины buf (в байтах).
Если буфер слишком короткий, возникает исключение
BufferTooShort, а полное сообщение становится доступным какe.args[0], гдеe— экземпляр исключения.
Изменено в версии 3.3: Теперь сами объекты соединения можно передавать между процессами с помощью
Connection.send()иConnection.recv().Теперь объекты соединения также поддерживают протокол управления контекстом — см. раздел Типы менеджеров контекста.
__enter__()возвращает объект соединения, а__exit__()вызываетclose(). -
Например:
>>> from multiprocessing import Pipe
>>> a, b = Pipe()
>>> a.send([1, 'hello', None])
>>> b.recv()
[1, 'hello', None]
>>> b.send_bytes(b'thank you')
>>> a.recv_bytes()
b'thank you'
>>> import array
>>> arr1 = array.array('i', range(5))
>>> arr2 = array.array('i', [0] * 10)
>>> a.send_bytes(arr1)
>>> count = b.recv_bytes_into(arr2)
>>> assert count == len(arr1) * arr1.itemsize
>>> arr2
array('i', [0, 1, 2, 3, 4, 0, 0, 0, 0, 0])
Предупреждение
Метод Connection.recv() автоматически распаковывает полученные данные с помощью pickle, что может представлять угрозу безопасности, если вы не доверяете процессу, отправившему сообщение.
Поэтому, если объект соединения не был создан с помощью Pipe(), методы recv() и send() следует использовать только после выполнения какой-либо аутентификации. См. раздел Ключи аутентификации.
Предупреждение
Если процесс завершается принудительно при попытке чтения из канала или записи в него, данные в канале, скорее всего, будут повреждены, поскольку определить границы сообщений может оказаться невозможно.
Примитивы синхронизации
В общем случае примитивы синхронизации не так необходимы в многопроцессной программе, как в многопоточной. См. документацию по модулю threading.
Обратите внимание, что примитивы синхронизации также можно создавать с помощью объекта-менеджера — см. Менеджеры.
-
class multiprocessing.Barrier(parties[, action[, timeout]]) -
Объект барьера: копия
threading.Barrier.Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.
Добавлено в версии 3.3.
-
class multiprocessing.BoundedSemaphore([value]) -
Объект ограниченного семафора: близкий аналог
threading.BoundedSemaphore.Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.
Есть лишь одно отличие от близкого аналога: первый аргумент метода
acquireназывается block, как и вLock.acquire().-
locked() -
Возвращает логическое значение, указывающее, заблокирован ли объект в данный момент.
Добавлено в версии 3.14.
Примечание
В macOS этот класс неотличим от
Semaphore, посколькуsem_getvalue()на этой платформе не реализован. -
-
class multiprocessing.Condition([lock]) -
Условная переменная: псевдоним для
threading.Condition.Если указан lock, это должен быть объект
LockилиRLockизmultiprocessing.Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.
Изменено в версии 3.3: Добавлен метод
wait_for().
-
class multiprocessing.Event -
Копия
threading.Event.Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.
-
class multiprocessing.Lock -
Нерекурсивный объект блокировки: близкий аналог
threading.Lock. После того как процесс или поток захватил блокировку, последующие попытки захватить её из любого процесса или потока будут блокироваться до её освобождения; освободить её может любой процесс или поток. Понятия и поведениеthreading.Lockприменительно к потокам здесь воспроизводятся вmultiprocessing.Lockприменительно как к процессам, так и к потокам, за исключением отмеченных случаев.Обратите внимание, что
Lockна самом деле является фабричной функцией, возвращающей экземплярmultiprocessing.synchronize.Lock, инициализированный с контекстом по умолчанию.Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.
Lockподдерживает протокол менеджера контекста и поэтому может использоваться в инструкцияхwith.-
acquire(block=True, timeout=None) -
Захватывает блокировку в блокирующем или неблокирующем режиме.
Если аргументу block задано значение
True(по умолчанию), вызов метода блокируется, пока блокировка не перейдёт в незаблокированное состояние, затем переводит её в заблокированное состояние и возвращаетTrue. Обратите внимание, что имя первого аргумента отличается от имени вthreading.Lock.acquire().Если аргументу block задано значение
False, вызов метода не блокируется. Если блокировка в данный момент заблокирована, возвращаетсяFalse; в противном случае блокировка переводится в заблокированное состояние и возвращаетсяTrue.При вызове с положительным значением timeout типа с плавающей точкой метод блокируется не более указанного в timeout количества секунд, пока не удастся захватить блокировку. Вызовы с отрицательным значением timeout эквивалентны вызовам с нулевым значением timeout. Вызовы со значением timeout
None(по умолчанию) устанавливают бесконечный период ожидания. Обратите внимание, что обработка отрицательных значений или значенийNoneдля timeout отличается от реализованного поведения вthreading.Lock.acquire(). Аргумент timeout не имеет практического значения, если для аргумента block заданоFalse, и в этом случае игнорируется. ВозвращаетTrue, если блокировка захвачена, илиFalse, если время ожидания истекло.
-
release() -
Освобождает блокировку. Этот метод можно вызвать из любого процесса или потока, а не только из того процесса или потока, который изначально захватил блокировку.
Поведение такое же, как у
threading.Lock.release(), за исключением того, что при вызове для незаблокированной блокировки возникает исключениеValueError.
-
locked() -
Возвращает логическое значение, указывающее, заблокирован ли объект в данный момент.
Добавлено в версии 3.14.
-
-
class multiprocessing.RLock -
Рекурсивный объект блокировки: близкий аналог
threading.RLock. Рекурсивная блокировка должна быть освобождена тем процессом или потоком, который её захватил. После того как процесс или поток захватил рекурсивную блокировку, тот же процесс или поток может захватить её снова без блокировки; этот процесс или поток должен освободить её столько раз, сколько он её захватывал.Обратите внимание, что
RLockна самом деле является фабричной функцией, возвращающей экземплярmultiprocessing.synchronize.RLock, инициализированный с контекстом по умолчанию.Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.
RLockподдерживает протокол менеджера контекста и поэтому может использоваться в инструкцияхwith.-
acquire(block=True, timeout=None) -
Захватывает блокировку в блокирующем или неблокирующем режиме.
При вызове со значением
Trueдля аргумента block метод блокируется, пока блокировка не перейдёт в незаблокированное состояние (то есть не будет принадлежать ни одному процессу или потоку), если только она уже не принадлежит текущему процессу или потоку. Затем текущий процесс или поток становится владельцем блокировки (если ещё не был им), а уровень рекурсии блокировки увеличивается на единицу; в результате возвращаетсяTrue. Обратите внимание, что поведение этого первого аргумента отличается от поведения в реализацииthreading.RLock.acquire()несколькими способами, начиная с самого имени аргумента.При вызове со значением
Falseдля аргумента block блокировка не выполняется. Если блокировка уже захвачена (и, следовательно, принадлежит) другому процессу или потоку, текущий процесс или поток не становится её владельцем, уровень рекурсии блокировки не изменяется, и возвращаетсяFalse. Если блокировка находится в незаблокированном состоянии, текущий процесс или поток становится её владельцем, уровень рекурсии увеличивается, и возвращаетсяTrue.Использование аргумента timeout и его поведение совпадают с описанными для
Lock.acquire(). Обратите внимание, что некоторые особенности поведения timeout отличаются от реализованного поведения вthreading.RLock.acquire().
-
release() -
Освобождает блокировку, уменьшая уровень рекурсии. Если после уменьшения уровень рекурсии равен нулю, блокировка переводится в незаблокированное состояние (то есть больше не принадлежит ни одному процессу или потоку). Если другие процессы или потоки заблокированы в ожидании разблокировки, продолжить работу разрешается ровно одному из них. Если после уменьшения уровень рекурсии остаётся ненулевым, блокировка остаётся заблокированной и принадлежит вызвавшему метод процессу или потоку.
Вызывайте этот метод только в том случае, если блокировка принадлежит вызывающему процессу или потоку. Если метод вызван процессом или потоком, отличным от владельца, либо блокировка находится в незаблокированном состоянии (не принадлежит никому), возникает исключение
AssertionError. Обратите внимание, что тип исключения, возникающего в такой ситуации, отличается от реализованного поведения вthreading.RLock.release().
-
locked() -
Возвращает логическое значение, указывающее, заблокирован ли объект в данный момент.
Добавлено в версии 3.14.
-
-
class multiprocessing.Semaphore([value]) -
Объект семафора: близкий аналог
threading.Semaphore.Создание экземпляра этого класса может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.
Есть лишь одно отличие от близкого аналога: первый аргумент метода
acquireназывается block, как и вLock.acquire().-
get_value() -
Возвращает текущее значение семафора.
Обратите внимание, что на таких платформах, как macOS, где
sem_getvalue()не реализован, может возникнуть исключениеNotImplementedError.
-
locked() -
Возвращает логическое значение, указывающее, заблокирован ли объект в данный момент.
Добавлено в версии 3.14.
-
Примечание
В macOS функция sem_timedwait не поддерживается, поэтому вызов acquire() с тайм-аутом имитирует поведение этой функции с помощью цикла с задержкой.
Примечание
Для работы некоторых функций этого пакета необходима реализация общих семафоров в операционной системе хоста. Если такой реализации нет, модуль multiprocessing.synchronize будет отключён, а попытка его импорта приведёт к возникновению исключения ImportError. Дополнительные сведения см. в bpo-3770.
Менеджеры
Менеджеры позволяют создавать данные, которыми можно совместно пользоваться в разных процессах, в том числе через сеть между процессами, работающими на разных машинах. Объект менеджера управляет серверным процессом, который обслуживает совместно используемые объекты. Другие процессы могут получать доступ к совместно используемым объектам с помощью прокси.
-
multiprocessing.Manager() -
Возвращает запущенный объект
SyncManager, который можно использовать для совместного использования объектов между процессами. Возвращённый объект менеджера соответствует порождённому дочернему процессу и имеет методы, создающие совместно используемые объекты и возвращающие соответствующие прокси.
Процессы менеджера завершаются, как только они собираются сборщиком мусора или завершается их родительский процесс. Классы менеджеров определены в модуле multiprocessing.managers:
-
class multiprocessing.managers.BaseManager(address=None, authkey=None, serializer='pickle', ctx=None, *, shutdown_timeout=1.0) -
Создаёт объект BaseManager.
После создания следует вызвать
start()илиget_server().serve_forever(), чтобы объект менеджера ссылался на запущенный процесс менеджера.address — адрес, на котором процесс менеджера ожидает новых подключений. Если address равен
None, выбирается произвольный адрес.authkey — ключ аутентификации, используемый для проверки допустимости входящих подключений к серверному процессу. Если authkey равен
None, используетсяcurrent_process().authkey. В противном случае используется authkey, который должен быть байтовой строкой.serializer должен быть равен
'pickle'(используется сериализацияpickle) или'xmlrpclib'(используется сериализацияxmlrpc.client).ctx — объект контекста или
None(используется текущий контекст). Если заданоNone, вызов может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.shutdown_timeout — время ожидания в секундах, в течение которого процесс, используемый менеджером, должен завершиться при вызове метода
shutdown(). Если время ожидания завершения истекает, процесс принудительно завершается. Если и принудительное завершение процесса занимает больше отведённого времени, процесс уничтожается.Изменено в версии 3.11: Добавлен параметр shutdown_timeout.
-
start([initializer[, initargs]]) -
Запускает подпроцесс для запуска менеджера. Если initializer не равен
None, при запуске подпроцесс вызоветinitializer(*initargs).
-
get_server() -
Возвращает объект
Server, представляющий реальный сервер, которым управляет менеджер. ОбъектServerподдерживает методserve_forever():>>> from multiprocessing.managers import BaseManager >>> manager = BaseManager(address=('', 50000), authkey=b'abc') >>> server = manager.get_server() >>> server.serve_forever()Кроме того, у
Serverесть атрибутaddress.
-
connect() -
Подключает локальный объект менеджера к удалённому процессу менеджера:
>>> from multiprocessing.managers import BaseManager >>> m = BaseManager(address=('127.0.0.1', 50000), authkey=b'abc') >>> m.connect()
-
shutdown() -
Останавливает процесс, используемый менеджером. Этот метод доступен только в том случае, если для запуска серверного процесса использовался
start().Метод можно вызывать несколько раз.
-
register(typeid[, callable[, proxytype[, exposed[, method_to_typeid[, create_method]]]]]) -
Метод класса, который можно использовать для регистрации типа или вызываемого объекта в классе менеджера.
typeid — «идентификатор типа», используемый для обозначения конкретного типа совместно используемого объекта. Должен быть строкой.
callable — вызываемый объект, создающий объекты для этого идентификатора типа. Если экземпляр менеджера будет подключён к серверу с помощью метода
connect()или если аргумент create_method равенFalse, этот аргумент можно оставить равнымNone.proxytype — подкласс
BaseProxy, используемый для создания прокси совместно используемых объектов с этим typeid. Если заданоNone, класс прокси создаётся автоматически.exposed используется для указания последовательности имён методов, к которым прокси этого typeid разрешено обращаться с помощью
BaseProxy._callmethod(). (Если exposed равенNone, вместо него используетсяproxytype._exposed_, если такой атрибут существует.) Если список доступных методов не указан, доступны все «открытые методы» совместно используемого объекта. (Здесь «открытый метод» означает любой атрибут, у которого есть метод__call__()и имя которого не начинается с'_'.)method_to_typeid — соответствие, используемое для указания типа возвращаемого значения тех доступных методов, которые должны возвращать прокси. Оно сопоставляет имена методов со строками typeid. (Если method_to_typeid равен
None, вместо него используетсяproxytype._method_to_typeid_, если такой атрибут существует.) Если имя метода не является ключом этого соответствия или соответствие равноNone, объект, возвращённый методом, будет скопирован по значению.create_method определяет, следует ли создать метод с именем typeid, который можно использовать, чтобы указать серверному процессу создать новый совместно используемый объект и вернуть прокси для него. По умолчанию он равен
True.
У экземпляров
BaseManagerтакже есть одно свойство только для чтения:-
address -
Адрес, используемый менеджером.
Изменено в версии 3.3: Объекты менеджера поддерживают протокол управления контекстом — см. Типы менеджеров контекста.
__enter__()запускает серверный процесс (если он ещё не запущен), а затем возвращает объект менеджера.__exit__()вызываетshutdown().В предыдущих версиях
__enter__()не запускал серверный процесс менеджера, если тот ещё не был запущен. -
-
class multiprocessing.managers.SyncManager -
Подкласс
BaseManager, предназначенный для синхронизации процессов. Объекты этого типа возвращаются функциейmultiprocessing.Manager().Его методы создают и возвращают объекты-прокси для ряда часто используемых типов данных, синхронизируемых между процессами. В частности, к ним относятся совместно используемые списки и словари.
-
Barrier(parties[, action[, timeout]]) -
Создаёт совместно используемый объект
threading.Barrierи возвращает прокси для него.Добавлено в версии 3.3.
-
BoundedSemaphore([value]) -
Создаёт совместно используемый объект
threading.BoundedSemaphoreи возвращает прокси для него.
-
Condition([lock]) -
Создаёт совместно используемый объект
threading.Conditionи возвращает прокси для него.Если передан аргумент lock, он должен быть прокси объекта
threading.Lockилиthreading.RLock.Изменено в версии 3.3: Добавлен метод
wait_for().
-
Event() -
Создаёт совместно используемый объект
threading.Eventи возвращает прокси для него.
-
Lock() -
Создаёт совместно используемый объект
threading.Lockи возвращает прокси для него.
-
Namespace() -
Создаёт совместно используемый объект
Namespaceи возвращает прокси для него.
-
Queue([maxsize]) -
Создаёт совместно используемый объект
queue.Queueи возвращает прокси для него.
-
RLock() -
Создаёт совместно используемый объект
threading.RLockи возвращает прокси для него.
-
Semaphore([value]) -
Создаёт совместно используемый объект
threading.Semaphoreи возвращает прокси для него.
-
Array(typecode, sequence) -
Создаёт массив и возвращает прокси для него.
-
Value(typecode, value) -
Создаёт объект с доступным для записи атрибутом
valueи возвращает прокси для него.
-
dict() - dict(mapping)
- dict(sequence)
-
Создаёт совместно используемый объект
dictи возвращает прокси для него.
-
list() - list(sequence)
-
Создаёт совместно используемый объект
listи возвращает прокси для него.
-
set() - set(sequence)
- set(mapping)
-
Создаёт совместно используемый объект
setи возвращает прокси для него.Добавлено в версии 3.14: Добавлена поддержка
set.
Изменено в версии 3.6: Совместно используемые объекты могут быть вложенными. Например, совместно используемый контейнер, такой как общий список, может содержать другие совместно используемые объекты, которыми также будет управлять и которые будет синхронизировать
SyncManager. -
-
class multiprocessing.managers.Namespace -
Тип, который можно зарегистрировать в
SyncManager.Объект пространства имён не имеет открытых методов, но содержит доступные для записи атрибуты. Его представление показывает значения этих атрибутов.
Однако при использовании прокси для объекта пространства имён атрибут, начинающийся с
'_', будет атрибутом прокси, а не объекта, на который оно ссылается:>>> mp_context = multiprocessing.get_context('spawn') >>> manager = mp_context.Manager() >>> Global = manager.Namespace() >>> Global.x = 10 >>> Global.y = 'hello' >>> Global._z = 12.3 # this is an attribute of the proxy >>> print(Global) Namespace(x=10, y='hello')
Настраиваемые менеджеры
Чтобы создать собственный менеджер, нужно создать подкласс BaseManager и использовать метод класса register() для регистрации новых типов или вызываемых объектов в классе менеджера. Например:
from multiprocessing.managers import BaseManager
class MathsClass:
def add(self, x, y):
return x + y
def mul(self, x, y):
return x * y
class MyManager(BaseManager):
pass
MyManager.register('Maths', MathsClass)
if __name__ == '__main__':
with MyManager() as manager:
maths = manager.Maths()
print(maths.add(4, 3)) # prints 7
print(maths.mul(7, 8)) # prints 56
Использование удалённого менеджера
Можно запустить сервер менеджера на одной машине, а клиентам — использовать его с других машин (если это допускают соответствующие брандмауэры).
Следующие команды создают сервер для одной совместно используемой очереди, к которой могут обращаться удалённые клиенты:
>>> from multiprocessing.managers import BaseManager
>>> from queue import Queue
>>> queue = Queue()
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue', callable=lambda:queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()
Один клиент может обратиться к серверу следующим образом:
>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.put('hello')
Эту очередь может использовать и другой клиент:
>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.get()
'hello'
К этой очереди могут обращаться и локальные процессы, используя приведённый выше код клиента для удалённого доступа:
>>> from multiprocessing import Process, Queue
>>> from multiprocessing.managers import BaseManager
>>> class Worker(Process):
... def __init__(self, q):
... self.q = q
... super().__init__()
... def run(self):
... self.q.put('local hello')
...
>>> queue = Queue()
>>> w = Worker(queue)
>>> w.start()
>>> class QueueManager(BaseManager): pass
...
>>> QueueManager.register('get_queue', callable=lambda: queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()
Объекты-прокси
Прокси — это объект, который ссылается на совместно используемый объект, находящийся (предположительно) в другом процессе. Совместно используемый объект называется целевым объектом прокси. У нескольких объектов-прокси может быть один и тот же целевой объект.
Объект-прокси имеет методы, вызывающие соответствующие методы целевого объекта (однако через прокси могут быть доступны не все методы целевого объекта). Таким образом, прокси можно использовать так же, как и сам целевой объект:
>>> mp_context = multiprocessing.get_context('spawn')
>>> manager = mp_context.Manager()
>>> l = manager.list([i*i for i in range(10)])
>>> print(l)
[0, 1, 4, 9, 16, 25, 36, 49, 64, 81]
>>> print(repr(l))
<ListProxy object, typeid 'list' at 0x...>
>>> l[4]
16
>>> l[2:5]
[4, 9, 16]
Обратите внимание: применение str() к прокси возвращает представление целевого объекта, тогда как применение repr() возвращает представление прокси.
Важное свойство объектов-прокси — возможность их сериализации с помощью pickle, благодаря чему их можно передавать между процессами. Поэтому целевой объект может содержать объекты-прокси. Это позволяет создавать вложенные управляемые списки, словари и другие объекты-прокси:
>>> a = manager.list()
>>> b = manager.list()
>>> a.append(b) # referent of a now contains referent of b
>>> print(a, b)
[<ListProxy object, typeid 'list' at ...>] []
>>> b.append('hello')
>>> print(a[0], b)
['hello'] ['hello']
Аналогично, прокси словарей и списков можно вкладывать друг в друга:
>>> l_outer = manager.list([ manager.dict() for i in range(2) ])
>>> d_first_inner = l_outer[0]
>>> d_first_inner['a'] = 1
>>> d_first_inner['b'] = 2
>>> l_outer[1]['c'] = 3
>>> l_outer[1]['z'] = 26
>>> print(l_outer[0])
{'a': 1, 'b': 2}
>>> print(l_outer[1])
{'c': 3, 'z': 26}
Если целевой объект содержит обычные (не прокси) объекты list или dict, изменения этих изменяемых значений не будут передаваться через менеджер, поскольку прокси не может определить, когда изменяются содержащиеся в них значения. Однако запись значения в прокси-контейнер (которая вызывает __setitem__ для объекта-прокси) передаёт изменение через менеджер. Поэтому, чтобы эффективно изменить такой элемент, можно присвоить изменённое значение прокси-контейнеру повторно:
# create a list proxy and append a mutable object (a dictionary)
lproxy = manager.list()
lproxy.append({})
# now mutate the dictionary
d = lproxy[0]
d['a'] = 1
d['b'] = 2
# at this point, the changes to d are not yet synced, but by
# updating the dictionary, the proxy is notified of the change
lproxy[0] = d
Для большинства случаев такой способ, возможно, менее удобен, чем использование вложенных объектов-прокси, но он также демонстрирует возможность управлять синхронизацией.
Примечание
Типы прокси в multiprocessing не поддерживают сравнение по значению. Например:
>>> manager.list([1,2,3]) == [1,2,3] False
При сравнении следует использовать копию целевого объекта.
-
class multiprocessing.managers.BaseProxy -
Объекты-прокси являются экземплярами подклассов
BaseProxy.-
_callmethod(methodname[, args[, kwds]]) -
Вызывает метод целевого объекта прокси и возвращает его результат.
Если
proxy— прокси, целевым объектом которого являетсяobj, то выражениеproxy._callmethod(methodname, args, kwds)
вычислит выражение
getattr(obj, methodname)(*args, **kwds)
в процессе менеджера.
Возвращённое значение будет копией результата вызова или прокси нового совместно используемого объекта — см. документацию аргумента method_to_typeid метода
BaseManager.register().Если при вызове возникает исключение,
_callmethod()повторно вызовет его. Если в процессе менеджера возникает другое исключение, оно преобразуется в исключениеRemoteErrorи вызывается методом_callmethod().Обратите особое внимание: исключение будет вызвано, если methodname не был открыт.
Пример использования
_callmethod():>>> l = manager.list(range(10)) >>> l._callmethod('__len__') 10 >>> l._callmethod('__getitem__', (slice(2, 7),)) # equivalent to l[2:7] [2, 3, 4, 5, 6] >>> l._callmethod('__getitem__', (20,)) # equivalent to l[20] Traceback (most recent call last): ... IndexError: list index out of range
-
_getvalue() -
Возвращает копию целевого объекта.
Если целевой объект нельзя сериализовать с помощью pickle, будет вызвано исключение.
-
__repr__() -
Возвращает представление объекта-прокси.
-
__str__() -
Возвращает представление целевого объекта.
-
Очистка
Объект-прокси использует обратный вызов weakref, чтобы при сборке объекта сборщиком мусора он удалял свою регистрацию из менеджера, которому принадлежит его целевой объект.
Совместно используемый объект удаляется из процесса менеджера, когда на него больше не ссылается ни один прокси.
Пулы процессов
Можно создать пул процессов, которые будут выполнять отправленные ему задачи, с помощью класса Pool.
-
class multiprocessing.pool.Pool([processes[, initializer[, initargs[, maxtasksperchild[, context]]]]]) -
Объект пула процессов, который управляет пулом рабочих процессов, которым можно отправлять задания. Он поддерживает асинхронные результаты с тайм-аутами и обратными вызовами, а также параллельную реализацию map.
processes — количество рабочих процессов. Если processes равно
None, используется число, возвращаемоеos.process_cpu_count().Если initializer не равно
None, каждый рабочий процесс вызоветinitializer(*initargs)при запуске.maxtasksperchild — количество задач, которые рабочий процесс может выполнить до завершения и замены новым рабочим процессом, что позволяет освободить неиспользуемые ресурсы. По умолчанию maxtasksperchild равно
None, то есть рабочие процессы будут работать, пока существует пул.Параметр context можно использовать для указания контекста, применяемого при запуске рабочих процессов. Обычно пул создаётся с помощью функции
multiprocessing.Pool()или методаPool()объекта контекста. В обоих случаях для context задаётся подходящее значение. ЕслиNone, вызов этой функции приведёт к побочному эффекту: будет задан текущий глобальный метод запуска, если он ещё не задан. См. функциюget_context().Обратите внимание: методы объекта пула должен вызывать только процесс, создавший этот пул.
Предупреждение
Объекты
multiprocessing.poolсодержат внутренние ресурсы, которыми необходимо должным образом управлять (как и любыми другими ресурсами): для этого используйте пул как менеджер контекста либо вручную вызывайтеclose()иterminate(). Если этого не сделать, при финализации процесс может зависнуть.Не следует полагаться на сборщик мусора при уничтожении пула: CPython не гарантирует вызов финализатора пула (подробнее см.
object.__del__()).Изменено в версии 3.2: Добавлен параметр maxtasksperchild.
Изменено в версии 3.4: Добавлен параметр context.
Изменено в версии 3.13: По умолчанию для processes используется
os.process_cpu_count()вместоos.cpu_count().Примечание
Рабочие процессы в
Poolобычно работают в течение всего времени обработки очереди задач пула. Во многих других системах (например, Apache, mod_wsgi и т. д.) распространён подход, позволяющий освободить ресурсы, занятые рабочими процессами: рабочему процессу в пуле разрешают выполнить только определённый объём работы, после чего он завершается, очищается, а на замену ему запускается новый процесс. Аргумент maxtasksperchild классаPoolпредоставляет эту возможность конечному пользователю.-
apply(func[, args[, kwds]]) -
Вызывает func с аргументами args и именованными аргументами kwds. Выполнение блокируется до готовности результата. Поскольку этот вызов блокирует выполнение, для параллельной работы лучше подходит
apply_async(). Кроме того, func выполняется только в одном из рабочих процессов пула.
-
apply_async(func[, args[, kwds[, callback[, error_callback]]]]) -
Вариант метода
apply(), который возвращает объектAsyncResult.Если указан callback, он должен быть вызываемым объектом, принимающим один аргумент. Когда результат будет готов, к нему будет применён callback, если только вызов не завершился ошибкой; в таком случае вместо него будет применён error_callback.
Если указан error_callback, он должен быть вызываемым объектом, принимающим один аргумент. Если целевая функция завершится ошибкой, error_callback будет вызван с экземпляром исключения.
Обратные вызовы должны завершаться немедленно, иначе поток, обрабатывающий результаты, будет заблокирован.
-
map(func, iterable[, chunksize]) -
Параллельный аналог встроенной функции
map()(поддерживается только один аргумент iterable; для нескольких последовательностей см.starmap()). Выполнение блокируется до готовности результата.Этот метод разбивает последовательность на несколько блоков и отправляет их в пул процессов как отдельные задачи. (Приблизительный) размер блоков можно задать, установив для chunksize положительное целое число.
Обратите внимание: для очень длинных последовательностей этот метод может потреблять много памяти. Для повышения эффективности рассмотрите возможность использовать
imap()илиimap_unordered(), явно задав параметр chunksize.
-
map_async(func, iterable[, chunksize[, callback[, error_callback]]]) -
Вариант метода
map(), который возвращает объектAsyncResult.Если указан callback, он должен быть вызываемым объектом, принимающим один аргумент. Когда результат будет готов, к нему будет применён callback, если только вызов не завершился ошибкой; в таком случае вместо него будет применён error_callback.
Если указан error_callback, он должен быть вызываемым объектом, принимающим один аргумент. Если целевая функция завершится ошибкой, error_callback будет вызван с экземпляром исключения.
Обратные вызовы должны завершаться немедленно, иначе поток, обрабатывающий результаты, будет заблокирован.
-
imap(func, iterable[, chunksize]) -
Более ленивый вариант
map().Аргумент chunksize такой же, как в методе
map(). Для очень длинных последовательностей большое значение chunksize может ускорить выполнение задачи значительно по сравнению со значением по умолчанию1.Кроме того, если chunksize равен
1, методnext()итератора, возвращённого методомimap(), принимает необязательный параметр timeout:next(timeout)вызоветmultiprocessing.TimeoutError, если результат не удастся вернуть в течение указанного количества секунд.
-
imap_unordered(func, iterable[, chunksize]) -
То же, что и
imap(), но порядок результатов в возвращаемом итераторе следует считать произвольным. (Порядок гарантированно будет «правильным» только при наличии одного рабочего процесса.)
-
starmap(func, iterable[, chunksize]) -
Подобен
map(), но предполагается, что элементы iterable являются последовательностями, которые распаковываются в аргументы.Таким образом, iterable, содержащий
[(1,2), (3, 4)], приводит к[func(1,2), func(3,4)].Добавлено в версии 3.3.
-
starmap_async(func, iterable[, chunksize[, callback[, error_callback]]]) -
Сочетает
starmap()иmap_async(): перебирает последовательность iterable, состоящую из последовательностей, и вызывает func, распаковывая эти последовательности в аргументы. Возвращает объект результата.Добавлено в версии 3.3.
-
close() -
Запрещает отправлять в пул новые задачи. После завершения всех задач рабочие процессы завершатся.
-
terminate() -
Немедленно останавливает рабочие процессы, не дожидаясь завершения выполняемой работы. При сборке мусора для объекта пула будет немедленно вызван
terminate().
-
join() -
Ожидает завершения рабочих процессов. Перед использованием
join()необходимо вызватьclose()илиterminate().
Изменено в версии 3.3: Объекты пула теперь поддерживают протокол менеджера контекста — см. Типы менеджеров контекста.
__enter__()возвращает объект пула, а__exit__()вызываетterminate(). -
-
class multiprocessing.pool.AsyncResult -
Класс результата, возвращаемого методами
Pool.apply_async()иPool.map_async().-
get([timeout]) -
Возвращает результат по его готовности. Если timeout не равен
Noneи результат не поступил в течение указанного количества секунд, возбуждаетсяmultiprocessing.TimeoutError. Если удалённый вызов вызвал исключение, оно будет повторно возбуждено методомget().
-
wait([timeout]) -
Ожидает появления результата или истечения указанного количества секунд timeout.
-
ready() -
Возвращает признак завершения вызова.
-
successful() -
Возвращает признак того, что вызов завершился без возбуждения исключения. Если результат ещё не готов, возбуждается
ValueError.Изменено в версии 3.7: Если результат ещё не готов, возбуждается
ValueErrorвместоAssertionError.
-
В следующем примере показано использование пула:
from multiprocessing import Pool
import time
def f(x):
return x*x
if __name__ == '__main__':
with Pool(processes=4) as pool: # start 4 worker processes
result = pool.apply_async(f, (10,)) # evaluate "f(10)" asynchronously in a single process
print(result.get(timeout=1)) # prints "100" unless your computer is *very* slow
print(pool.map(f, range(10))) # prints "[0, 1, 4,..., 81]"
it = pool.imap(f, range(10))
print(next(it)) # prints "0"
print(next(it)) # prints "1"
print(it.next(timeout=1)) # prints "4" unless your computer is *very* slow
result = pool.apply_async(time.sleep, (10,))
print(result.get(timeout=1)) # raises multiprocessing.TimeoutError
Слушатели и клиенты
Обычно обмен сообщениями между процессами выполняется с помощью очередей или объектов Connection, возвращаемых функцией Pipe().
Однако модуль multiprocessing.connection предоставляет дополнительные возможности. По сути, он предлагает высокоуровневый API для обмена сообщениями через сокеты или именованные каналы Windows. Он также поддерживает дайджест-аутентификацию с помощью модуля hmac и опрос нескольких соединений одновременно.
-
multiprocessing.connection.deliver_challenge(connection, authkey) -
Отправляет другому концу соединения случайно сгенерированное сообщение и ожидает ответа.
Если ответ совпадает с дайджестом сообщения, вычисленным с использованием authkey в качестве ключа, другому концу соединения отправляется приветственное сообщение. В противном случае возбуждается
AuthenticationError.
-
multiprocessing.connection.answer_challenge(connection, authkey) -
Получает сообщение, вычисляет его дайджест с использованием authkey в качестве ключа и отправляет дайджест обратно.
Если приветственное сообщение не получено, возбуждается
AuthenticationError.
-
multiprocessing.connection.Client(address[, family[, authkey]]) -
Пытается установить соединение со слушателем, использующим адрес address, и возвращает объект
Connection.Тип соединения определяется аргументом family, но обычно его можно не указывать, поскольку тип обычно можно определить по формату address. (См. раздел Форматы адресов.)
Если задано значение authkey, отличное от
None, оно должно быть байтовой строкой и будет использоваться как секретный ключ в процессе аутентификации на основе HMAC. Если authkey равноNone, аутентификация не выполняется. Если аутентификация завершится неудачей, возбуждаетсяAuthenticationError. См. раздел Ключи аутентификации.
-
class multiprocessing.connection.Listener([address[, family[, backlog[, authkey]]]]) -
Обёртка для связанного сокета или именованного канала Windows, ожидающего соединений.
address — адрес связанного сокета или именованного канала объекта-слушателя.
Примечание
Если используется адрес «0.0.0.0», в Windows он не будет доступной для подключения конечной точкой. Если нужна конечная точка, доступная для подключения, используйте «127.0.0.1».
family — тип используемого сокета (или именованного канала). Это может быть одна из строк
'AF_INET'(для TCP-сокета),'AF_UNIX'(для доменного сокета Unix) или'AF_PIPE'(для именованного канала Windows). Из этих вариантов гарантированно доступен только первый. Если family равноNone, тип определяется по формату address. Если address также равноNone, выбирается значение по умолчанию. Это тип, который считается самым быстрым из доступных. См. раздел Форматы адресов. Обратите внимание: если family равно'AF_UNIX', а address равноNone, сокет будет создан в закрытом временном каталоге, созданном с помощьюtempfile.mkstemp().Если объект-слушатель использует сокет, параметр backlog (по умолчанию равный 1) передаётся методу
listen()сокета после его связывания.Если задано значение authkey, отличное от
None, оно должно быть байтовой строкой и будет использоваться как секретный ключ в процессе аутентификации на основе HMAC. Если authkey равноNone, аутентификация не выполняется. Если аутентификация завершится неудачей, возбуждаетсяAuthenticationError. См. раздел Ключи аутентификации.-
accept() -
Принимает соединение через связанный сокет или именованный канал объекта-слушателя и возвращает объект
Connection. Если выполняется аутентификация и она завершается неудачей, возбуждаетсяAuthenticationError.
-
close() -
Закрывает связанный сокет или именованный канал объекта-слушателя. Этот метод вызывается автоматически при сборке мусора для слушателя. Тем не менее рекомендуется вызывать его явно.
Объекты-слушатели имеют следующие атрибуты, доступные только для чтения:
-
address -
Адрес, используемый объектом Listener.
-
last_accepted -
Адрес, с которого поступило последнее принятое соединение. Если он недоступен, значение равно
None.
Изменено в версии 3.3: Объекты Listener теперь поддерживают протокол менеджера контекста — см. Типы менеджеров контекста.
__enter__()возвращает объект-слушатель, а__exit__()вызываетclose(). -
-
multiprocessing.connection.wait(object_list, timeout=None) -
Ожидает готовности одного из объектов в object_list. Возвращает список готовых объектов из object_list. Если timeout — число с плавающей точкой, вызов блокируется не дольше указанного количества секунд. Если timeout равно
None, блокировка будет длиться неограниченно. Отрицательное значение тайм-аута эквивалентно нулевому.И в POSIX, и в Windows объект можно включить в object_list, если это:
- доступный для чтения объект
Connection; - подключённый и доступный для чтения объект
socket.socket; или - атрибут
sentinelобъектаProcess.
Объект соединения или сокета считается готовым, если из него можно прочитать данные или если другой конец соединения был закрыт.
POSIX:
wait(object_list, timeout)почти эквивалентенselect.select(object_list, [], [], timeout). Разница заключается в том, что еслиselect.select()прерван сигналом, он может вызватьOSErrorс номером ошибкиEINTR, тогда какwait()этого не сделает.Windows: элемент object_list должен быть либо целочисленным дескриптором, поддерживающим ожидание (согласно определению в документации функции Win32
WaitForMultipleObjects()), либо объектом с методомfileno(), который возвращает дескриптор сокета или канала. (Обратите внимание: дескрипторы каналов и сокетов не поддерживают ожидание.)Добавлено в версии 3.3.
- доступный для чтения объект
Примеры
Следующий код сервера создаёт слушателя, использующего 'secret password' в качестве ключа аутентификации. Затем он ожидает соединения и отправляет клиенту данные:
from multiprocessing.connection import Listener
from array import array
address = ('localhost', 6000) # family is deduced to be 'AF_INET'
with Listener(address, authkey=b'secret password') as listener:
with listener.accept() as conn:
print('connection accepted from', listener.last_accepted)
conn.send([2.25, None, 'junk', float])
conn.send_bytes(b'hello')
conn.send_bytes(array('i', [42, 1729]))
Следующий код подключается к серверу и получает от него данные:
from multiprocessing.connection import Client
from array import array
address = ('localhost', 6000)
with Client(address, authkey=b'secret password') as conn:
print(conn.recv()) # => [2.25, None, 'junk', float]
print(conn.recv_bytes()) # => 'hello'
arr = array('i', [0, 0, 0, 0, 0])
print(conn.recv_bytes_into(arr)) # => 8
print(arr) # => array('i', [42, 1729, 0, 0, 0])
Следующий код использует wait(), чтобы одновременно ожидать сообщений от нескольких процессов:
from multiprocessing import Process, Pipe, current_process
from multiprocessing.connection import wait
def foo(w):
for i in range(10):
w.send((i, current_process().name))
w.close()
if __name__ == '__main__':
readers = []
for i in range(4):
r, w = Pipe(duplex=False)
readers.append(r)
p = Process(target=foo, args=(w,))
p.start()
# We close the writable end of the pipe now to be sure that
# p is the only process which owns a handle for it. This
# ensures that when p closes its handle for the writable end,
# wait() will promptly report the readable end as being ready.
w.close()
while readers:
for r in wait(readers):
try:
msg = r.recv()
except EOFError:
readers.remove(r)
else:
print(msg)
Форматы адресов
- Адрес
'AF_INET'— это кортеж вида(hostname, port), где hostname — строка, а port — целое число. - Адрес
'AF_UNIX'— это строка, представляющая имя файла в файловой системе. - Адрес
'AF_PIPE'— это строка видаr'\\.\pipe\PipeName'. Чтобы подключиться с помощьюClient()к именованному каналу на удалённом компьютере с именем ServerName, следует использовать адрес видаr'\\ServerName\pipe\PipeName'.
Обратите внимание: по умолчанию любая строка, начинающаяся с двух обратных косых черт, считается адресом 'AF_PIPE', а не адресом 'AF_UNIX'.
Ключи аутентификации
При использовании Connection.recv полученные данные автоматически распаковываются с помощью pickle. К сожалению, распаковка данных из ненадёжного источника представляет угрозу безопасности. Поэтому Listener и Client() используют модуль hmac для дайджест-аутентификации.
Ключ аутентификации — это байтовая строка, которую можно считать паролем: после установления соединения каждая сторона будет требовать от другой доказательства знания ключа аутентификации. (Подтверждение того, что обе стороны используют один и тот же ключ, не предполагает передачу ключа по соединению.)
Если запрошена аутентификация, но ключ аутентификации не указан, используется возвращаемое значение current_process().authkey (см. Process). Это значение автоматически наследуется любым объектом Process, создаваемым текущим процессом. Это означает, что (по умолчанию) все процессы многопроцессной программы используют общий ключ аутентификации, который можно применять при установлении соединений между ними.
Подходящие ключи аутентификации также можно создать с помощью os.urandom().
Эта аутентификация защищает соединения Listener и Client(), доступные по адресу. Она не применяется к анонимным каналам, создаваемым функцией Pipe() или используемым внутри Queue. multiprocessing считает доверенными все локальные процессы, запущенные от имени одного пользователя; в большинстве операционных систем такие процессы в любом случае могут получать доступ к файловым дескрипторам каналов друг друга. Приложениям, которым требуется изоляция процессов одного пользователя, необходимо обеспечить её на уровне операционной системы — например, запускать рабочие процессы от имени другой учётной записи или в песочнице.
Ведение журнала
Доступна некоторая поддержка ведения журнала. Однако обратите внимание, что пакет logging не использует блокировки, совместно используемые процессами, поэтому сообщения разных процессов могут перемешиваться (в зависимости от типа обработчика).
-
multiprocessing.get_logger() -
Возвращает регистратор, используемый
multiprocessing. При необходимости будет создан новый.При первом создании регистратор имеет уровень
logging.NOTSETи не имеет обработчика по умолчанию. Сообщения, отправленные этому регистратору, по умолчанию не передаются корневому регистратору.Обратите внимание, что в Windows дочерние процессы наследуют только уровень регистратора родительского процесса — другие настройки регистратора не наследуются.
-
multiprocessing.log_to_stderr(level=None) -
Эта функция вызывает
get_logger(), но помимо возврата регистратора, созданного get_logger, добавляет обработчик, который отправляет вывод вsys.stderrв формате'[%(levelname)s/%(processName)s] %(message)s'. Вы можете изменитьlevelnameрегистратора, передав аргументlevel.
Ниже приведён пример сеанса с включённым ведением журнала:
>>> import multiprocessing, logging
>>> logger = multiprocessing.log_to_stderr()
>>> logger.setLevel(logging.INFO)
>>> logger.warning('doomed')
[WARNING/MainProcess] doomed
>>> m = multiprocessing.Manager()
[INFO/SyncManager-...] child process calling self.run()
[INFO/SyncManager-...] created temp directory /.../pymp-...
[INFO/SyncManager-...] manager serving at '/.../listener-...'
>>> del m
[INFO/MainProcess] sending shutdown message to manager
[INFO/SyncManager-...] manager exiting with exitcode 0
Полную таблицу уровней ведения журнала см. в модуле logging.
Модуль multiprocessing.dummy
multiprocessing.dummy воспроизводит API multiprocessing, но представляет собой лишь оболочку над модулем threading.
В частности, функция Pool, предоставляемая multiprocessing.dummy, возвращает экземпляр ThreadPool — подкласса Pool, поддерживающего те же вызовы методов, но использующего пул рабочих потоков вместо рабочих процессов.
-
class multiprocessing.pool.ThreadPool([processes[, initializer[, initargs]]]) -
Объект пула потоков, управляющий пулом рабочих потоков, которым можно отправлять задания. Экземпляры
ThreadPoolполностью совместимы по интерфейсу с экземплярамиPool, и их ресурсы также необходимо правильно освобождать: либо используя пул как менеджер контекста, либо вручную вызываяclose()иterminate().processes — количество рабочих потоков. Если processes равно
None, используется значение, возвращаемое функциейos.process_cpu_count().Если initializer не равно
None, при запуске каждый рабочий процесс вызоветinitializer(*initargs).В отличие от
Pool, параметры maxtasksperchild и context указывать нельзя.Примечание
Объект
ThreadPoolимеет тот же интерфейс, что иPool, который изначально рассчитан на пул процессов и появился раньше модуляconcurrent.futures. Поэтому он наследует некоторые операции, не имеющие смысла для пула на основе потоков, а также использует собственный тип для представления состояния асинхронных заданий —AsyncResult, который не распознаётся другими библиотеками.Как правило, пользователям следует предпочесть
concurrent.futures.ThreadPoolExecutor: его более простой интерфейс изначально разработан для потоков, а возвращаемые им экземплярыconcurrent.futures.Futureсовместимы со многими другими библиотеками, в том числе сasyncio.
Рекомендации по программированию
При использовании multiprocessing следует соблюдать определённые рекомендации и соглашения.
Все методы запуска
Следующие рекомендации относятся ко всем методам запуска.
Избегайте совместного состояния
По возможности старайтесь избегать передачи больших объёмов данных между процессами.
Для межпроцессного взаимодействия, вероятно, лучше использовать очереди или каналы, а не низкоуровневые примитивы синхронизации.
Возможность сериализации с помощью pickle
Убедитесь, что аргументы методов прокси можно сериализовать с помощью pickle.
Потокобезопасность прокси
Не используйте объект прокси из нескольких потоков, если только вы не защищаете его блокировкой.
(Никаких проблем не возникает, если разные процессы используют один и тот же прокси.)
Присоединение процессов-зомби
В POSIX завершившийся, но не присоединённый процесс становится зомби. Их не должно накапливаться много, поскольку при каждом запуске нового процесса (или вызове active_children()) присоединяются все завершившиеся процессы, которые ещё не были присоединены. Кроме того, вызов Process.is_alive для завершившегося процесса присоединит его. Тем не менее, хорошей практикой будет явно присоединять все запущенные вами процессы.
Предпочитайте наследование сериализации и десериализации с помощью pickle
При использовании методов запуска spawn или forkserver многие типы из multiprocessing должны поддерживать сериализацию с помощью pickle, чтобы дочерние процессы могли их использовать. Однако обычно следует избегать передачи совместно используемых объектов другим процессам через каналы или очереди. Вместо этого организуйте программу так, чтобы процесс, которому нужен доступ к совместно используемому ресурсу, созданному в другом месте, мог унаследовать его от процесса-предка.
Избегайте принудительного завершения процессов
Использование метода Process.terminate для остановки процесса может привести к повреждению или недоступности для других процессов любых совместно используемых ресурсов (например, блокировок, семафоров, каналов и очередей), которые использовались этим процессом.
Поэтому, вероятно, лучше использовать Process.terminate только для процессов, которые никогда не используют совместно используемые ресурсы.
Присоединение процессов, использующих очереди
Помните, что процесс, поместивший элементы в очередь, будет ожидать перед завершением, пока поток «подачи» не передаст все буферизованные элементы в базовый канал. (Чтобы избежать такого поведения, дочерний процесс может вызвать метод очереди Queue.cancel_join_thread.)
Это означает, что при использовании очереди необходимо убедиться, что все помещённые в неё элементы будут извлечены до присоединения процесса. В противном случае нельзя быть уверенным, что процессы, поместившие элементы в очередь, завершатся. Также помните, что процессы, не являющиеся демонами, будут присоединены автоматически.
Следующий пример приведёт к взаимной блокировке:
from multiprocessing import Process, Queue
def f(q):
q.put('X' * 1000000)
if __name__ == '__main__':
queue = Queue()
p = Process(target=f, args=(queue,))
p.start()
p.join() # this deadlocks
obj = queue.get()
Исправить это можно, поменяв местами две последние строки (или просто удалив строку p.join()).
Явно передавайте ресурсы дочерним процессам
В POSIX при использовании метода запуска fork дочерний процесс может получить доступ к совместно используемому ресурсу, созданному в родительском процессе, используя глобальный ресурс. Однако лучше передавать объект как аргумент конструктору дочернего процесса.
Помимо того, что это (потенциально) обеспечивает совместимость кода с Windows и другими методами запуска, это также гарантирует, что объект не будет удалён сборщиком мусора в родительском процессе, пока дочерний процесс остаётся активным. Это может быть важно, если при удалении объекта сборщиком мусора в родительском процессе освобождается какой-либо ресурс.
Например,
from multiprocessing import Process, Lock
def f():
... do something using "lock" ...
if __name__ == '__main__':
lock = Lock()
for i in range(10):
Process(target=f).start()
следует переписать так:
from multiprocessing import Process, Lock
def f(l):
... do something using "l" ...
if __name__ == '__main__':
lock = Lock()
for i in range(10):
Process(target=f, args=(lock,)).start()
Будьте осторожны, заменяя sys.stdin «файлоподобным объектом»
Изначально multiprocessing безусловно вызывал:
os.close(sys.stdin.fileno())
в методе multiprocessing.Process._bootstrap() — это приводило к проблемам при вложенных процессах. Теперь используется следующий вызов:
sys.stdin.close() sys.stdin = open(os.open(os.devnull, os.O_RDONLY), closefd=False)
Это решает основную проблему взаимного влияния процессов, приводившую к ошибке неверного файлового дескриптора, но создаёт потенциальную опасность для приложений, заменяющих sys.stdin() «файлоподобным объектом» с буферизацией вывода. Опасность заключается в том, что если несколько процессов вызовут close() для этого файлоподобного объекта, одни и те же данные могут быть сброшены в объект несколько раз, что приведёт к повреждению данных.
Если вы создаёте файлоподобный объект и реализуете собственное кэширование, его можно сделать безопасным при fork: сохраняйте PID при каждом добавлении данных в кэш и очищайте кэш при изменении PID. Например:
@property
def cache(self):
pid = os.getpid()
if pid != self._pid:
self._pid = pid
self._cache = []
return self._cache
Дополнительную информацию см. в bpo-5155, bpo-5313 и bpo-5331
Методы запуска spawn и forkserver
Существуют дополнительные ограничения, не применимые к методу запуска fork.
Дополнительные требования к сериализации с помощью pickle
Убедитесь, что все аргументы Process можно сериализовать с помощью pickle. Кроме того, если вы создаёте подкласс Process.__init__, необходимо убедиться, что его экземпляры можно сериализовать с помощью pickle при вызове метода Process.start.
Глобальные переменные
Помните, что если код, выполняемый в дочернем процессе, попытается обратиться к глобальной переменной, её значение (если оно есть) может отличаться от значения в родительском процессе на момент вызова Process.start.
Однако глобальные переменные, представляющие собой константы на уровне модуля, не вызывают проблем.
Безопасный импорт главного модуля
Убедитесь, что главный модуль можно безопасно импортировать новым интерпретатором Python без непреднамеренных побочных эффектов (например, запуска нового процесса).
Например, при использовании метода запуска spawn или forkserver выполнение следующего модуля завершится ошибкой RuntimeError:
from multiprocessing import Process
def foo():
print('hello')
p = Process(target=foo)
p.start()
Вместо этого следует защитить «точку входа» программы с помощью if
__name__ == '__main__': следующим образом:
from multiprocessing import Process, freeze_support, set_start_method
def foo():
print('hello')
if __name__ == '__main__':
freeze_support()
set_start_method('spawn')
p = Process(target=foo)
p.start()
(Строку freeze_support() можно опустить, если программа будет запускаться обычным способом, а не в замороженном виде.)
Это позволяет вновь созданному интерпретатору Python безопасно импортировать модуль, а затем выполнить функцию foo() этого модуля.
Аналогичные ограничения действуют, если пул или менеджер создаётся в главном модуле.
Примеры
Пример создания и использования пользовательских менеджеров и прокси:
from multiprocessing import freeze_support
from multiprocessing.managers import BaseManager, BaseProxy
import operator
##
class Foo:
def f(self):
print('you called Foo.f()')
def g(self):
print('you called Foo.g()')
def _h(self):
print('you called Foo._h()')
# A simple generator function
def baz():
for i in range(10):
yield i*i
# Proxy type for generator objects
class GeneratorProxy(BaseProxy):
_exposed_ = ['__next__']
def __iter__(self):
return self
def __next__(self):
return self._callmethod('__next__')
# Function to return the operator module
def get_operator_module():
return operator
##
class MyManager(BaseManager):
pass
# register the Foo class; make `f()` and `g()` accessible via proxy
MyManager.register('Foo1', Foo)
# register the Foo class; make `g()` and `_h()` accessible via proxy
MyManager.register('Foo2', Foo, exposed=('g', '_h'))
# register the generator function baz; use `GeneratorProxy` to make proxies
MyManager.register('baz', baz, proxytype=GeneratorProxy)
# register get_operator_module(); make public functions accessible via proxy
MyManager.register('operator', get_operator_module)
##
def test():
manager = MyManager()
manager.start()
print('-' * 20)
f1 = manager.Foo1()
f1.f()
f1.g()
assert not hasattr(f1, '_h')
assert sorted(f1._exposed_) == sorted(['f', 'g'])
print('-' * 20)
f2 = manager.Foo2()
f2.g()
f2._h()
assert not hasattr(f2, 'f')
assert sorted(f2._exposed_) == sorted(['g', '_h'])
print('-' * 20)
it = manager.baz()
for i in it:
print('<%d>' % i, end=' ')
print()
print('-' * 20)
op = manager.operator()
print('op.add(23, 45) =', op.add(23, 45))
print('op.pow(2, 94) =', op.pow(2, 94))
print('op._exposed_ =', op._exposed_)
##
if __name__ == '__main__':
freeze_support()
test()
Использование Pool:
import multiprocessing
import time
import random
import sys
#
# Functions used by test code
#
def calculate(func, args):
result = func(*args)
return '%s says that %s%s = %s' % (
multiprocessing.current_process().name,
func.__name__, args, result
)
def calculatestar(args):
return calculate(*args)
def mul(a, b):
time.sleep(0.5 * random.random())
return a * b
def plus(a, b):
time.sleep(0.5 * random.random())
return a + b
def f(x):
return 1.0 / (x - 5.0)
def pow3(x):
return x ** 3
def noop(x):
pass
#
# Test code
#
def test():
PROCESSES = 4
print('Creating pool with %d processes\n' % PROCESSES)
with multiprocessing.Pool(PROCESSES) as pool:
#
# Tests
#
TASKS = [(mul, (i, 7)) for i in range(10)] + \
[(plus, (i, 8)) for i in range(10)]
results = [pool.apply_async(calculate, t) for t in TASKS]
imap_it = pool.imap(calculatestar, TASKS)
imap_unordered_it = pool.imap_unordered(calculatestar, TASKS)
print('Ordered results using pool.apply_async():')
for r in results:
print('\t', r.get())
print()
print('Ordered results using pool.imap():')
for x in imap_it:
print('\t', x)
print()
print('Unordered results using pool.imap_unordered():')
for x in imap_unordered_it:
print('\t', x)
print()
print('Ordered results using pool.map() --- will block till complete:')
for x in pool.map(calculatestar, TASKS):
print('\t', x)
print()
#
# Test error handling
#
print('Testing error handling:')
try:
print(pool.apply(f, (5,)))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from pool.apply()')
else:
raise AssertionError('expected ZeroDivisionError')
try:
print(pool.map(f, list(range(10))))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from pool.map()')
else:
raise AssertionError('expected ZeroDivisionError')
try:
print(list(pool.imap(f, list(range(10)))))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from list(pool.imap())')
else:
raise AssertionError('expected ZeroDivisionError')
it = pool.imap(f, list(range(10)))
for i in range(10):
try:
x = next(it)
except ZeroDivisionError:
if i == 5:
pass
except StopIteration:
break
else:
if i == 5:
raise AssertionError('expected ZeroDivisionError')
assert i == 9
print('\tGot ZeroDivisionError as expected from IMapIterator.next()')
print()
#
# Testing timeouts
#
print('Testing ApplyResult.get() with timeout:', end=' ')
res = pool.apply_async(calculate, TASKS[0])
while 1:
sys.stdout.flush()
try:
sys.stdout.write('\n\t%s' % res.get(0.02))
break
except multiprocessing.TimeoutError:
sys.stdout.write('.')
print()
print()
print('Testing IMapIterator.next() with timeout:', end=' ')
it = pool.imap(calculatestar, TASKS)
while 1:
sys.stdout.flush()
try:
sys.stdout.write('\n\t%s' % it.next(0.02))
except StopIteration:
break
except multiprocessing.TimeoutError:
sys.stdout.write('.')
print()
print()
if __name__ == '__main__':
multiprocessing.freeze_support()
test()
Пример использования очередей для передачи заданий набору рабочих процессов и сбора результатов:
import time
import random
from multiprocessing import Process, Queue, current_process, freeze_support
#
# Function run by worker processes
#
def worker(input, output):
for func, args in iter(input.get, 'STOP'):
result = calculate(func, args)
output.put(result)
#
# Function used to calculate result
#
def calculate(func, args):
result = func(*args)
return '%s says that %s%s = %s' % \
(current_process().name, func.__name__, args, result)
#
# Functions referenced by tasks
#
def mul(a, b):
time.sleep(0.5*random.random())
return a * b
def plus(a, b):
time.sleep(0.5*random.random())
return a + b
#
#
#
def test():
NUMBER_OF_PROCESSES = 4
TASKS1 = [(mul, (i, 7)) for i in range(20)]
TASKS2 = [(plus, (i, 8)) for i in range(10)]
# Create queues
task_queue = Queue()
done_queue = Queue()
# Submit tasks
for task in TASKS1:
task_queue.put(task)
# Start worker processes
for i in range(NUMBER_OF_PROCESSES):
Process(target=worker, args=(task_queue, done_queue)).start()
# Get and print results
print('Unordered results:')
for i in range(len(TASKS1)):
print('\t', done_queue.get())
# Add more tasks using `put()`
for task in TASKS2:
task_queue.put(task)
# Get and print some more results
for i in range(len(TASKS2)):
print('\t', done_queue.get())
# Tell child processes to stop
for i in range(NUMBER_OF_PROCESSES):
task_queue.put('STOP')
if __name__ == '__main__':
freeze_support()
test()
© 2001 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.14/library/multiprocessing.html