Spec-Zone.ru › Python 3.14

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()

Без блокировки вывод разных процессов может смешиваться.

Совместное использование состояния между процессами

Как упоминалось выше, при параллельном программировании обычно лучше по возможности избегать совместно используемого состояния. Это особенно важно при использовании нескольких процессов.

Однако если вам действительно необходимо использовать общие данные, multiprocessing предоставляет несколько способов сделать это.

Разделяемая память

Данные можно хранить в карте разделяемой памяти с помощью Value или Array. Например, следующий код

from multiprocessing import Process, Value, Array

def f(n, a):
    n.value = 3.1415927
    for i in range(len(a)):
        a[i] = -a[i]

if __name__ == '__main__':
    num = Value('d', 0.0)
    arr = Array('i', range(10))

    p = Process(target=f, args=(num, arr))
    p.start()
    p.join()

    print(num.value)
    print(arr[:])

выведет

3.1415927
[0, -1, -2, -3, -4, -5, -6, -7, -8, -9]

Аргументы 'd' и 'i', используемые при создании num и arr, представляют собой коды типов, используемые модулем array: 'd' обозначает число с плавающей точкой двойной точности, а 'i' — целое число со знаком. Эти общие объекты безопасны для процессов и потоков.

Для более гибкой работы с разделяемой памятью можно использовать модуль multiprocessing.sharedctypes, который поддерживает создание произвольных объектов ctypes, размещённых в разделяемой памяти.

Серверный процесс

Объект менеджера, возвращаемый функцией Manager(), управляет серверным процессом, который хранит объекты Python и позволяет другим процессам работать с ними через прокси.

Менеджер, возвращаемый функцией Manager(), поддерживает типы list, dict, set, Namespace, Lock, RLock, Semaphore, BoundedSemaphore, Condition, Event, Barrier, Queue, Value и Array. Например:

from multiprocessing import Process, Manager

def f(d, l, s):
    d[1] = '1'
    d['2'] = 2
    d[0.25] = None
    l.reverse()
    s.add('a')
    s.add('b')

if __name__ == '__main__':
    with Manager() as manager:
        d = manager.dict()
        l = manager.list(range(10))
        s = manager.set()

        p = Process(target=f, args=(d, l, s))
        p.start()
        p.join()

        print(d)
        print(l)
        print(s)

выведет

{0.25: None, 1: '1', '2': 2}
[9, 8, 7, 6, 5, 4, 3, 2, 1, 0]
{'a', 'b'}

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

Использование пула рабочих процессов

Класс 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", однако начиная с Python 3.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, а фоновый поток позднее передаёт сериализованные данные в нижележащий канал. Это имеет некоторые, немного неожиданные последствия, однако они не должны вызывать практических затруднений. Если же они действительно вас беспокоят, можно вместо этого использовать очередь, созданную с помощью менеджера.

  1. После помещения объекта в пустую очередь может пройти ничтожно малое время, прежде чем метод очереди empty() вернёт False, а вызов get_nowait() сможет завершиться без исключения queue.Empty.
  2. Если несколько процессов помещают объекты в очередь, то на другой стороне они могут быть получены не в том порядке. Однако объекты, добавленные одним и тем же процессом, всегда будут расположены в ожидаемом порядке относительно друг друга.

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

Если процесс завершён с помощью 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.

См. также

os.cpu_count() os.process_cpu_count()

Изменено в версии 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, после чего соединение больше нельзя будет читать.

Изменено в версии 3.3: Раньше эта функция вызывала исключение IOError, которое теперь является псевдонимом 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.

Объекты ctypes в общей памяти

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

multiprocessing.Value(typecode_or_type, *args, lock=True)

Возвращает объект ctypes, выделенный в общей памяти. По умолчанию возвращаемое значение на самом деле представляет собой синхронизированную оболочку объекта. Доступ к самому объекту можно получить через атрибут value объекта Value.

typecode_or_type определяет тип возвращаемого объекта: это либо тип ctypes, либо однобуквенный код типа, используемый модулем array. *args передаются конструктору типа.

Если lock имеет значение True (по умолчанию), для синхронизации доступа к значению создаётся новый рекурсивный объект блокировки. Если lock — объект Lock или RLock, он будет использоваться для синхронизации доступа к значению. Если lock имеет значение False, доступ к возвращаемому объекту не будет автоматически защищён блокировкой, поэтому он не обязательно будет «безопасен для процессов».

Такие операции, как +=, включающие чтение и запись, не являются атомарными. Поэтому, если, например, требуется атомарно увеличить общее значение, недостаточно просто выполнить

counter.value += 1

Если связанная блокировка рекурсивная (а по умолчанию это так), вместо этого можно выполнить

with counter.get_lock():
    counter.value += 1

Обратите внимание, что lock — это аргумент, который можно передать только по имени.

multiprocessing.Array(typecode_or_type, size_or_initializer, *, lock=True)

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

typecode_or_type определяет тип элементов возвращаемого массива: это либо тип ctypes, либо однобуквенный код типа, используемый модулем array, за исключением 'w', который не поддерживается. Кроме того, код типа 'c' является псевдонимом для ctypes.c_char. Если size_or_initializer — целое число, оно определяет длину массива, который изначально заполняется нулями. В противном случае size_or_initializer — последовательность, используемая для инициализации массива; её длина определяет длину массива.

Если lock имеет значение True (по умолчанию), для синхронизации доступа к значению создаётся новый объект блокировки. Если lock — объект Lock или RLock, он будет использоваться для синхронизации доступа к значению. Если lock имеет значение False, доступ к возвращаемому объекту не будет автоматически защищён блокировкой, поэтому он не обязательно будет «безопасен для процессов».

Обратите внимание, что lock — это аргумент, который можно передать только по имени.

Обратите внимание, что массив ctypes.c_char имеет атрибуты value и raw, которые можно использовать для хранения и получения байтовых строк. Атрибут raw позволяет взаимодействовать с объектом bytes размером со весь массив, тогда как чтение value завершается после нулевого байта, как это обычно происходит при обработке строк в большинстве языков программирования.

Модуль multiprocessing.sharedctypes

Модуль multiprocessing.sharedctypes предоставляет функции для выделения объектов ctypes в общей памяти, чтобы дочерние процессы могли их наследовать.

Примечание

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

multiprocessing.sharedctypes.RawArray(typecode_or_type, size_or_initializer)

Возвращает массив ctypes, выделенный в общей памяти.

typecode_or_type определяет тип элементов возвращаемого массива: это либо тип ctypes, либо однобуквенный код типа, используемый модулем array. Если size_or_initializer — целое число, оно определяет длину массива, который изначально заполняется нулями. В противном случае size_or_initializer — последовательность, используемая для инициализации массива; её длина определяет длину массива.

Обратите внимание, что запись и чтение элемента потенциально не являются атомарными операциями; используйте вместо этого Array(), чтобы обеспечить автоматическую синхронизацию доступа с помощью блокировки.

multiprocessing.sharedctypes.RawValue(typecode_or_type, *args)

Возвращает объект ctypes, выделенный в общей памяти.

typecode_or_type определяет тип возвращаемого объекта: это либо тип ctypes, либо однобуквенный код типа, используемый модулем array. *args передаются конструктору типа.

Обратите внимание, что запись и чтение значения потенциально не являются атомарными операциями; используйте вместо этого Value(), чтобы обеспечить автоматическую синхронизацию доступа с помощью блокировки.

Обратите внимание, что массив ctypes.c_char имеет атрибуты value и raw, позволяющие использовать его для хранения и получения строк; см. документацию по ctypes.

multiprocessing.sharedctypes.Array(typecode_or_type, size_or_initializer, *, lock=True, ctx=None)

То же, что и RawArray(), за исключением того, что в зависимости от значения lock вместо необёрнутого массива ctypes может возвращаться синхронизированная оболочка, безопасная для процессов.

Если lock имеет значение True (по умолчанию), для синхронизации доступа к значению создаётся новый объект блокировки. Если lock — объект Lock или RLock, он будет использоваться для синхронизации доступа к значению. Если lock имеет значение False, доступ к возвращаемому объекту не будет автоматически защищён блокировкой, поэтому он не обязательно будет «безопасен для процессов».

ctx — это объект контекста либо None (используется текущий контекст). Если None, вызов этой функции может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.

Обратите внимание, что lock и ctx — это параметры, которые можно передать только по имени.

multiprocessing.sharedctypes.Value(typecode_or_type, *args, lock=True, ctx=None)

То же, что и RawValue(), за исключением того, что в зависимости от значения lock вместо необёрнутого объекта ctypes может возвращаться синхронизированная оболочка, безопасная для процессов.

Если lock имеет значение True (по умолчанию), для синхронизации доступа к значению создаётся новый объект блокировки. Если lock — объект Lock или RLock, он будет использоваться для синхронизации доступа к значению. Если lock имеет значение False, доступ к возвращаемому объекту не будет автоматически защищён блокировкой, поэтому он не обязательно будет «безопасен для процессов».

ctx — это объект контекста либо None (используется текущий контекст). Если None, вызов этой функции может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.

Обратите внимание, что lock и ctx — это параметры, которые можно передать только по имени.

multiprocessing.sharedctypes.copy(obj)

Возвращает объект ctypes, выделенный в общей памяти и являющийся копией объекта ctypes obj.

multiprocessing.sharedctypes.synchronized(obj, lock=None, ctx=None)

Возвращает оболочку объекта ctypes, безопасную для процессов и использующую lock для синхронизации доступа. Если lock имеет значение None (по умолчанию), объект multiprocessing.RLock создаётся автоматически.

ctx — это объект контекста либо None (используется текущий контекст). Если None, вызов этой функции может установить глобальный метод запуска. Подробнее см. в разделе Глобальный метод запуска.

У синхронизированной оболочки будут два метода в дополнение к методам обёрнутого объекта: get_obj() возвращает обёрнутый объект, а get_lock() возвращает объект блокировки, используемый для синхронизации.

Обратите внимание, что доступ к объекту ctypes через оболочку может быть намного медленнее, чем доступ к необёрнутому объекту ctypes.

Изменено в версии 3.5: Синхронизированные объекты поддерживают протокол менеджера контекста.

В таблице ниже сравнивается синтаксис создания объектов ctypes в общей памяти с обычным синтаксисом ctypes. (В таблице MyStruct — некоторый подкласс ctypes.Structure.)

ctypes

sharedctypes с использованием типа

sharedctypes с использованием кода типа

c_double(2.4)

RawValue(c_double, 2.4)

RawValue(‘d’, 2.4)

MyStruct(4, 6)

RawValue(MyStruct, 4, 6)

(c_short * 7)()

RawArray(c_short, 7)

RawArray(‘h’, 7)

(c_int * 3)(9, 2, 8)

RawArray(c_int, (9, 2, 8))

RawArray(‘i’, (9, 2, 8))

Ниже приведён пример изменения нескольких объектов ctypes дочерним процессом:

from multiprocessing import Process, Lock
from multiprocessing.sharedctypes import Value, Array
from ctypes import Structure, c_double

class Point(Structure):
    _fields_ = [('x', c_double), ('y', c_double)]

def modify(n, x, s, A):
    n.value **= 2
    x.value **= 2
    s.value = s.value.upper()
    for a in A:
        a.x **= 2
        a.y **= 2

if __name__ == '__main__':
    lock = Lock()

    n = Value('i', 7)
    x = Value(c_double, 1.0/3.0, lock=False)
    s = Array('c', b'hello world', lock=lock)
    A = Array(Point, [(1.875,-6.25), (-5.75,2.0), (2.375,9.5)], lock=lock)

    p = Process(target=modify, args=(n, x, s, A))
    p.start()
    p.join()

    print(n.value)
    print(x.value)
    print(s.value)
    print([(a.x, a.y) for a in A])

Вывод будет следующим:

49
0.1111111111111111
HELLO WORLD
[(3.515625, 39.0625), (33.0625, 4.0), (5.640625, 90.25)]

Менеджеры

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

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

Spec-Zone.ru

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