Spec-Zone.ru › Python 3.10

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

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

Введение

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

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

from multiprocessing import Pool

def f(x):
    return x*x

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

будет выводить на стандартный вывод

[1, 4, 9]

См. также

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

Класс Process

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

from multiprocessing import Process

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

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

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

from multiprocessing import Process
import os

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

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

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

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

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

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

spawn

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

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

fork

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

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

forkserver

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

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

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

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

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

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

import multiprocessing as mp

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

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

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

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

import multiprocessing as mp

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

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

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

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

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

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

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

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

Очереди

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

from multiprocessing import Process, Queue

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

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

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

Трубы

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

from multiprocessing import Process, Pipe

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

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

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

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

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

from multiprocessing import Process, Lock

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

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

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

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

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

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

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

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

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

from multiprocessing import Process, Value, Array

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

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

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

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

выведет

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

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

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

Процесс-сервер

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

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

from multiprocessing import Process, Manager

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

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

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

        print(d)
        print(l)

выведет

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

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

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

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

Например:

from multiprocessing import Pool, TimeoutError
import time
import os

def f(x):
    return x*x

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

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

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

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

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

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

        # make a single worker sleep for 10 secs
        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):
AttributeError: 'module' object has no attribute 'f'
AttributeError: 'module' object has no attribute 'f'
AttributeError: 'module' object has no attribute 'f'

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

Справочник

Пакет 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 не передаются аргументы.

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

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

run()

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

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

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

Новое в версии 3.3.

terminate()

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

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

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

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

kill()

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

Новое в версии 3.7.

close()

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

Новое в версии 3.7.

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

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

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

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

exception multiprocessing.BufferTooShort

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

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

exception multiprocessing.AuthenticationError

Возбуждается при ошибке аутентификации.

exception multiprocessing.TimeoutError

Возбуждается методами с таймаутом при истечении таймаута.

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

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

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

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

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

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

Примечание

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

Примечание

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

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

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

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

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

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

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

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

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

multiprocessing.Pipe([duplex])

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

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

class multiprocessing.Queue([maxsize])

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

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

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

qsize()

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

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

empty()

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

full()

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

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

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

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

put_nowait(obj)

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

get([block[, timeout]])

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

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

get_nowait()

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

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

close()

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

join_thread()

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

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

cancel_join_thread()

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

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

Примечание

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

class multiprocessing.SimpleQueue

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

close()

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

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

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

empty()

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

get()

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

put(item)

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

class multiprocessing.JoinableQueue([maxsize])

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

task_done()

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

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

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

join()

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

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

Разное

multiprocessing.active_children()

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

Вызов этого метода имеет побочный эффект «присоединения» к завершенным процессам.

multiprocessing.cpu_count()

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

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

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

См. также

os.cpu_count()

multiprocessing.current_process()

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

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

multiprocessing.parent_process()

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

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

multiprocessing.freeze_support()

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

Эта функция должна вызываться сразу после строки if __name__ == '__main__' в главном модуле. Например:

from multiprocessing import Process, freeze_support

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

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

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

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

multiprocessing.get_all_start_methods()

Возвращает список поддерживаемых методов запуска, первый из которых является стандартным. Возможные методы запуска: 'fork', 'spawn' и 'forkserver'. В Windows доступен только 'spawn'. В Unix 'fork' и 'spawn' всегда поддерживаются, 'fork' является стандартным.

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

multiprocessing.get_context(method=None)

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

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

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

multiprocessing.get_start_method(allow_none=False)

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

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

Возвращаемое значение может быть 'fork', 'spawn', 'forkserver' или None. 'fork' является стандартным в Unix, а 'spawn' — в Windows и macOS.

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

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

multiprocessing.set_executable(executable)

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

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

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

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

multiprocessing.set_start_method(method)

Устанавливает метод, который должен использоваться для запуска дочерних процессов. method может быть 'fork', 'spawn' или 'forkserver'.

Обратите внимание, что это должно вызываться не более одного раза, и это должно быть защищено внутри 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]])

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

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

recv_bytes([maxlength])

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

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

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

recv_bytes_into(buffer[, offset])

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

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

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

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

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

Например:

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

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

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

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

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

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

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

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

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

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

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

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

class multiprocessing.BoundedSemaphore([value])

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

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

Примечание

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

class multiprocessing.Condition([lock])

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

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

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

class multiprocessing.Event

Копия threading.Event.

class multiprocessing.Lock

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

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

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

acquire(block=True, timeout=None)

Получить блокировку, блокирующим или неблокирующим способом.

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

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

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

release()

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

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

class multiprocessing.RLock

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

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

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

acquire(block=True, timeout=None)

Получить блокировку, блокирующим или неблокирующим способом.

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

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

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

release()

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

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

class multiprocessing.Semaphore([value])

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

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

Примечание

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

Примечание

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

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

Примечание

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

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

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

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

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

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

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

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

counter.value += 1

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

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

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

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

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

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

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

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

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

Модуль multiprocessing.sharedctypes

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

Примечание

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

multiprocessing.sharedctypes.RawArray(typecode_or_type, size_or_initializer)

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

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

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

multiprocessing.sharedctypes.RawValue(typecode_or_type, *args)

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

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

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

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

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

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

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

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

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

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

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

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

multiprocessing.sharedctypes.copy(obj)

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

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

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

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

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

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

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

ctypes

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

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

c_double(2.4)

RawValue(c_double, 2.4)

RawValue(‘d’, 2.4)

MyStruct(4, 6)

RawValue(MyStruct, 4, 6)

(c_short * 7)()

RawArray(c_short, 7)

RawArray(‘h’, 7)

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

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

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

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

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

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

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

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

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

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

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

Выводимые результаты:

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

Менеджеры

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

multiprocessing.Manager()

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

Процессы менеджеров будут закрыты, как только они будут удалены сборщиком мусора или родительский процесс завершит работу. Классы менеджеров определены в модуле multiprocessing.managers:

class multiprocessing.managers.BaseManager([address[, authkey]])

Создает объект BaseManager.

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

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

authkey — ключ аутентификации, который будет использоваться для проверки подлинности входящих подключений к серверному процессу. Если authkey — None , то используется current_process().authkey. В противном случае используется authkey, и он должен быть строкой байтов.

start([initializer[, initargs]])

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

get_server()

Возвращает объект Server , представляющий фактический сервер под управлением менеджера. Объект Server поддерживает метод serve_forever():

>>> from multiprocessing.managers import BaseManager
>>> manager = BaseManager(address=('', 50000), authkey=b'abc')
>>> server = manager.get_server()
>>> server.serve_forever()

Server дополнительно имеет атрибут address.

connect()

Подключить локальный объект менеджера к удаленному процессу менеджера:

>>> from multiprocessing.managers import BaseManager
>>> m = BaseManager(address=('127.0.0.1', 50000), authkey=b'abc')
>>> m.connect()
shutdown()

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

Это можно вызывать несколько раз.

register(typeid[, callable[, proxytype[, exposed[, method_to_typeid[, create_method]]]]])

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

typeid — «идентификатор типа», используемый для идентификации конкретного типа общего объекта. Это должна быть строка.

callable — вызываемый объект, используемый для создания объектов для этого идентификатора типа. Если экземпляр менеджера будет подключен к серверу с помощью метода connect(), или если аргумент create_method — False , то это можно оставить как None.

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

exposed используется для указания последовательности имён методов, которые прокси для этого typeid должны иметь право доступа с помощью BaseProxy._callmethod(). (Если exposed — None , то proxytype._exposed_ используется вместо него, если он существует.) В случае, если список exposed не указан, все «публичные методы» общего объекта будут доступны. (Здесь «публичный метод» означает любой атрибут, у которого есть метод __call__() и имя которого не начинается с '_'.)

method_to_typeid — отображение, используемое для указания типа возвращаемого значения тех методов с раскрытием, которые должны возвращать прокси. Оно сопоставляет имена методов строкам typeid. (Если method_to_typeid — None , то proxytype._method_to_typeid_ используется вместо него, если он существует.) Если имя метода не является ключом в этом отображении или если отображение — None , то объект, возвращаемый методом, будет скопирован по значению.

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

BaseManager экземпляры также имеют одно свойство только для чтения:

address

Адрес, используемый менеджером.

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

В предыдущих версиях __enter__() не запускал серверный процесс менеджера, если он еще не был запущен.

class multiprocessing.managers.SyncManager

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

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

Barrier(parties[, action[, timeout]])

Создает общий объект threading.Barrier и возвращает прокси для него.

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

BoundedSemaphore([value])

Создает общий объект threading.BoundedSemaphore и возвращает прокси для него.

Condition([lock])

Создает общий объект threading.Condition и возвращает прокси для него.

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

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

Event()

Создает общий объект threading.Event и возвращает прокси для него.

Lock()

Создает общий объект threading.Lock и возвращает прокси для него.

Namespace()

Создает общий объект Namespace и возвращает прокси для него.

Queue([maxsize])

Создает общий объект queue.Queue и возвращает прокси для него.

RLock()

Создает общий объект threading.RLock и возвращает прокси для него.

Semaphore([value])

Создает общий объект threading.Semaphore и возвращает прокси для него.

Array(typecode, sequence)

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

Value(typecode, value)

Создает объект с изменяемым атрибутом value и возвращает прокси для него.

dict()
dict(mapping)
dict(sequence)

Создает общий объект dict и возвращает прокси для него.

list()
list(sequence)

Создает общий объект list и возвращает прокси для него.

Изменено в версии 3.6: Общие объекты могут быть вложены. Например, общий контейнер, такой как общий список, может содержать другие общие объекты, которые будут управляться и синхронизироваться менеджером SyncManager.

class multiprocessing.managers.Namespace

Тип, который может регистрироваться в SyncManager.

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

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

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

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

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

from multiprocessing.managers import BaseManager

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

class MyManager(BaseManager):
    pass

MyManager.register('Maths', MathsClass)

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Примечание

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

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

Вместо этого для сравнения следует использовать копию объекта-референта.

class multiprocessing.managers.BaseProxy

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

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

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

Если proxy — это прокси, чьим объектом-референтом является obj , то выражение

proxy._callmethod(methodname, args, kwds)

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

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

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

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

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

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

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

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

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

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

__repr__()

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

__str__()

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

Очистка

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Добавлена в версии 3.2: maxtasksperchild

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

Примечание

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

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

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

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

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

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

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

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

map(func, iterable[, chunksize])

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

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

Обратите внимание, что это может привести к высокому потреблению памяти для очень длинных итерируемых объектов. Рассмотрите использование imap() или imap_unordered() с явным параметром chunksize для повышения эффективности.

map_async(func, iterable[, chunksize[, callback[, error_callback]]])

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

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

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

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

imap(func, iterable[, chunksize])

Более ленивая версия map().

Аргумент chunksize такой же, как и в методе map(). Для очень длинных итерируемых объектов использование большого значения для chunksize может значительно ускорить выполнение задания по сравнению со значением по умолчанию 1.

Также, если chunksize равно 1 , метод next() итератора, возвращаемого методом imap(), имеет необязательный параметр timeout: next(timeout) будет поднимать multiprocessing.TimeoutError, если результат не может быть возвращён в течение timeout секунд.

imap_unordered(func, iterable[, chunksize])

То же самое, что и imap(), за исключением того, что порядок результатов из возвращаемого итератора должен считаться произвольным. (Только когда есть только один рабочий процесс, порядок гарантируется «правильным».)

starmap(func, iterable[, chunksize])

Подобно map(), за исключением того, что элементы iterable ожидаются как итерируемые объекты, которые распаковываются в качестве аргументов.

Таким образом, iterable из [(1,2), (3, 4)] приводит к [func(1,2), func(3,4)].

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

starmap_async(func, iterable[, chunksize[, callback[, error_callback]]])

Комбинация starmap() и map_async(), которая итерируется по iterable из итерируемых объектов и вызывает func с распакованными итерируемыми объектами. Возвращает объект результата.

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

close()

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

terminate()

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

join()

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

New in version 3.3: Objects of the Pool class now support the context management protocol – see Типы менеджеров контекста. __enter__() возвращает объект пула, а __exit__() вызывает terminate().

class multiprocessing.pool.AsyncResult

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

get([timeout])

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

wait([timeout])

Ожидает доступности результата или истечения timeout секунд.

ready()

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

successful()

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

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

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

from multiprocessing import Pool
import time

def f(x):
    return x*x

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

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

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

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

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

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

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

multiprocessing.connection.deliver_challenge(connection, authkey)

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

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

multiprocessing.connection.answer_challenge(connection, authkey)

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

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

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

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

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

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

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

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

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

Примечание

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

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

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

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

accept()

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

close()

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

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

address

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

last_accepted

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

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

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

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

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

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

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

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

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

Новое в версии 3.3.

Примеры

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

from multiprocessing.connection import Listener
from array import array

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

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

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

        conn.send_bytes(b'hello')

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

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

from multiprocessing.connection import Client
from array import array

address = ('localhost', 6000)

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

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

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

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

import time, random
from multiprocessing import Process, Pipe, current_process
from multiprocessing.connection import wait

def foo(w):
    for i in range(10):
        w.send((i, current_process().name))
    w.close()

if __name__ == '__main__':
    readers = []

    for i in range(4):
        r, w = Pipe(duplex=False)
        readers.append(r)
        p = Process(target=foo, args=(w,))
        p.start()
        # We close the writable end of the pipe now to be sure that
        # p is the only process which owns a handle for it.  This
        # ensures that when p closes its handle for the writable end,
        # wait() will promptly report the readable end as being ready.
        w.close()

    while readers:
        for r in wait(readers):
            try:
                msg = r.recv()
            except EOFError:
                readers.remove(r)
            else:
                print(msg)

Форматы адресов

  • Адрес 'AF_INET' представляет собой кортеж вида (hostname, port), где hostname — строка, а port — целое число.
  • Адрес 'AF_UNIX' — это строка, представляющая имя файла в файловой системе.
  • Адрес 'AF_PIPE' — это строка вида r'\\.\pipe\PipeName'. Для подключения к именованной pipe на удалённом компьютере под названием ServerName следует использовать адрес вида r'\\ServerName\pipe\PipeName'.

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

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

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

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

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

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

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

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

multiprocessing.get_logger()

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

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

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

multiprocessing.log_to_stderr(level=None)

Эта функция выполняет вызов get_logger(), но дополнительно добавляет обработчик, который отправляет вывод в sys.stderr с форматом '[%(levelname)s/%(processName)s] %(message)s'. Вы можете изменить levelname логировщика, передав аргумент level.

Ниже приведен пример сеанса с включённым ведением журнала:

>>> import multiprocessing, logging
>>> logger = multiprocessing.log_to_stderr()
>>> logger.setLevel(logging.INFO)
>>> logger.warning('doomed')
[WARNING/MainProcess] doomed
>>> m = multiprocessing.Manager()
[INFO/SyncManager-...] child process calling self.run()
[INFO/SyncManager-...] created temp directory /.../pymp-...
[INFO/SyncManager-...] manager serving at '/.../listener-...'
>>> del m
[INFO/MainProcess] sending shutdown message to manager
[INFO/SyncManager-...] manager exiting with exitcode 0

Полную таблицу уровней ведения журнала см. в модуле logging.

Модуль multiprocessing.dummy

Модуль multiprocessing.dummy дублирует API модуля multiprocessing, но представляет собой не более чем обёртку вокруг модуля threading.

В частности, функция Pool модуля multiprocessing.dummy возвращает экземпляр ThreadPool, который является подклассом Pool и поддерживает все те же вызовы методов, но использует пул потоков-рабочих, а не процессов-рабочих.

class multiprocessing.pool.ThreadPool([processes[, initializer[, initargs]]])

Объект пула потоков, который управляет пулом потоков-рабочих, которым можно передавать задания. Экземпляры ThreadPool полностью совместимы по интерфейсу с экземплярами Pool, и их ресурсы также должны быть должным образом управляемы, либо используя пул как менеджер контекста, либо вызывая close() и terminate() вручную.

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

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

В отличие от Pool, maxtasksperchild и context нельзя указать.

Примечание

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

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

Рекомендации по программированию

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

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

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

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

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

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

Поддержка пиклирования

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

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

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

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

Обработка завершённых процессов-зомби

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

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

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

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

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

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

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

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

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

Пример, который может привести к тупику:

from multiprocessing import Process, Queue

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

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

Решение в этом случае — поменять местами две последние строки (или просто удалить строку p.join()).

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

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

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

Например

from multiprocessing import Process, Lock

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

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

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

from multiprocessing import Process, Lock

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

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

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

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

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

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

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

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

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

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

Дополнительную информацию можно найти в bpo-5155, bpo-5313 и bpo-5331

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

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

Более строгие требования к пиклированию

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

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

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

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

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

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

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

from multiprocessing import Process

def foo():
    print('hello')

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

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

from multiprocessing import Process, freeze_support, set_start_method

def foo():
    print('hello')

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

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

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

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

Примеры

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

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

##

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

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

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

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

##

class MyManager(BaseManager):
    pass

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

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

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

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

##

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

    print('-' * 20)

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

    print('-' * 20)

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

    print('-' * 20)

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

    print('-' * 20)

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

##

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

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

import multiprocessing
import time
import random
import sys

#
# Functions used by test code
#

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

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

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

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

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

def pow3(x):
    return x ** 3

def noop(x):
    pass

#
# Test code
#

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

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

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

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

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

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

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

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

        #
        # Test error handling
        #

        print('Testing error handling:')

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

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

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

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

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

        #
        # Testing timeouts
        #

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

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


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

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

import time
import random

from multiprocessing import Process, Queue, current_process, freeze_support

#
# Function run by worker processes
#

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

#
# Function used to calculate result
#

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

#
# Functions referenced by tasks
#

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

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

#
#
#

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

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

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

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

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

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

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

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


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

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

Spec-Zone.ru

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