Spec-Zone.ru › Python 3.11

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

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

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

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

Введение

multiprocessing — это пакет, который поддерживает запуск процессов с помощью API, аналогичного модулю threading. Пакет multiprocessing обеспечивает как локальное, так и удалённое конкурентное выполнение, эффективно обходя Глобальную блокировку интерпретатора путём использования дочерних процессов вместо потоков. Благодаря этому, модуль multiprocessing позволяет программисту в полной мере использовать несколько процессоров на данном компьютере. Он работает как на Unix, так и на 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.

Доступно на Unix и Windows. По умолчанию на Windows и macOS.

fork

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

Доступно только на Unix. По умолчанию на Unix.

forkserver

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

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

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

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

На Unix с использованием методов запуска 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) на Unix. Метод запуска '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()

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

Трубы

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

Синхронизация между процессами

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__' (built-in)>
AttributeError: Can't get attribute 'f' on <module '__main__' (built-in)>
AttributeError: Can't get attribute 'f' on <module '__main__' (built-in)>

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

Справочник

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

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

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

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

terminate()

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

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

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

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

kill()

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

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

close()

Закройте объект Process, освободив все ресурсы, связанные с ним. ValueError возникает, если дочерний процесс всё ещё выполняется. После успешного возврата из close(), большинство других методов и атрибутов объекта Process будут вызывать ValueError.

Новая функция в версии 3.7.

Обратите внимание, что методы start(), join(), is_alive(), terminate() и exitcode должны вызываться только процессом, который создал объект процесса.

Пример использования некоторых методов объекта Process:

>>> import multiprocessing, time, signal
>>> p = multiprocessing.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() для каждой задачи, удаленной из очереди, в противном случае семафор, используемый для подсчёта количества незавершенных задач, может переполниться, вызвав исключение.

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

Примечание

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 — только для отправки сообщений.

class multiprocessing.Queue([maxsize])

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

Обычно исключения queue.Empty и queue.Full из модуля стандартной библиотеки queue используются для сигнализации о таймаутах.

Queue реализует все методы queue.Queue, за исключением task_done() и join().

qsize()

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

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

empty()

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

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

get()

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

put(item)

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

class multiprocessing.JoinableQueue([maxsize])

JoinableQueue, a Queue subclass, is a queue which additionally has task_done() and join() methods.

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', 'forkserver'. В Windows доступен только 'spawn'. В Unix всегда поддерживаются 'fork' и 'spawn', при этом 'fork' является значением по умолчанию.

Введено в версии 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 ложно, то метод запуска устанавливается по умолчанию, и возвращается его имя. Если метод запуска не определён и allow_none истинно, то возвращается None. Возвращаемое значение может быть 'fork', 'spawn', 'forkserver' или None. 'fork' — значение по умолчанию в Unix, а 'spawn' — в Windows и macOS.

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

Введено в версии 3.4.

multiprocessing.set_executable(executable)

Устанавливает путь к интерпретатору Python, который будет использоваться при запуске дочерних процессов. (По умолчанию используется sys.executable). Встраивающим системам, скорее всего, потребуется что-то вроде

set_executable(os.path.join(sys.exec_prefix, 'pythonw.exe'))

перед созданием дочерних процессов.

Изменено в версии 3.4: Теперь поддерживается в Unix при использовании метода запуска 'spawn'.

Изменено в версии 3.11: Принимает объект, похожий на путь.

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

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

Примечание

Если сигнал SIGINT, сгенерированный комбинацией Ctrl-C, поступит во время блокировки основного потока вызовом BoundedSemaphore.acquire(), Lock.acquire(), RLock.acquire(), Semaphore.acquire(), Condition.acquire() или Condition.wait(), то вызов будет немедленно прерван, и будет поднято исключение KeyboardInterrupt.

Это отличается от поведения модуля threading, где SIGINT будет проигнорирован во время выполнения эквивалентных блокирующих вызовов.

Примечание

Некоторые функции этого пакета требуют работоспособной реализации общих семафоров в операционной системе хоста. Без неё модуль 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, которые позволяют использовать его для хранения и извлечения строк.

Модуль multiprocessing.sharedctypes

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

Примечание

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

multiprocessing.sharedctypes.RawArray(typecode_or_type, size_or_initializer)

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

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

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

multiprocessing.sharedctypes.RawValue(typecode_or_type, *args)

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

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

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

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

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

То же, что и 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 с использованием кода типа

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.

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

Однако при использовании прокси для объекта пространства имён атрибут, начинающийся с '_' , будет атрибутом прокси, а не атрибутом объекта-оригинала:

>>> manager = multiprocessing.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()

Объекты-прокси

Прокси — это объект, который ссылается на общий объект, который (предположительно) существует в другом процессе. Общий объект называется объектом-ссылочным для прокси. Несколько объектов-прокси могут иметь один и тот же объект-ссылочный.

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

>>> from multiprocessing import Manager
>>> manager = 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() вернёт представление прокси.

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

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

Возвращает копию объекта-ссылочного.

Если объект-ссылочный не пикелируем, то это вызовет исключение.

__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() вызывается немедленно.

END_OF_DOCUMENT_MARKER
join()

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

New in version 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. ValueError будет выброшено вместо AssertionError если результат не готов.

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

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

from multiprocessing import Pool
import time

def f(x):
    return x*x

if __name__ == '__main__':
    with Pool(processes=4) as pool:         # start 4 worker processes
        result = pool.apply_async(f, (10,)) # evaluate "f(10)" asynchronously in a single process
        print(result.get(timeout=1))        # prints "100" unless your computer is *very* slow

        print(pool.map(f, range(10)))       # prints "[0, 1, 4,..., 81]"

        it = pool.imap(f, range(10))
        print(next(it))                     # prints "0"
        print(next(it))                     # prints "1"
        print(it.next(timeout=1))           # prints "4" unless your computer is *very* slow

        result = pool.apply_async(time.sleep, (10,))
        print(result.get(timeout=1))        # raises multiprocessing.TimeoutError

Слушатели и клиенты

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

Однако модуль multiprocessing.connection позволяет добавить некоторую гибкость. Он в основном предоставляет API ориентированный на сообщения высокого уровня для работы с сокетами или именованными каналами Windows. Он также поддерживает *аутентификацию по хешу* с использованием модуля hmac и для опроса нескольких подключений одновременно.

multiprocessing.connection.deliver_challenge(connection, authkey)

Отправить случайно сгенерированное сообщение на другой конец соединения и дождаться ответа.

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

multiprocessing.connection.answer_challenge(connection, authkey)

Получить сообщение, вычислить его хеш с использованием authkey в качестве ключа и затем отправить хеш обратно.

Если приветственное сообщение не получено, то возбуждается исключение AuthenticationError.

multiprocessing.connection.Client(address[, family[, authkey]])

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

Тип соединения определяется аргументом family, но его можно обычно опустить, так как его обычно можно вывести из формата address. (См. Форматы адресов)

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

class multiprocessing.connection.Listener([address[, family[, backlog[, authkey]]]])

Обёртка для сокета или именованного канала Windows, который «слушает» подключения.

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

Примечание

Если используется адрес ‘0.0.0.0’, адрес не будет являться соединяемой точкой на Windows. Если вам требуется соединяемая точка, используйте ‘127.0.0.1’.

family — тип сокета (или именованного канала) для использования. Может быть одной из строк 'AF_INET' (для TCP сокета), 'AF_UNIX' (для Unix-доменного сокета) или 'AF_PIPE' (для именованного канала Windows). Из них только первый гарантированно доступен. Если family равно None, тип определяется из формата address. Если address также None, выбирается значение по умолчанию. Это значение по умолчанию — тип, который предположительно является самым быстрым доступным. См. Форматы адресов. Обратите внимание, что если family равно 'AF_UNIX' и address равно None, сокет будет создан в частном временном каталоге, созданном с помощью tempfile.mkstemp().

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

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

accept()

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

close()

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

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

address

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

last_accepted

Адрес, с которого пришло последнее принятое подключение. Если он недоступен, то None.

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

multiprocessing.connection.wait(object_list, timeout=None)

Дождаться, пока объект в object_list не подготовится. Возвращает список объектов в object_list, которые готовы. Если timeout — число с плавающей точкой, то вызов блокируется не более чем на указанное количество секунд. Если timeout равно None, то вызов блокируется неограниченное время. Отрицательное значение timeout эквивалентно нулевому.

Для Unix и Windows объект может появиться в object_list, если он:

  • читаемый объект Connection;
  • подключённый и читаемый объект socket.socket; или
  • атрибут sentinel объекта Process.

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

Unix: 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'. Для подключения к именованной трубе на удалённом компьютере под названием 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 следует придерживаться определённых рекомендаций и паттернов.

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

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

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

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

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

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

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

Потокобезопасность прокси

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

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

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

В Unix, если процесс завершается, но не был объединён, он становится «зомби». Их не должно быть слишком много, так как каждый раз при запуске нового процесса (или вызове 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()).

Явное передача ресурсов дочерним процессам

В Unix при использовании метода запуска 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–2023 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.11/library/multiprocessing.html

Spec-Zone.ru

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