Spec-Zone.ru › Python 3.13

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

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

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

Этот модуль не поддерживается на мобильных платформах или платформах WebAssembly.

Введение

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

Модуль multiprocessing также вводит API, у которого нет аналогов в модуле threading. Ярким примером является объект Pool, который предоставляет удобный способ распараллеливания выполнения функции по нескольким входным значениям, распределяя входные данные по процессам (параллелизм данных). Следующий пример демонстрирует распространённую практику определения таких функций в модуле, чтобы дочерние процессы могли успешно импортировать этот модуль. Этот базовый пример параллелизма данных с использованием Pool,

from multiprocessing import Pool

def f(x):
    return x*x

if __name__ == '__main__':
    with Pool(5) as p:
        print(p.map(f, [1, 2, 3]))

выведет на стандартный вывод

[1, 4, 9]

См. также

concurrent.futures.ProcessPoolExecutor предлагает интерфейс более высокого уровня для передачи задач в фоновый процесс без блокировки выполнения вызывающего процесса. По сравнению с прямым использованием интерфейса Pool, API concurrent.futures более удобно позволяет отделить отправку работы в подпроцесс от ожидания результатов.

Класс Process

В multiprocessing процессы порождаются путём создания объекта Process, а затем вызова его метода start(). Process следует API threading.Thread. Тривиальным примером многопроцессной программы является

from multiprocessing import Process

def f(name):
    print('hello', name)

if __name__ == '__main__':
    p = Process(target=f, args=('bob',))
    p.start()
    p.join()

Чтобы показать отдельные идентификаторы процессов, вот расширенный пример:

from multiprocessing import Process
import os

def info(title):
    print(title)
    print('module name:', __name__)
    print('parent process:', os.getppid())
    print('process id:', os.getpid())

def f(name):
    info('function f')
    print('hello', name)

if __name__ == '__main__':
    info('main line')
    p = Process(target=f, args=('bob',))
    p.start()
    p.join()

Для объяснения того, почему необходима часть if __name__ == '__main__', см. Руководство по программированию.

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

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

spawn

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

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

fork

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

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

Примечание

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

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

forkserver

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

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

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

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

На платформах POSIX при использовании методов запуска spawn или forkserver также запускается процесс отслеживания ресурсов, который отслеживает неразрывные именованные системные ресурсы (например, именованные семафоры или SharedMemory объекты) созданные процессами программы. Когда все процессы завершат работу, процесс отслеживания ресурсов разрывает любые оставшиеся отслеживаемые объекты. Обычно их не должно быть, но если процесс был убит сигналом, могут быть некоторые «утекшие» ресурсы. (Ни утекшие семафоры, ни сегменты общей памяти не будут автоматически разрываться до следующей перезагрузки. Это проблема для обоих объектов, поскольку система позволяет только ограниченное количество именованных семафоров, а сегменты общей памяти занимают некоторое место в основной памяти.)

Чтобы выбрать метод запуска, используйте set_start_method() в разделе if __name__ == '__main__' основного модуля. Например:

import multiprocessing as mp

def foo(q):
    q.put('hello')

if __name__ == '__main__':
    mp.set_start_method('spawn')
    q = mp.Queue()
    p = mp.Process(target=foo, args=(q,))
    p.start()
    print(q.get())
    p.join()

set_start_method() не должен использоваться более одного раза в программе.

В качестве альтернативы, вы можете использовать get_context() для получения объекта контекста. Объекты контекста имеют тот же API, что и модуль multiprocessing, и позволяют использовать несколько методов запуска в одной программе.

import multiprocessing as mp

def foo(q):
    q.put('hello')

if __name__ == '__main__':
    ctx = mp.get_context('spawn')
    q = ctx.Queue()
    p = ctx.Process(target=foo, args=(q,))
    p.start()
    print(q.get())
    p.join()

Обратите внимание, что объекты, относящиеся к одному контексту, могут быть несовместимы с процессами другого контекста. В частности, блокировки, созданные с использованием контекста fork, нельзя передавать процессам, запущенным с методами запуска spawn или forkserver.

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

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

Методы запуска 'spawn' и 'forkserver' обычно нельзя использовать с «замороженными» исполняемыми файлами (т. е., двоичными файлами, созданными такими пакетами, как PyInstaller и cx_Freeze) на системах POSIX. Метод запуска 'fork' может работать, если код не использует потоки.

Обмен объектами между процессами

multiprocessing поддерживает два типа каналов связи между процессами:

Очереди

Класс Queue является почти точной копией класса queue.Queue. Например:

from multiprocessing import Process, Queue

def f(q):
    q.put([42, None, 'hello'])

if __name__ == '__main__':
    q = Queue()
    p = Process(target=f, args=(q,))
    p.start()
    print(q.get())    # prints "[42, None, 'hello']"
    p.join()

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

Трубы

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

from multiprocessing import Process, Pipe

def f(conn):
    conn.send([42, None, 'hello'])
    conn.close()

if __name__ == '__main__':
    parent_conn, child_conn = Pipe()
    p = Process(target=f, args=(child_conn,))
    p.start()
    print(parent_conn.recv())   # prints "[42, None, 'hello']"
    p.join()

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

Метод send() сериализует объект, а метод recv() восстанавливает объект.

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

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

from multiprocessing import Process, Lock

def f(l, i):
    l.acquire()
    try:
        print('hello world', i)
    finally:
        l.release()

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

    for num in range(10):
        Process(target=f, args=(lock, num)).start()

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

Обмен состоянием между процессами

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

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

Общий доступ к памяти

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

from multiprocessing import Process, Value, Array

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

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

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

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

будет выводить

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

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

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

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

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

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

from multiprocessing import Process, Manager

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

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

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

        print(d)
        print(l)

будет выводить

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

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

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

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

Например:

from multiprocessing import Pool, TimeoutError
import time
import os

def f(x):
    return x*x

if __name__ == '__main__':
    # start 4 worker processes
    with Pool(processes=4) as pool:

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

        # print same numbers in arbitrary order
        for i in pool.imap_unordered(f, range(10)):
            print(i)

        # evaluate "f(20)" asynchronously
        res = pool.apply_async(f, (20,))      # runs in *only* one process
        print(res.get(timeout=1))             # prints "400"

        # evaluate "os.getpid()" asynchronously
        res = pool.apply_async(os.getpid, ()) # runs in *only* one process
        print(res.get(timeout=1))             # prints the PID of that process

        # launching multiple evaluations asynchronously *may* use more processes
        multiple_results = [pool.apply_async(os.getpid, ()) for i in range(4)]
        print([res.get(timeout=1) for res in multiple_results])

        # make a single worker sleep for 10 seconds
        res = pool.apply_async(time.sleep, (10,))
        try:
            print(res.get(timeout=1))
        except TimeoutError:
            print("We lacked patience and got a multiprocessing.TimeoutError")

        print("For the moment, the pool remains available for more work")

    # exiting the 'with'-block has stopped the pool
    print("Now the pool is closed and no longer available")

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

Примечание

Функциональность этого пакета требует, чтобы модуль __main__ был импортируемым для дочерних процессов. Это описывается в Руководстве по программированию, но стоит отметить здесь. Это означает, что некоторые примеры, такие как примеры multiprocessing.pool.Pool, не будут работать в интерактивном интерпретаторе. Например:

>>> from multiprocessing import Pool
>>> p = Pool(5)
>>> def f(x):
...     return x*x
...
>>> with p:
...     p.map(f, [1,2,3])
Process PoolWorker-1:
Process PoolWorker-2:
Process PoolWorker-3:
Traceback (most recent call last):
Traceback (most recent call last):
Traceback (most recent call last):
AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)>
AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)>
AttributeError: Can't get attribute 'f' on <module '__main__' (<class '_frozen_importlib.BuiltinImporter'>)>

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

Ссылки

Пакет multiprocessing в основном дублирует API модуля threading.

Process и исключения

class multiprocessing.Process(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)

Объекты Process представляют активность, выполняемую в отдельном процессе. Класс Process имеет эквиваленты всех методов класса threading.Thread.

Конструктор всегда должен вызываться с ключевыми аргументами. group всегда должен быть None; он существует только для совместимости с threading.Thread. target — вызываемый объект, который будет вызван методом run(). По умолчанию он равен None, что означает, что ничего не вызывается. name — имя процесса (подробнее см. name). args — кортеж аргументов для вызова целевого объекта. kwargs — словарь ключевых аргументов для вызова целевого объекта. Если указан, ключевой аргумент daemon устанавливает флаг процесса daemon в значение True или False. Если None (значение по умолчанию), этот флаг будет унаследован от создающего процесса.

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

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

terminate()

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

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

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

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

kill()

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

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

close()

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

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

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

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

>>> import multiprocessing, time, signal
>>> mp_context = multiprocessing.get_context('spawn')
>>> p = mp_context.Process(target=time.sleep, args=(1000,))
>>> print(p, p.is_alive())
<...Process ... initial> False
>>> p.start()
>>> print(p, p.is_alive())
<...Process ... started> True
>>> p.terminate()
>>> time.sleep(0.1)
>>> print(p, p.is_alive())
<...Process ... stopped exitcode=-SIGTERM> False
>>> p.exitcode == -signal.SIGTERM
True
exception multiprocessing.ProcessError

Базовый класс всех исключений multiprocessing.

exception multiprocessing.BufferTooShort

Исключение, генерируемое Connection.recv_bytes_into() , когда предоставленный буфер слишком мал для чтения сообщения.

Если e является экземпляром BufferTooShort, то e.args[0] вернёт сообщение в виде байтовой строки.

exception multiprocessing.AuthenticationError

Генерируется при ошибке аутентификации.

exception multiprocessing.TimeoutError

Генерируется методами с таймаутом при истечении таймаута.

Каналы и очереди

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

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

Типы Queue, SimpleQueue и JoinableQueue — это многопоточные очереди «первый пришёл — первый ушёл» (FIFO), моделируемые по классу queue.Queue стандартной библиотеки. Они отличаются тем, что в Queue отсутствуют методы task_done() и join(), введённые в класс queue.Queue Python 2.5.

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

Отличие от других реализаций очередей Python заключается в том, что в multiprocessing все объекты, помещаемые в очереди, сериализуются с помощью pickle. Возвращаемый методом get объект — это воссозданный объект, который не использует общую память с оригинальным объектом.

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

Примечание

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

Примечание

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

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

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

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

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

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

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

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

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

multiprocessing.Pipe([duplex])

Возвращает пару (conn1, conn2) объектов Connection, представляющих концы канала.

Если duplex — True (по умолчанию), то канал двунаправленный. Если duplex — False , то канал однонаправленный: conn1 может использоваться только для получения сообщений, а conn2 — только для отправки сообщений.

Метод send() сериализует объект с помощью pickle, а метод recv() воссоздаёт объект.

END_OF_DOCUMENT_MARKER
class multiprocessing.Queue([maxsize])

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

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

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

qsize()

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

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

empty()

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

Может вызвать OSError для закрытых очередей. (не гарантируется)

full()

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

put(obj[, block[, timeout]])

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

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

put_nowait(obj)

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

get([block[, timeout]])

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

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

get_nowait()

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

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

close()

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

join_thread()

Присоединиться к фоновой нити. Этот метод можно использовать только после вызова close(). Блокируется, пока фоновая нить не завершит работу, гарантируя, что все данные в буфере отправлены по каналу.

По умолчанию, если процесс не является создателем очереди, при выходе он попытается присоединиться к фоновой нити очереди. Процесс может вызвать cancel_join_thread() для отмены действия join_thread().

cancel_join_thread()

Предотвратить блокировку join_thread(). В частности, это предотвращает автоматическое присоединение фоновой нити к процессу при выходе — см. join_thread().

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

Примечание

Функциональность этого класса требует работоспособной реализации общего семафора в операционной системе. Без неё функциональность этого класса будет отключена, и попытки создать Queue приведут к ImportError. См. bpo-3770 для дополнительной информации. То же самое относится к любым специализированным типам очередей, перечисленных ниже.

class multiprocessing.SimpleQueue

Это упрощённый тип Queue, очень похожий на защищённый Pipe.

close()

Закрыть очередь: освободить внутренние ресурсы.

Очередь больше нельзя использовать после закрытия. Например, методы get(), put() и empty() больше вызывать нельзя.

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

empty()

Возвращает True, если очередь пуста, False, в противном случае.

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

get()

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

put(item)

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

class multiprocessing.JoinableQueue([maxsize])

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

task_done()

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

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

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

join()

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

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

Разное

multiprocessing.active_children()

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

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

multiprocessing.cpu_count()

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

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

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

См. также

os.cpu_count() os.process_cpu_count()

Изменено в версии 3.13: Значение возвращаемого значения также может быть переопределено с помощью флага -X cpu_count или переменной среды PYTHON_CPU_COUNT, так как это просто обертка вокруг API подсчета ЦП из модуля os.

multiprocessing.current_process()

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

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

multiprocessing.parent_process()

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

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

multiprocessing.freeze_support()

Добавляет поддержку для случаев, когда программа, использующая multiprocessing, была заморожена для создания исполняемого файла Windows. (Тестировалось с py2exe, PyInstaller и cx_Freeze.)

Необходимо вызвать эту функцию сразу после строки if __name__ == '__main__' в основном модуле. Например:

from multiprocessing import Process, freeze_support

def f():
    print('hello world!')

if __name__ == '__main__':
    freeze_support()
    Process(target=f).start()

Если строка freeze_support() пропущена, то при запуске замороженного исполняемого файла будет выброшено исключение RuntimeError.

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

multiprocessing.get_all_start_methods()

Возвращает список поддерживаемых методов запуска, первый из которых является по умолчанию. Возможные методы запуска: 'fork', 'spawn' и 'forkserver'. Не все платформы поддерживают все методы. См. Контексты и методы запуска.

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

multiprocessing.get_context(method=None)

Возвращает объект контекста, у которого есть те же атрибуты, что и у модуля multiprocessing.

Если method равно None, возвращается контекст по умолчанию. В противном случае method должен быть 'fork', 'spawn', 'forkserver'. Если указанный метод запуска недоступен, генерируется ValueError. См. Контексты и методы запуска.

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

multiprocessing.get_start_method(allow_none=False)

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

Если метод запуска не был задан и allow_none равно false, то метод запуска фиксируется по умолчанию, и возвращается его имя. Если метод запуска не был задан и allow_none равно true, то возвращается None.

Возвращаемое значение может быть 'fork', 'spawn', 'forkserver' или None. См. Контексты и методы запуска.

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

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

multiprocessing.set_executable(executable)

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

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

прежде чем они смогут создать дочерние процессы.

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

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

multiprocessing.set_forkserver_preload(module_names)

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

Для работы необходимо вызвать его до запуска сервера виртуального процесса (до создания Pool или запуска Process).

Актуально только при использовании метода запуска 'forkserver'. См. Контексты и методы запуска.

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

multiprocessing.set_start_method(method, force=False)

Устанавливает метод, который должен использоваться для запуска дочерних процессов. Аргумент method может быть 'fork', 'spawn' или 'forkserver'. Генерирует RuntimeError, если метод запуска уже задан, а force не равно True. Если method равно None и force равно True, то метод запуска устанавливается в None. Если method равно None и force равно False, то контекст устанавливается в контекст по умолчанию.

Обратите внимание, что это следует вызывать не более одного раза, и это должно быть защищено внутри if __name__ == '__main__' блока основного модуля.

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

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

Примечание

multiprocessing не содержит аналогов threading.active_count(), threading.enumerate(), threading.settrace(), threading.setprofile(), threading.Timer или threading.local.

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

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

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

class multiprocessing.connection.Connection
send(obj)

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

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

recv()

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

fileno()

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

close()

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

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

poll([timeout])

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

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

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

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

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

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

recv_bytes([maxlength])

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

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

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

recv_bytes_into(buffer[, offset])

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

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

Если буфер слишком короткий, то возникает исключение BufferTooShort , и полное сообщение доступно в e.args[0], где e — экземпляр исключения.

Изменено в версии 3.3: Теперь сами объекты соединений могут передаваться между процессами с помощью Connection.send() и Connection.recv().

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

Например:

>>> from multiprocessing import Pipe
>>> a, b = Pipe()
>>> a.send([1, 'hello', None])
>>> b.recv()
[1, 'hello', None]
>>> b.send_bytes(b'thank you')
>>> a.recv_bytes()
b'thank you'
>>> import array
>>> arr1 = array.array('i', range(5))
>>> arr2 = array.array('i', [0] * 10)
>>> a.send_bytes(arr1)
>>> count = b.recv_bytes_into(arr2)
>>> assert count == len(arr1) * arr1.itemsize
>>> arr2
array('i', [0, 1, 2, 3, 4, 0, 0, 0, 0, 0])

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

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

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

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

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

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

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

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

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

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

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

class multiprocessing.BoundedSemaphore([value])

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

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

Примечание

В macOS он не отличается от Semaphore, так как sem_getvalue() не реализован на этой платформе.

class multiprocessing.Condition([lock])

Переменная условия: псевдоним для threading.Condition.

Если задан lock, он должен быть объектом Lock или RLock из модуля multiprocessing.

Изменено в версии 3.3: Добавлен метод wait_for().

class multiprocessing.Event

Копия объекта threading.Event.

class multiprocessing.Lock

Объект нерекурсивной блокировки: близкий аналог threading.Lock. После того, как процесс или поток приобрели блокировку, последующие попытки получить ее из любого процесса или потока будут блокироваться до ее освобождения; любой процесс или поток может ее освободить. Концепции и поведение threading.Lock по отношению к потокам воспроизводятся здесь в multiprocessing.Lock по отношению к процессам или потокам, за исключением случаев, оговоренных отдельно.

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

Lock поддерживает протокол менеджера контекста и, таким образом, может использоваться в операторах with.

acquire(block=True, timeout=None)

Получение блокировки, блокирующее или неблокирующее.

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

При значении аргумента block, равном False, вызов метода не будет блокироваться. Если блокировка в настоящее время находится в заблокированном состоянии, возвращает False; в противном случае устанавливает блокировку в заблокированное состояние и возвращает True.

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

release()

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

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

class multiprocessing.RLock

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

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

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

acquire(block=True, timeout=None)

Получение блокировки, блокирующее или неблокирующее.

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

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

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

release()

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

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

class multiprocessing.Semaphore([value])

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

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

Примечание

В macOS функция sem_timedwait не поддерживается, поэтому вызов acquire() с таймаутом эмулирует поведение этой функции с помощью цикла ожидания.

Примечание

Некоторые функции этого пакета требуют работоспособной реализации семафора общего доступа в операционной системе хоста. Без неё модуль multiprocessing.synchronize будет отключён, и попытки импортировать его приведут к ImportError. Дополнительную информацию см. на bpo-3770.

Объекты ctypes общего доступа

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

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

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

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

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

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

counter.value += 1

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

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

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

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

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

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

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

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

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

Модуль 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: Объекты синхронизации поддерживают протокол менеджера контекста.

В таблице ниже сравнивается синтаксис создания объектов shared 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.

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

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

>>> mp_context = multiprocessing.get_context('spawn')
>>> manager = mp_context.Manager()
>>> Global = manager.Namespace()
>>> Global.x = 10
>>> Global.y = 'hello'
>>> Global._z = 12.3    # this is an attribute of the proxy
>>> print(Global)
Namespace(x=10, y='hello')

Настраиваемые менеджеры

Чтобы создать собственный менеджер, необходимо создать подкласс BaseManager и использовать метод класса register() для регистрации новых типов или вызываемых объектов в классе менеджера. Например:

from multiprocessing.managers import BaseManager

class MathsClass:
    def add(self, x, y):
        return x + y
    def mul(self, x, y):
        return x * y

class MyManager(BaseManager):
    pass

MyManager.register('Maths', MathsClass)

if __name__ == '__main__':
    with MyManager() as manager:
        maths = manager.Maths()
        print(maths.add(4, 3))         # prints 7
        print(maths.mul(7, 8))         # prints 56

Использование удалённого менеджера

Можно запустить сервер менеджера на одном компьютере и использовать его клиентами с других компьютеров (при условии, что соответствующие брандмауэры это допускают).

Выполнение следующих команд создает сервер для единственной общей очереди, к которой могут получить доступ удалённые клиенты:

>>> from multiprocessing.managers import BaseManager
>>> from queue import Queue
>>> queue = Queue()
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue', callable=lambda:queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()

Один клиент может получить доступ к серверу следующим образом:

>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.put('hello')

Другой клиент также может использовать его:

>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.get()
'hello'

Локальные процессы также могут получить доступ к этой очереди, используя код из приведённого выше клиентского кода для удалённого доступа:

>>> from multiprocessing import Process, Queue
>>> from multiprocessing.managers import BaseManager
>>> class Worker(Process):
...     def __init__(self, q):
...         self.q = q
...         super().__init__()
...     def run(self):
...         self.q.put('local hello')
...
>>> queue = Queue()
>>> w = Worker(queue)
>>> w.start()
>>> class QueueManager(BaseManager): pass
...
>>> QueueManager.register('get_queue', callable=lambda: queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()

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

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

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

>>> mp_context = multiprocessing.get_context('spawn')
>>> manager = mp_context.Manager()
>>> l = manager.list([i*i for i in range(10)])
>>> print(l)
[0, 1, 4, 9, 16, 25, 36, 49, 64, 81]
>>> print(repr(l))
<ListProxy object, typeid 'list' at 0x...>
>>> l[4]
16
>>> l[2:5]
[4, 9, 16]

Обратите внимание, что применение str() к прокси вернёт представление объекта-ссылка, а применение repr() вернёт представление прокси.

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

>>> a = manager.list()
>>> b = manager.list()
>>> a.append(b)         # referent of a now contains referent of b
>>> print(a, b)
[<ListProxy object, typeid 'list' at ...>] []
>>> b.append('hello')
>>> print(a[0], b)
['hello'] ['hello']

Аналогично, прокси словарей и списков могут быть вложены друг в друга:

>>> l_outer = manager.list([ manager.dict() for i in range(2) ])
>>> d_first_inner = l_outer[0]
>>> d_first_inner['a'] = 1
>>> d_first_inner['b'] = 2
>>> l_outer[1]['c'] = 3
>>> l_outer[1]['z'] = 26
>>> print(l_outer[0])
{'a': 1, 'b': 2}
>>> print(l_outer[1])
{'c': 3, 'z': 26}

Если стандартные (не-прокси) list или dict объекты содержатся в объекте-ссылке, изменения этих изменяемых значений не будут распространяться через менеджер, потому что у прокси нет возможности знать, когда значения внутри изменяются. Однако сохранение значения в контейнерном прокси (что вызывает __setitem__ для объекта-прокси) распространяется через менеджер, и поэтому для эффективного изменения такого элемента можно переназначить изменённое значение контейнерному прокси:

# create a list proxy and append a mutable object (a dictionary)
lproxy = manager.list()
lproxy.append({})
# now mutate the dictionary
d = lproxy[0]
d['a'] = 1
d['b'] = 2
# at this point, the changes to d are not yet synced, but by
# updating the dictionary, the proxy is notified of the change
lproxy[0] = d

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

Примечание

Типы прокси в multiprocessing не поддерживают сравнения по значению. Например:

>>> manager.list([1,2,3]) == [1,2,3]
False

При сравнении следует использовать копию объекта-ссылка.

class multiprocessing.managers.BaseProxy

Объекты-прокси являются экземплярами подклассов BaseProxy.

_callmethod(methodname[, args[, kwds]])

Вызов и возврат результата метода объекта-ссылка прокси.

Если proxy является прокси, чей объект-ссылка obj , то выражение

proxy._callmethod(methodname, args, kwds)

вычислит выражение

getattr(obj, methodname)(*args, **kwds)

в процессе менеджера.

Возвращаемое значение будет копией результата вызова или прокси для нового общего объекта — см. документацию к аргументу method_to_typeid в BaseManager.register().

Если во время вызова возникло исключение, оно перевыбрасывается методом _callmethod(). Если в процессе менеджера возникло другое исключение, оно преобразуется в исключение RemoteError и перевыбрасывается методом _callmethod().

Обратите внимание, что исключение будет возбуждено, если methodname не экспонировано.

Пример использования метода _callmethod():

>>> l = manager.list(range(10))
>>> l._callmethod('__len__')
10
>>> l._callmethod('__getitem__', (slice(2, 7),)) # equivalent to l[2:7]
[2, 3, 4, 5, 6]
>>> l._callmethod('__getitem__', (20,))          # equivalent to l[20]
Traceback (most recent call last):
...
IndexError: list index out of range
_getvalue()

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

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

__repr__()

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

__str__()

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

Очистка

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

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

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

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

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

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

processes — количество рабочих процессов. Если processes равно None, то используется количество, возвращаемое os.process_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.

Изменено в версии 3.13: Параметр processes по умолчанию использует os.process_cpu_count() вместо os.cpu_count().

Примечание

Рабочие процессы в рамках Pool обычно существуют на протяжении всего периода работы очереди задач пула. Частая схема в других системах (например, Apache, mod_wsgi и т.д.) для освобождения ресурсов, удерживаемых рабочими процессами, заключается в том, чтобы разрешить рабочему процессу в пуле выполнить только определённое количество работы, прежде чем он завершится, очистится, а новый процесс будет запущен для замены старого. Аргумент maxtasksperchild в Pool предоставляет пользователю возможность использовать эту возможность.

apply(func[, args[, kwds]])

Вызывает func с аргументами args и ключевыми аргументами kwds. Блокирует выполнение, пока результат не будет готов. Учитывая, что это блокирует выполнение, apply_async() более подходит для выполнения работы параллельно. Кроме того, func выполняется только в одном из рабочих процессов пула.

apply_async(func[, args[, kwds[, callback[, error_callback]]]])

Вариант метода apply(), который возвращает объект AsyncResult.

Если указан callback, он должен быть вызываемым объектом, принимающим один аргумент. Когда результат становится готовым, к нему применяется callback, если только вызов не завершился ошибкой, в этом случае применяется error_callback.

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

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

map(func, iterable[, chunksize])

Параллельный аналог встроенной функции map() (хотя она поддерживает только один аргумент iterable, для нескольких аргументов iterable см. starmap()). Блокирует выполнение, пока результат не будет готов.

Этот метод делит iterable на несколько фрагментов, которые он отправляет в пул процессов как отдельные задачи. Размер этих фрагментов можно приблизительно указать, установив chunksize в положительное целое число.

Обратите внимание, что это может привести к высокому использованию памяти для очень длинных iterable. Рассмотрите использование 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(). Для очень длинных iterable использование большого значения для 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 из iterable и вызывает func с распакованными iterable. Возвращает объект результата.

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

close()

Запрещает дальнейшую отправку задач в пул. После завершения всех задач рабочие процессы завершатся.

terminate()

Немедленно останавливает рабочие процессы, не завершая текущую работу. При сборе мусора объекта пула terminate() будет вызван немедленно.

join()

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

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

class multiprocessing.pool.AsyncResult

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

get([timeout])

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

wait([timeout])

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

ready()

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

successful()

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

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

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

from multiprocessing import Pool
import time

def f(x):
    return x*x

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

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

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

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

Клиенты и прослушиватели

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

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

multiprocessing.connection.deliver_challenge(connection, authkey)

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

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

multiprocessing.connection.answer_challenge(connection, authkey)

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

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

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

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

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

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

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

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

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

Примечание

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

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

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

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

accept()

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

close()

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

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

address

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

last_accepted

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

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

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

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

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

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

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

POSIX: wait(object_list, timeout) почти эквивалентно select.select(object_list, [], [], timeout). Разница заключается в том, что если вызов select.select() прерывается сигналом, он может генерировать исключение OSError с кодом ошибки EINTR, в то время как wait() этого не делает.

Windows: Элемент в object_list должен быть либо целочисленным дескриптором, который можно ожидать (в соответствии с определением, используемым документацией функции Win32 WaitForMultipleObjects()) , либо объектом с методом fileno(), который возвращает дескриптор сокета или канала.

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

Примеры

Следующий код сервера создает прослушиватель, использующий 'secret password' в качестве ключа аутентификации. Затем он ожидает подключения и отправляет некоторые данные клиенту:

from multiprocessing.connection import Listener
from array import array

address = ('localhost', 6000)     # family is deduced to be 'AF_INET'

with Listener(address, authkey=b'secret password') as listener:
    with listener.accept() as conn:
        print('connection accepted from', listener.last_accepted)

        conn.send([2.25, None, 'junk', float])

        conn.send_bytes(b'hello')

        conn.send_bytes(array('i', [42, 1729]))

Следующий код подключается к серверу и получает некоторые данные от сервера:

from multiprocessing.connection import Client
from array import array

address = ('localhost', 6000)

with Client(address, authkey=b'secret password') as conn:
    print(conn.recv())                  # => [2.25, None, 'junk', float]

    print(conn.recv_bytes())            # => 'hello'

    arr = array('i', [0, 0, 0, 0, 0])
    print(conn.recv_bytes_into(arr))    # => 8
    print(arr)                          # => array('i', [42, 1729, 0, 0, 0])

Следующий код использует wait() для одновременного ожидания сообщений от нескольких процессов:

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

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

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

Примечание

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

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

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

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

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

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

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

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

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

Возможность сериализации

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

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

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

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

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

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

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

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

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

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

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

Объединение процессов, использующих очереди

Помните, что процесс, который поместил элементы в очередь, будет ожидать завершения, пока все буферизованные элементы не будут переданы потоком «подачи» в подлежащий канал. (Дочерний процесс может вызвать метод Queue.cancel_join_thread очереди, чтобы избежать этого поведения.)

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

Пример, который приведёт к тупику:

from multiprocessing import Process, Queue

def f(q):
    q.put('X' * 1000000)

if __name__ == '__main__':
    queue = Queue()
    p = Process(target=f, args=(queue,))
    p.start()
    p.join()                    # this deadlocks
    obj = queue.get()

Исправление в данном случае заключалось бы в перестановке последних двух строк (или просто удалении строки p.join()).

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

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

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

Например

from multiprocessing import Process, Lock

def f():
    ... do something using "lock" ...

if __name__ == '__main__':
    lock = Lock()
    for i in range(10):
        Process(target=f).start()

следует переписать как

from multiprocessing import Process, Lock

def f(l):
    ... do something using "l" ...

if __name__ == '__main__':
    lock = Lock()
    for i in range(10):
        Process(target=f, args=(lock,)).start()

Будьте осторожны при замене sys.stdin «объектом файла»

multiprocessing изначально безусловно вызывал:

os.close(sys.stdin.fileno())

в методе multiprocessing.Process._bootstrap() — это приводило к проблемам с процессами в процессах. Это было изменено на:

sys.stdin.close()
sys.stdin = open(os.open(os.devnull, os.O_RDONLY), closefd=False)

Что решает основную проблему столкновения процессов, приводящую к ошибке неверного дескриптора файла, но вводит потенциальную опасность для приложений, которые заменяют sys.stdin() «объектом файла» с буферизацией вывода. Эта опасность заключается в том, что если несколько процессов вызывают close() для этого объекта файла, это может привести к многократному сбросу одних и тех же данных в объект, что приведёт к повреждению.

Если вы пишете объект-файл и реализуете собственное кэширование, вы можете сделать его безопасным для форка, сохраняя идентификатор процесса (PID) при добавлении в кэш и удаляя кэш при изменении PID. Например:

@property
def cache(self):
    pid = os.getpid()
    if pid != self._pid:
        self._pid = pid
        self._cache = []
    return self._cache

Для получения дополнительной информации см. bpo-5155, bpo-5313 и bpo-5331

Методы запуска spawn и forkserver

Существует несколько дополнительных ограничений, которые не применяются к методу запуска fork.

Большая возможность сериализации

Убедитесь, что все аргументы для Process.__init__() могут быть сериализованы. Также, если вы создаёте подкласс Process, убедитесь, что экземпляры будут сериализуемыми, когда вызывается метод Process.start.

Глобальные переменные

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

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

Безопасный импорт главного модуля

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

Например, при использовании метода запуска spawn или forkserver запуск следующего модуля завершится ошибкой RuntimeError:

from multiprocessing import Process

def foo():
    print('hello')

p = Process(target=foo)
p.start()

Вместо этого следует защитить «точку входа» программы, используя if __name__ == '__main__': следующим образом:

from multiprocessing import Process, freeze_support, set_start_method

def foo():
    print('hello')

if __name__ == '__main__':
    freeze_support()
    set_start_method('spawn')
    p = Process(target=foo)
    p.start()

(Строка freeze_support() можно опустить, если программа будет запущена в обычном режиме, а не в замороженном.)

Это позволяет новому запущенному интерпретатору Python безопасно импортировать модуль и затем выполнить функцию foo() модуля.

Аналогичные ограничения применяются, если пул или менеджер создаются в основном модуле.

Примеры

Демонстрация того, как создавать и использовать настраиваемые менеджеры и прокси:

from multiprocessing import freeze_support
from multiprocessing.managers import BaseManager, BaseProxy
import operator

##

class Foo:
    def f(self):
        print('you called Foo.f()')
    def g(self):
        print('you called Foo.g()')
    def _h(self):
        print('you called Foo._h()')

# A simple generator function
def baz():
    for i in range(10):
        yield i*i

# Proxy type for generator objects
class GeneratorProxy(BaseProxy):
    _exposed_ = ['__next__']
    def __iter__(self):
        return self
    def __next__(self):
        return self._callmethod('__next__')

# Function to return the operator module
def get_operator_module():
    return operator

##

class MyManager(BaseManager):
    pass

# register the Foo class; make `f()` and `g()` accessible via proxy
MyManager.register('Foo1', Foo)

# register the Foo class; make `g()` and `_h()` accessible via proxy
MyManager.register('Foo2', Foo, exposed=('g', '_h'))

# register the generator function baz; use `GeneratorProxy` to make proxies
MyManager.register('baz', baz, proxytype=GeneratorProxy)

# register get_operator_module(); make public functions accessible via proxy
MyManager.register('operator', get_operator_module)

##

def test():
    manager = MyManager()
    manager.start()

    print('-' * 20)

    f1 = manager.Foo1()
    f1.f()
    f1.g()
    assert not hasattr(f1, '_h')
    assert sorted(f1._exposed_) == sorted(['f', 'g'])

    print('-' * 20)

    f2 = manager.Foo2()
    f2.g()
    f2._h()
    assert not hasattr(f2, 'f')
    assert sorted(f2._exposed_) == sorted(['g', '_h'])

    print('-' * 20)

    it = manager.baz()
    for i in it:
        print('<%d>' % i, end=' ')
    print()

    print('-' * 20)

    op = manager.operator()
    print('op.add(23, 45) =', op.add(23, 45))
    print('op.pow(2, 94) =', op.pow(2, 94))
    print('op._exposed_ =', op._exposed_)

##

if __name__ == '__main__':
    freeze_support()
    test()

Использование Pool:

import multiprocessing
import time
import random
import sys

#
# Functions used by test code
#

def calculate(func, args):
    result = func(*args)
    return '%s says that %s%s = %s' % (
        multiprocessing.current_process().name,
        func.__name__, args, result
        )

def calculatestar(args):
    return calculate(*args)

def mul(a, b):
    time.sleep(0.5 * random.random())
    return a * b

def plus(a, b):
    time.sleep(0.5 * random.random())
    return a + b

def f(x):
    return 1.0 / (x - 5.0)

def pow3(x):
    return x ** 3

def noop(x):
    pass

#
# Test code
#

def test():
    PROCESSES = 4
    print('Creating pool with %d processes\n' % PROCESSES)

    with multiprocessing.Pool(PROCESSES) as pool:
        #
        # Tests
        #

        TASKS = [(mul, (i, 7)) for i in range(10)] + \
                [(plus, (i, 8)) for i in range(10)]

        results = [pool.apply_async(calculate, t) for t in TASKS]
        imap_it = pool.imap(calculatestar, TASKS)
        imap_unordered_it = pool.imap_unordered(calculatestar, TASKS)

        print('Ordered results using pool.apply_async():')
        for r in results:
            print('\t', r.get())
        print()

        print('Ordered results using pool.imap():')
        for x in imap_it:
            print('\t', x)
        print()

        print('Unordered results using pool.imap_unordered():')
        for x in imap_unordered_it:
            print('\t', x)
        print()

        print('Ordered results using pool.map() --- will block till complete:')
        for x in pool.map(calculatestar, TASKS):
            print('\t', x)
        print()

        #
        # Test error handling
        #

        print('Testing error handling:')

        try:
            print(pool.apply(f, (5,)))
        except ZeroDivisionError:
            print('\tGot ZeroDivisionError as expected from pool.apply()')
        else:
            raise AssertionError('expected ZeroDivisionError')

        try:
            print(pool.map(f, list(range(10))))
        except ZeroDivisionError:
            print('\tGot ZeroDivisionError as expected from pool.map()')
        else:
            raise AssertionError('expected ZeroDivisionError')

        try:
            print(list(pool.imap(f, list(range(10)))))
        except ZeroDivisionError:
            print('\tGot ZeroDivisionError as expected from list(pool.imap())')
        else:
            raise AssertionError('expected ZeroDivisionError')

        it = pool.imap(f, list(range(10)))
        for i in range(10):
            try:
                x = next(it)
            except ZeroDivisionError:
                if i == 5:
                    pass
            except StopIteration:
                break
            else:
                if i == 5:
                    raise AssertionError('expected ZeroDivisionError')

        assert i == 9
        print('\tGot ZeroDivisionError as expected from IMapIterator.next()')
        print()

        #
        # Testing timeouts
        #

        print('Testing ApplyResult.get() with timeout:', end=' ')
        res = pool.apply_async(calculate, TASKS[0])
        while 1:
            sys.stdout.flush()
            try:
                sys.stdout.write('\n\t%s' % res.get(0.02))
                break
            except multiprocessing.TimeoutError:
                sys.stdout.write('.')
        print()
        print()

        print('Testing IMapIterator.next() with timeout:', end=' ')
        it = pool.imap(calculatestar, TASKS)
        while 1:
            sys.stdout.flush()
            try:
                sys.stdout.write('\n\t%s' % it.next(0.02))
            except StopIteration:
                break
            except multiprocessing.TimeoutError:
                sys.stdout.write('.')
        print()
        print()


if __name__ == '__main__':
    multiprocessing.freeze_support()
    test()

Пример, демонстрирующий использование очередей для подачи задач в набор рабочих процессов и сбора результатов:

import time
import random

from multiprocessing import Process, Queue, current_process, freeze_support

#
# Function run by worker processes
#

def worker(input, output):
    for func, args in iter(input.get, 'STOP'):
        result = calculate(func, args)
        output.put(result)

#
# Function used to calculate result
#

def calculate(func, args):
    result = func(*args)
    return '%s says that %s%s = %s' % \
        (current_process().name, func.__name__, args, result)

#
# Functions referenced by tasks
#

def mul(a, b):
    time.sleep(0.5*random.random())
    return a * b

def plus(a, b):
    time.sleep(0.5*random.random())
    return a + b

#
#
#

def test():
    NUMBER_OF_PROCESSES = 4
    TASKS1 = [(mul, (i, 7)) for i in range(20)]
    TASKS2 = [(plus, (i, 8)) for i in range(10)]

    # Create queues
    task_queue = Queue()
    done_queue = Queue()

    # Submit tasks
    for task in TASKS1:
        task_queue.put(task)

    # Start worker processes
    for i in range(NUMBER_OF_PROCESSES):
        Process(target=worker, args=(task_queue, done_queue)).start()

    # Get and print results
    print('Unordered results:')
    for i in range(len(TASKS1)):
        print('\t', done_queue.get())

    # Add more tasks using `put()`
    for task in TASKS2:
        task_queue.put(task)

    # Get and print some more results
    for i in range(len(TASKS2)):
        print('\t', done_queue.get())

    # Tell child processes to stop
    for i in range(NUMBER_OF_PROCESSES):
        task_queue.put('STOP')


if __name__ == '__main__':
    freeze_support()
    test()

© 2001–2024 Python Software Foundation
Licensed under the PSF License.
https://docs.python.org/3.13/library/multiprocessing.html

Spec-Zone.ru

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