Spec-Zone.ru › Python 3.12

multiprocessing — Параллелизм на основе процессов

Исходный код: Lib/multiprocessing/

Доступность: не Emscripten, не WASI.

Этот модуль не работает или недоступен на платформах WebAssembly wasm32-emscripten и wasm32-wasi. См. Платформы WebAssembly для получения дополнительной информации.

Введение

multiprocessing — это пакет, который поддерживает запуск процессов с помощью API, аналогичного модулю threading. Пакет multiprocessing обеспечивает как локальную, так и удалённую конкуретность, эффективно обходя глобальную блокировку интерпретатора, используя подпроцессы вместо потоков. Благодаря этому, модуль multiprocessing позволяет программисту в полной мере использовать несколько процессоров на данном компьютере. Он работает как на платформах POSIX, так и на Windows.

Модуль multiprocessing также вводит API, у которых нет аналогов в модуле threading. Ярким примером этого является объект 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]

См. также

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__', см. Руководство по программированию.

Контексты и методы запуска

В зависимости от платформы, multiprocessing поддерживает три способа запуска процесса. Эти методы запуска —

spawn

Родительский процесс запускает новый процесс Python-интерпретатора. Дочерний процесс унаследует только те ресурсы, необходимые для выполнения метода run() объекта процесса. В частности, ненужные дескрипторы файлов и дескрипторы из родительского процесса не будут унаследованы. Запуск процесса с помощью этого метода относительно медленный по сравнению с использованием fork или forkserver.

Доступно на платформах POSIX и Windows. По умолчанию на Windows и macOS.

fork

Родительский процесс использует os.fork() для разветвления Python-интерпретатора. Дочерний процесс, когда он запускается, фактически идентичен родительскому процессу. Все ресурсы родительского процесса наследуются дочерним процессом. Обратите внимание, что безопасное разветвление многопоточного процесса проблематично.

Доступно на системах POSIX. В настоящее время по умолчанию на системах POSIX за исключением macOS.

Примечание

Метод запуска по умолчанию будет изменён с fork в Python 3.14. Код, которому требуется fork, должен явно указать это с помощью get_context() или set_start_method().

Изменено в версии 3.12: Если Python обнаруживает, что ваш процесс имеет несколько потоков, функция os.fork(), которая используется этим методом запуска, будет поднимать DeprecationWarning. Используйте другой метод запуска. Обратитесь к документации os.fork() для получения дополнительной информации.

forkserver

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

Доступно на платформах POSIX, которые поддерживают передачу дескрипторов файлов через Unix-каналы, например, Linux.

Изменено в версии 3.4: spawn добавлен на все платформы POSIX, а forkserver добавлен для некоторых платформ POSIX. Дочерние процессы больше не наследуют все наследуемые дескрипторы родительского процесса на Windows.

Изменено в версии 3.8: В macOS метод запуска spawn теперь является по умолчанию. Метод запуска fork следует считать небезопасным, так как он может привести к сбою подпроцесса, так как системные библиотеки macOS могут запускать потоки. См. bpo-33725.

На POSIX при использовании методов запуска spawn или forkserver также запускается процесс слежения за ресурсами, который отслеживает неразъединённые именованные системные ресурсы (такие как именованные семафоры или 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.

Библиотека, которая хочет использовать определённый метод запуска, должна использовать get_context(), чтобы избежать вмешательства в выбор пользователя библиотеки.

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

Методы запуска '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, Namespace, Lock, RLock, Semaphore, BoundedSemaphore, Condition, Event, Barrier, Queue, Value и Array. Например,

from multiprocessing import Process, Manager

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

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

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

        print(d)
        print(l)

выведет

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

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

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

Класс 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.

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. Если не указан (по умолчанию), этот флаг будет унаследован от родительского процесса.

По умолчанию целевому объекту не передаются никакие аргументы. Аргумент args, по умолчанию равный (), может быть использован для указания списка или кортежа аргументов для передачи целевому объекту.

Если подкласс переопределяет конструктор, он должен убедиться, что вызывает конструктор базового класса (Process.__init__()) перед выполнением каких-либо других действий с процессом.

Изменено в версии 3.3: Добавлен параметр daemon.

run()

Метод, представляющий активность процесса.

Вы можете переопределить этот метод в подклассе. Стандартный метод run() вызывает вызываемый объект, переданный в конструктор объекта в качестве аргумента target, при наличии, с последовательными и ключевыми аргументами, взятыми из аргументов 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 — положительное число, то он блокируется не более 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.

terminate()

Завершить процесс. В POSIX это делается с помощью сигнала SIGTERM; в Windows используется TerminateProcess(). Обратите внимание, что обработчики выхода и блоки finally и т. д. не будут выполнены.

Обратите внимание, что дочерние процессы процесса не будут завершены — они просто станут сиротами.

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

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

kill()

То же самое, что и terminate(), но с использованием сигнала SIGKILL в POSIX.

Добавлен в версии 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.

Примечание

Когда объект помещается в очередь, он сериализуется, а фоновый поток позже записывает сериализованные данные в подлежащую трубу. Это имеет некоторые последствия, которые могут быть немного неожиданными, но не должны создавать практических трудностей — если они вас действительно беспокоят, вы можете вместо этого использовать очередь, созданную с помощью менеджера.

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

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

Если процесс убивается с помощью Process.terminate() или os.kill(), когда он пытается использовать Queue, данные в очереди могут быть повреждены. Это может вызвать исключение у других процессов при попытке использовать очередь позднее.

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

Как упоминалось выше, если дочерний процесс поместил элементы в очередь (и не использовал JoinableQueue.cancel_join_thread), этот процесс не завершится, пока все буферизованные элементы не будут отправлены в трубу.

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

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

Пример использования очередей для межпроцессного обмена данными см. в разделе Примеры.

multiprocessing.Pipe([duplex])

Возвращает пару (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().

qsize()

Возвращает приблизительный размер очереди. Из-за семантики многопоточности/многопроцессорности это число ненадежно.

Обратите внимание, что это может вызвать NotImplementedError на платформах, таких как macOS, где sem_getvalue() не реализован.

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: Если очередь закрыта, то поднимается ValueError вместо AssertionError.

put_nowait(obj)

Эквивалентно put(obj, False).

get([block[, timeout]])

Удалить и вернуть элемент из очереди. Если необязательные аргументы block равны True (по умолчанию) и timeout равен None (по умолчанию), то заблокировать, если необходимо, до тех пор, пока элемент не появится. Если timeout — положительное число, то блокировать не более timeout секунд и поднять исключение queue.Empty, если элемент не появился в течение этого времени. В противном случае (block равен False), вернуть элемент, если он доступен сразу, иначе поднять исключение queue.Empty (timeout игнорируется в этом случае).

Изменено в версии 3.8: Если очередь закрыта, то поднимается ValueError вместо OSError.

get_nowait()

Эквивалентно get(False).

multiprocessing.Queue имеет несколько дополнительных методов, отсутствующих в queue.Queue. Эти методы обычно не нужны для большинства кодов:

close()

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

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, в противном случае.

Всегда вызывает OSError, если SimpleQueue закрыта.

get()

Удалить и вернуть элемент из очереди.

put(item)

Поместить item в очередь.

class multiprocessing.JoinableQueue([maxsize])

JoinableQueue, подкласс Queue, — это очередь, которая дополнительно имеет методы task_done() и join().

task_done()

Указывает, что ранее помещенная в очередь задача завершена. Используется потребителями очереди. За каждый вызов get() для извлечения задачи, последующий вызов task_done() сообщает очереди, что обработка задачи завершена.

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

Вызывает исключение ValueError, если его вызывают больше раз, чем было помещено элементов в очередь.

join()

Блокируется, пока все элементы в очереди не будут извлечены и обработаны.

Счетчик незавершенных задач увеличивается при добавлении элемента в очередь. Он уменьшается, когда потребитель вызывает task_done(), чтобы указать, что элемент был извлечен и вся работа с ним завершена. Когда счетчик незавершенных задач достигает нуля, join() разблокируется.

Разное

multiprocessing.active_children()

Возвращает список всех активных дочерних процессов текущего процесса.

Вызов этой функции также «присоединяет» любые процессы, которые уже завершились.

multiprocessing.cpu_count()

Возвращает количество ЦП в системе.

Это число не эквивалентно количеству ЦП, которые может использовать текущий процесс. Количество доступных ЦП можно получить с помощью len(os.sched_getaffinity(0))

Если количество ЦП определить невозможно, генерируется исключение NotImplementedError.

См. также

os.cpu_count()

multiprocessing.current_process()

Возвращает объект Process, соответствующий текущему процессу.

Аналог threading.current_thread().

multiprocessing.parent_process()

Возвращает объект Process, соответствующий родительскому процессу current_process(). Для основного процесса parent_process будет None.

Добавлена в версии 3.8.

multiprocessing.freeze_support()

Добавляет поддержку для программ, использующих multiprocessing, которые были заморожены для создания исполняемого файла для Windows. (Тестировалось с 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() не имеет эффекта при вызове на любой операционной системе, кроме Windows. Кроме того, если модуль запускается обычным образом интерпретатором Python на Windows (программа не заморожена), то 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, то метод запуска устанавливается по умолчанию, и возвращается его имя. Если метод запуска не задан и allow_none равно true, то возвращается None

Возвращаемое значение может быть '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'. Поднимает исключение RuntimeError, если метод запуска уже установлен и force не True. Если 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.

Объекты соединения

Объекты соединения позволяют отправлять и получать сериализуемые объекты или строки. Их можно рассматривать как сокеты с ориентированным на сообщения соединением.

Объекты соединения обычно создаются с помощью Pipe — см. также Слушатели и Клиенты.

class multiprocessing.connection.Connection
send(obj)

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

Объект должен быть сериализуем. Очень большие сериализованные объекты (приблизительно 32 МБ+, хотя это зависит от ОС) могут вызвать исключение ValueError.

recv()

Возвращает объект, отправленный с другого конца соединения с помощью send(). Блокируется до тех пор, пока не будет чего-то для получения. Вызывает EOFError, если ничего не осталось для получения, и другой конец был закрыт.

fileno()

Возвращает дескриптор файла или дескриптор, используемый соединением.

close()

Закрыть соединение.

Это вызывается автоматически при сборке мусора соединения.

poll([timeout])

Возвращает, доступны ли какие-либо данные для чтения.

Если timeout не указан, то он возвращается немедленно. Если timeout является числом, то это указывает максимальное время в секундах для блокировки. Если timeout равен None , то используется бесконечный таймаут.

Обратите внимание, что несколько объектов соединения могут быть опрошены сразу с помощью multiprocessing.connection.wait().

send_bytes(buffer[, offset[, size]])

Отправить данные байтов из объекта типа байты как полное сообщение.

Если задан offset, данные читаются с этой позиции в buffer. Если задан size, то столько байтов будет прочитано из буфера. Очень большие буферы (приблизительно 32 МБ+, хотя это зависит от ОС) могут вызвать исключение ValueError

recv_bytes([maxlength])

Возвращает полное сообщение данных байтов, отправленное с другого конца соединения, как строку. Блокируется до тех пор, пока не будет чего-то для получения. Вызывает EOFError, если ничего не осталось для получения, и другой конец закрыт.

Если задан maxlength, и сообщение длиннее, чем maxlength, то возникает OSError, и соединение больше не будет читаемым.

Изменено в версии 3.3: Эта функция раньше вызывала IOError, которая теперь является псевдонимом OSError.

recv_bytes_into(buffer[, offset])

Прочитать в buffer полное сообщение данных байтов, отправленное с другого конца соединения, и вернуть количество байтов в сообщении. Блокируется до тех пор, пока не будет чего-то для получения. Вызывает EOFError, если ничего не осталось для получения, и другой конец был закрыт.

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

Если буфер слишком короткий, то возникает исключение 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() автоматически рассириализует полученные данные, что может представлять собой риск безопасности, если вы не доверяете процессу, который отправил сообщение.

Поэтому, если объект соединения не был получен с помощью Pipe() , вы должны использовать только методы recv() и send() после выполнения некоторой аутентификации. См. Ключи аутентификации.

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

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

Примитивы синхронизации

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

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

class multiprocessing.Barrier(parties[, action[, timeout]])

Объект барьера: копия объекта threading.Barrier.

Добавлен в версии 3.3.

class multiprocessing.BoundedSemaphore([value])

Объект ограниченного семафора: близкий аналог threading.BoundedSemaphore.

Единственное различие от его близкого аналога: первый аргумент метода acquire называется block, как и в Lock.acquire().

Примечание

В 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, равным None (по умолчанию), устанавливают временной интервал таймаута в бесконечность. Обратите внимание, что обработка отрицательных или None значений для timeout отличается от реализованного поведения в threading.Lock.acquire(). Аргумент timeout не имеет практического значения, если аргумент block установлен в False и поэтому игнорируется. Возвращает True если блокировка была получена или False если истек временной интервал таймаута.

release()

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

Поведение такое же, как в threading.Lock.release(), за исключением того, что при вызове для разблокированной блокировки генерируется ValueError.

class multiprocessing.RLock

Объект рекурсивной блокировки: близкий аналог threading.RLock. Рекурсивную блокировку должен освободить процесс или поток, который ее получил. После того как процесс или поток получил рекурсивную блокировку, тот же процесс или поток может получить ее снова без блокировки; этот процесс или поток должен освободить ее один раз для каждого раза, когда она была получена.

Обратите внимание, что RLock фактически является функцией-фабрикой, которая возвращает экземпляр multiprocessing.synchronize.RLock, инициализированный по умолчанию.

RLock поддерживает протокол менеджера контекста и поэтому может использоваться в операторах with.

acquire(block=True, timeout=None)

Получить блокировку, блокируя или не блокируя.

При вызове с аргументом block, равным True, заблокироваться до тех пор, пока блокировка не будет разблокирована (не принадлежит ни одному процессу или потоку), если блокировка уже не принадлежит текущему процессу или потоку. Текущий процесс или поток затем получает владение блокировкой (если он еще не владеет ею) и уровень рекурсии внутри блокировки увеличивается на единицу, что приводит к возвращаемому значению True. Обратите внимание, что поведение этого первого аргумента отличается от реализации threading.RLock.acquire(), начиная с самого имени аргумента.

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

Использование и поведение аргумента timeout такие же, как в Lock.acquire(). Обратите внимание, что некоторые из этих поведений timeout отличаются от реализованных поведений в threading.RLock.acquire().

release()

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

Вызывайте этот метод только тогда, когда вызывающий процесс или поток владеет блокировкой. Если этот метод вызван процессом или потоком, отличным от владельца, или если блокировка находится в разблокированном (не принадлежащем) состоянии, возникает AssertionError. Обратите внимание, что тип исключения, генерируемого в этой ситуации, отличается от реализованного поведения в threading.RLock.release().

END_OF_DOCUMENT_MARKER
class multiprocessing.Semaphore([value])

Объект семафора: близкий аналог threading.Semaphore.

Единственное отличие от его близкого аналога состоит в том, что аргумент метода acquire называется block, как и в Lock.acquire().

Примечание

В 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. Если size_or_initializer — целое число, то оно определяет длину массива, и массив будет первоначально заполнен нулями. В противном случае size_or_initializer — это последовательность, которая используется для инициализации массива, а её длина определяет длину массива.

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

Обратите внимание, что lock — это ключевое слово-только аргумент.

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

The multiprocessing.sharedctypes module

The multiprocessing.sharedctypes module provides functions for allocating ctypes objects from shared memory which can be inherited by child processes.

Примечание

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

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)

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

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

Обратите внимание, что lock — ключевой аргумент.

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

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

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

Обратите внимание, что lock — ключевой аргумент.

multiprocessing.sharedctypes.copy(obj)

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

multiprocessing.sharedctypes.synchronized(obj[, lock])

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

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

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

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

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

ctypes

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

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

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 (использовать текущий контекст). См. функцию get_context().

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_ , если оно существует.) В случае, если список 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 и возвращает прокси для него.

Изменено в версии 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() вернёт представление прокси.

Важной особенностью объектов-прокси является то, что они сериализуемы (picklable), поэтому их можно передавать между процессами. Таким образом, референт может содержать объекты-прокси. Это позволяет вкладывать эти управляемые списки, словари и другие объекты-прокси:

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

Возврат копии референта.

Если референт не сериализуем (unpicklable), это вызовет исключение.

__repr__()

Возврат представления объекта-прокси.

__str__()

Возврат представления референта.

Очистка

Объект-прокси использует обратный вызов weakref, чтобы при удалении его из мусора он дерегистрировался у менеджера, который владеет его референтом.

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

Пулы процессов

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

class multiprocessing.pool.Pool([processes[, initializer[, initargs[, maxtasksperchild[, context]]]]])

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

processes — количество рабочих процессов для использования. Если processes равно None , то используется число, возвращаемое os.cpu_count().

Если initializer не None , то каждый рабочий процесс вызовет initializer(*initargs) при запуске.

maxtasksperchild — количество задач, которые может выполнить рабочий процесс, прежде чем он завершит работу и будет заменён новым, чтобы освободить неиспользуемые ресурсы. По умолчанию maxtasksperchild равно None, что означает, что рабочие процессы будут существовать до тех пор, пока существует пул.

context может использоваться для указания контекста, используемого для запуска рабочих процессов. Обычно пул создаётся с помощью функции multiprocessing.Pool() или метода Pool() объекта контекста. В обоих случаях context устанавливается соответствующим образом.

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

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

multiprocessing.pool объекты имеют внутренние ресурсы, которые необходимо правильно управлять (как и любыми другими ресурсами), используя пул как контекстный менеджер или вручную вызывая close() и terminate(). Несоблюдение этого может привести к зависанию процесса при завершении.

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

Изменено в версии 3.2: Добавлен параметр maxtasksperchild.

Изменено в версии 3.4: Добавлен параметр context.

Примечание

Рабочие процессы в 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, если результат не может быть возвращён в течение timeout секунд.

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

Подождите, пока рабочие процессы завершат свою работу. Необходимо вызвать close() или terminate() перед использованием join().

Изменено в версии 3.3: Объекты Pool теперь поддерживают протокол управления контекстом — см. Типы менеджеров контекста. __enter__() возвращает объект пула, а __exit__() вызывает terminate().

class multiprocessing.pool.AsyncResult

Класс результата, возвращаемого Pool.apply_async() и Pool.map_async().

get([timeout])

Возвращает результат при его поступлении. Если timeout не None и результат не поступает в течение timeout секунд, то будет поднято исключение multiprocessing.TimeoutError. Если при удалённом вызове возникло исключение, то это исключение будет повторно поднято методом get().

wait([timeout])

Ожидает, пока результат не станет доступным или пока не пройдёт timeout секунд.

ready()

Возвращает, завершился ли вызов.

successful()

Возвращает, завершился ли вызов без возникновения исключения. Если результат не готов, будет поднято исключение ValueError.

Изменено в версии 3.7: Если результат не готов, то вместо AssertionError будет поднято исключение ValueError.

Следующий пример демонстрирует использование пула:

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' и адрес None, то сокет будет создан в личном временном каталоге, созданном с помощью tempfile.mkstemp().

Если объект слушателя использует сокет, то backlog (по умолчанию 1) передаётся методу listen() сокета после его привязки.

Если authkey задан и не None, он должен быть строкой байтов и будет использоваться в качестве секретного ключа для вызова аутентификации на основе HMAC. Если authkey None, аутентификация не выполняется. AuthenticationError генерируется, если аутентификация не удалась. См. Ключи аутентификации.

accept()

Принимает соединение на связанном сокете или именованном канале объекта слушателя и возвращает объект Connection. Если попытка аутентификации не удалась, генерируется исключение AuthenticationError.

close()

Закрывает привязанный сокет или именованный канал объекта слушателя. Это вызывается автоматически при сборе мусора слушателя. Однако рекомендуется вызывать его явно.

Объекты слушателей имеют следующие только для чтения свойства:

address

Адрес, используемый объектом слушателя.

last_accepted

Адрес, из которого поступило последнее принятое соединение. Если это недоступно, то None.

Изменено в версии 3.3: Объекты слушателей теперь поддерживают протокол управления контекстом — см. Типы менеджеров контекста. __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() для одновременного ожидания сообщений от нескольких процессов:

import time, random
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'. Для подключения к именованной pipe на удаленном компьютере с именем ServerName следует использовать адрес вида r'\\ServerName\pipe\PipeName'.

Обратите внимание, что любая строка, начинающаяся с двух обратных слешей, по умолчанию рассматривается как адрес 'AF_PIPE', а не как адрес 'AF_UNIX'.

Ключи аутентификации

При использовании Connection.recv, полученные данные автоматически распаковываются. К сожалению, распаковка данных из ненадежного источника представляет собой риск для безопасности. Поэтому Listener и Client() используют модуль hmac для аутентификации с помощью дайджеста.

Ключ аутентификации — это строка байтов, которую можно рассматривать как пароль: после установления соединения оба конца запросят подтверждение, что другой знает ключ аутентификации. (Демонстрация того, что оба конца используют один и тот же ключ, не подразумевает отправку ключа по соединению.)

Если запрошена аутентификация, но ключ аутентификации не указан, используется возвращаемое значение current_process().authkey (см. Process). Это значение автоматически наследуется любым объектом Process, который создает текущий процесс. Это означает, что (по умолчанию) все процессы многопроцессной программы будут использовать один и тот же ключ аутентификации, который можно использовать при установлении соединений между ними.

Подходящие ключи аутентификации также можно сгенерировать с помощью os.urandom().

Ведение журнала

Доступна некоторая поддержка ведения журнала. Обратите внимание, что пакет logging не использует общие блокировки процессов, поэтому (в зависимости от типа обработчика) сообщения из разных процессов могут быть перемешаны.

multiprocessing.get_logger()

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

При первом создании логгер имеет уровень logging.NOTSET и не имеет обработчика по умолчанию. Сообщения, отправленные в этот логгер, по умолчанию не передаются в корневой логгер.

Обратите внимание, что в системах Windows дочерние процессы будут наследовать только уровень логгера родительского процесса — любые другие настройки логгера не будут унаследованы.

multiprocessing.log_to_stderr(level=None)

Эта функция выполняет вызов 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.cpu_count().

Если initializer не None, каждый поток-работник будет вызывать initializer(*initargs) при запуске.

В отличие от Pool, maxtasksperchild и context предоставить нельзя.

Примечание

ThreadPool разделяет тот же интерфейс, что и Pool, который разработан для пула процессов и предшествовал появлению модуля concurrent.futures. Вследствие этого он наследует некоторые операции, которые не имеют смысла для пула, управляемого потоками, и имеет собственный тип для представления состояния асинхронных заданий, AsyncResult, который не понятен другим библиотекам.

Пользователям рекомендуется использовать concurrent.futures.ThreadPoolExecutor, который имеет более простой интерфейс, разработанный с самого начала для работы с потоками, и который возвращает экземпляры concurrent.futures.Future, совместимые со многими другими библиотеками, включая asyncio.

Руководящие принципы программирования

Существуют определенные рекомендации и приемы, которых следует придерживаться при использовании multiprocessing.

Все методы запуска

Следующее относится ко всем методам запуска.

Избегайте общих состояний

По возможности следует избегать передачи больших объемов данных между процессами.

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

Пиклируемость

Убедитесь, что аргументы методов прокси-объектов пиклируемы.

Безопасность потоков прокси-объектов

Не используйте объект прокси из нескольких потоков, если не защищаете его с помощью блокировки.

(Проблема с использованием одного и того же прокси-объекта различными процессами никогда не возникает.)

Объединение процессов-зомби

В системах POSIX, когда процесс завершается, но не объединён, он становится процессом-зомби. Их не должно быть слишком много, поскольку каждый раз при запуске нового процесса (или вызове active_children()) все завершенные, но не объединённые процессы будут объединены. Также вызов метода Process.is_alive завершенного процесса объединит этот процесс. Тем не менее, вероятно, хорошей практикой является явное объединение всех запущенных процессов.

Наследование предпочтительнее пиклирования/депиклирования

При использовании методов запуска spawn или forkserver многие типы из multiprocessing должны быть пиклируемыми, чтобы дочерние процессы могли их использовать. Однако следует избегать отправки общих объектов в другие процессы с помощью каналов или очередей. Вместо этого следует организовать программу таким образом, чтобы процесс, которому нужен доступ к общему ресурсу, созданному где-то в другом месте, унаследовал его от родительского процесса.

Избегайте завершения процессов

Использование метода 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() для этого объекта типа файла, это может привести к многократному сбросу одних и тех же данных в объект, что может привести к повреждению данных.

Если вы пишете объект типа файла и реализуете свой кэширование, вы можете сделать его безопасным при разветвлении, сохраняя 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.

Более высокая пиклируемость

Убедитесь, что все аргументы для Process.__init__() пиклируемы. Также, если вы наследуете от Process, убедитесь, что экземпляры будут пиклируемыми при вызове метода 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–2024 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.12/library/multiprocessing.html

Spec-Zone.ru

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