Spec-Zone.ru › Python 3.7

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

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

Введение

multiprocessing — пакет, поддерживающий создание процессов с помощью API, похожим на 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]

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

fork

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

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

forkserver

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

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

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

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

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

По умолчанию целевому объекту 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 , если процесс еще не завершился. Отрицательное значение -N указывает, что дочерний процесс был завершен сигналом N.

authkey

Ключ аутентификации процесса (строка байтов).

При инициализации multiprocessing главному процессу присваивается случайная строка с использованием os.urandom().

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

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

sentinel

Числовой дескриптор системного объекта, который станет «готовым» при завершении процесса.

Вы можете использовать это значение, если хотите дождаться сразу нескольких событий, используя multiprocessing.connection.wait(). В противном случае вызов join() проще.

В Windows это системный дескриптор, используемый с семейством API-функций WaitForSingleObject и WaitForMultipleObjects. В Unix это дескриптор файла, используемый с примитивами из модуля select.

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

terminate()

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

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

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

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

kill()

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

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

close()

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

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

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

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

>>> import multiprocessing, time, signal
>>> p = multiprocessing.Process(target=time.sleep, args=(1000,))
>>> print(p, p.is_alive())
<Process(Process-1, initial)> False
>>> p.start()
>>> print(p, p.is_alive())
<Process(Process-1, started)> True
>>> p.terminate()
>>> time.sleep(0.1)
>>> print(p, p.is_alive())
<Process(Process-1, stopped[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, таких как Mac OS X, где sem_getvalue() не реализован.

empty()

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

full()

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

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

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

put_nowait(obj)

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

get([block[, timeout]])

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

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.

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

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

multiprocessing.set_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]])

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

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

recv_bytes([maxlength])

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

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

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

recv_bytes_into(buffer[, offset])

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

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

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

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

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

Примечание

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

class multiprocessing.Condition([lock])

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

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

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

class multiprocessing.Event

Клон threading.Event.

class multiprocessing.Lock

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

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

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

acquire(block=True, timeout=None)

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

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

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

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

release()

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

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

class multiprocessing.RLock

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

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

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

acquire(block=True, timeout=None)

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

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

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

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

release()

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

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

class multiprocessing.Semaphore([value])

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

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

Примечание

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

Примечание

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

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

Примечание

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

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

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

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

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

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

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

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

counter.value += 1

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

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

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

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

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

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

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

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

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

Модуль multiprocessing.sharedctypes

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

Примечание

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

multiprocessing.sharedctypes.RawArray(typecode_or_type, size_or_initializer)

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

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

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

multiprocessing.sharedctypes.RawValue(typecode_or_type, *args)

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

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

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

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

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

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

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

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

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

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

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

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

multiprocessing.sharedctypes.copy(obj)

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

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

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

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

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

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

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

ctypes

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

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

c_double(2.4)

RawValue(c_double, 2.4)

RawValue(‘d’, 2.4)

MyStruct(4, 6)

RawValue(MyStruct, 4, 6)

(c_short * 7)()

RawArray(c_short, 7)

RawArray(‘h’, 7)

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

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

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

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

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

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

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

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

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

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

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

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

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

Менеджеры

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

multiprocessing.Manager()

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

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

class multiprocessing.managers.BaseManager([address[, 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__() не запускал серверный процесс менеджера, если он не был уже запущен.

END_OF_DOCUMENT_MARKER
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(Worker, self).__init__()
...     def run(self):
...         self.q.put('local hello')
...
>>> queue = Queue()
>>> w = Worker(queue)
>>> w.start()
>>> class QueueManager(BaseManager): pass
...
>>> QueueManager.register('get_queue', callable=lambda: queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()

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

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

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

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

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

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

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

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

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

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

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

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

Примечание

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

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

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

class multiprocessing.managers.BaseProxy

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

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

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

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

proxy._callmethod(methodname, args, kwds)

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

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

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

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

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

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

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

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

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

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

__repr__()

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

__str__()

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

Очистка

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Новое в версии 3.2: maxtasksperchild

Новое в версии 3.4: context

Примечание

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

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

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

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

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

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

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

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

map(func, iterable[, chunksize])

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

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

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

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

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

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

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

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

imap(func, iterable[, chunksize])

Более ленивый вариант map().

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

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

imap_unordered(func, iterable[, chunksize])

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

starmap(func, iterable[, chunksize])

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

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

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

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

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

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

close()

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

terminate()

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

join()

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

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

class multiprocessing.pool.AsyncResult

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

get([timeout])

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

wait([timeout])

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

ready()

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

successful()

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

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

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' представляет собой строку, представляющую имя файла в файловой системе.
  • An 'AF_PIPE' address is a string of the form

    r'\.\pipe{PipeName}'. Чтобы подключиться к именованной трубе на удалённом компьютере под названием ServerName с помощью Client(), необходимо использовать адрес вида 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()

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

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

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

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

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

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

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

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

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

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

Пригодность для сериализации

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Следующий пример приведёт к тупику:

from multiprocessing import Process, Queue

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

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

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

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

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

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

Например

from multiprocessing import Process, Lock

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

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

должно быть переписано как

sys.stdin

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

from multiprocessing import Process

def foo():
    print('hello')

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

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

from multiprocessing import Process, freeze_support, set_start_method

def foo():
    print('hello')

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

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

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

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

Примеры

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

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

##

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

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

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

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

##

class MyManager(BaseManager):
    pass

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

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

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

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

##

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

    print('-' * 20)

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

    print('-' * 20)

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

    print('-' * 20)

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

    print('-' * 20)

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

##

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

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

import multiprocessing
import time
import random
import sys

#
# Functions used by test code
#

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

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

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

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

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

def pow3(x):
    return x ** 3

def noop(x):
    pass

#
# Test code
#

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

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

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

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

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

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

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

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

        #
        # Test error handling
        #

        print('Testing error handling:')

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

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

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

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

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

        #
        # Testing timeouts
        #

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

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


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

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

import time
import random

from multiprocessing import Process, Queue, current_process, freeze_support

#
# Function run by worker processes
#

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

#
# Function used to calculate result
#

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

#
# Functions referenced by tasks
#

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

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

#
#
#

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

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

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

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

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

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

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

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


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

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

Spec-Zone.ru

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