Spec-Zone.ru › Python 3.8

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

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

Введение

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

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

from multiprocessing import Pool

def f(x):
    return x*x

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

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

[1, 4, 9]

Класс Process

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

from multiprocessing import Process

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

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

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

from multiprocessing import Process
import os

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

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

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

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

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

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

spawn

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

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

fork

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

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

forkserver

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

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

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

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

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

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

import multiprocessing as mp

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

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

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

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

import multiprocessing as mp

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

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

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

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

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

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

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

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

Очереди

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

from multiprocessing import Process, Queue

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

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

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

Трубы

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

from multiprocessing import Process, Pipe

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

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

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

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

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

from multiprocessing import Process, Lock

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

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

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

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

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

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

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

Общая память

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

from multiprocessing import Process, Value, Array

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

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

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

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

выведет

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

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

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

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

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

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

from multiprocessing import Process, Manager

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

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

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

        print(d)
        print(l)

выведет

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

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

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

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

Например:

from multiprocessing import Pool, TimeoutError
import time
import os

def f(x):
    return x*x

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

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

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

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

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

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

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

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

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

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

Примечание

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

>>> from multiprocessing import Pool
>>> p = Pool(5)
>>> def f(x):
...     return x*x
...
>>> with p:
...   p.map(f, [1,2,3])
Process PoolWorker-1:
Process PoolWorker-2:
Process PoolWorker-3:
Traceback (most recent call last):
AttributeError: 'module' object has no attribute 'f'
AttributeError: 'module' object has no attribute 'f'
AttributeError: 'module' object has no attribute 'f'

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

Справочник

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

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

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

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

Конструктор всегда должен вызываться с именованными аргументами. Аргумент group всегда должен быть None; он существует только для совместимости с threading.Thread. target — вызываемый объект, который будет вызван методом run(). По умолчанию он равен None, что означает, что ничего не вызывается. name — имя процесса (см. name для получения дополнительной информации). args — кортеж аргументов для вызова target. kwargs — словарь именованных аргументов для вызова target. Если указан, ключевой аргумент 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 ... initial> False
 >>> p.start()
 >>> print(p, p.is_alive())
 <Process ... started> True
 >>> p.terminate()
 >>> time.sleep(0.1)
 >>> print(p, p.is_alive())
 <Process ... stopped exitcode=-SIGTERM> False
 >>> p.exitcode == -signal.SIGTERM
 True
exception multiprocessing.ProcessError

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

exception multiprocessing.BufferTooShort

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

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

exception multiprocessing.AuthenticationError

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

exception multiprocessing.TimeoutError

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

Трубы и очереди

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

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

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

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

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

Примечание

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

Примечание

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

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

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

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

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

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

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

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

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

multiprocessing.Pipe([duplex])

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

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

class multiprocessing.Queue([maxsize])

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

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

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

qsize()

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

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

empty()

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

full()

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

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

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

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

put_nowait(obj)

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

get([block[, timeout]])

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

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

get_nowait()

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

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

close()

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

join_thread()

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

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

cancel_join_thread()

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

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

Примечание

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

class multiprocessing.SimpleQueue

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

empty()

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

get()

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

put(item)

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

class multiprocessing.JoinableQueue([maxsize])

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

task_done()

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

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

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

join()

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

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

Разное

multiprocessing.active_children()

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

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

multiprocessing.cpu_count()

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

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

Может вызвать исключение NotImplementedError.

См. также

os.cpu_count()

multiprocessing.current_process()

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

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

multiprocessing.parent_process()

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

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

multiprocessing.freeze_support()

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

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

from multiprocessing import Process, freeze_support

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

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

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

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

multiprocessing.get_all_start_methods()

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

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

multiprocessing.get_context(method=None)

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

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

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

multiprocessing.get_start_method(allow_none=False)

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

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

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

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

Введено в версии 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 указан, столько байтов будут считаны из buffer. Очень большие буферы (приблизительно 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: Объекты соединения теперь поддерживают протокол управления контекстом – см. Типы менеджеров контекста. __enter__() возвращает объект соединения, а __exit__() вызывает close().

Например:

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

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

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

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

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

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

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

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

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

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

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

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

class multiprocessing.BoundedSemaphore([value])

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

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

Примечание

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

class multiprocessing.Condition([lock])

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

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

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

class multiprocessing.Event

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

class multiprocessing.Lock

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

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

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

acquire(block=True, timeout=None)

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

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

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

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

release()

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

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

class multiprocessing.RLock

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

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

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

acquire(block=True, timeout=None)

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

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

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

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

release()

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

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

class multiprocessing.Semaphore([value])

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

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

END_OF_DOCUMENT_MARKER

Примечание

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

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

ctypes

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

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

c_double(2.4)

RawValue(c_double, 2.4)

RawValue(‘d’, 2.4)

MyStruct(4, 6)

RawValue(MyStruct, 4, 6)

(c_short * 7)()

RawArray(c_short, 7)

RawArray(‘h’, 7)

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

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

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

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

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

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

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

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

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

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

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

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

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

Менеджеры

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

multiprocessing.Manager()

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

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

class multiprocessing.managers.BaseManager([address[, 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 — отображение, используемое для указания типа возвращаемого значения тех методов exposed, которые должны возвращать прокси. Оно сопоставляет имена методов с строками typeid. (Если method_to_typeid None , то вместо этого используется proxytype._method_to_typeid_ , если оно существует.) Если имя метода не является ключом этого отображения или отображение None , то возвращаемый объектом метод будет скопирован по значению.

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

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

address

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

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

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

class multiprocessing.managers.SyncManager

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

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

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

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

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

BoundedSemaphore([value])

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

Condition([lock])

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

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

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

Event()

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

Lock()

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

Namespace()

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

Queue([maxsize])

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

RLock()

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

Semaphore([value])

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

Array(typecode, sequence)

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

Value(typecode, value)

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

dict()
dict(mapping)
dict(sequence)

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

list()
list(sequence)

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

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

class multiprocessing.managers.Namespace

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

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

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

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

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

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

from multiprocessing.managers import BaseManager

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

class MyManager(BaseManager):
    pass

MyManager.register('Maths', MathsClass)

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Примечание

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

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

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

class multiprocessing.managers.BaseProxy

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

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

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

Если proxy — это прокси, чья ссылочная сущность — obj , то выражение

proxy._callmethod(methodname, args, kwds)

оценит выражение

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

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

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

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

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

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

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

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

Если ссылочная сущность не пикл-сериализуема, это вызовет исключение.

__repr__()

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

__str__()

Возвращает представление ссылочной сущности.

Очистка

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Введено в версии 3.2: maxtasksperchild

Введено в версии 3.4: context

Примечание

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

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

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

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

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

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

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

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

map(func, iterable[, chunksize])

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

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

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

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

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

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

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

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

imap(func, iterable[, chunksize])

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

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

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

imap_unordered(func, iterable[, chunksize])

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

starmap(func, iterable[, chunksize])

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

Следовательно, iterable из [(1,2), (3, 4)] элементов приводит к [func(1,2), func(3,4)].

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

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

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

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

close()

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

terminate()

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

join()

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

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

class multiprocessing.pool.AsyncResult

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

get([timeout])

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

wait([timeout])

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

ready()

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

successful()

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

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

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

from multiprocessing import Pool
import time

def f(x):
    return x*x

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

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

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

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

Серверы и клиенты

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

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

multiprocessing.connection.deliver_challenge(connection, authkey)

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

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

multiprocessing.connection.answer_challenge(connection, authkey)

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

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

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

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

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

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

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

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

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

Примечание

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

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

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

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

accept()

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

close()

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

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

address

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

last_accepted

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

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

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

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

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

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

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

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

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

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

Примеры

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

from multiprocessing.connection import Listener
from array import array

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

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

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

        conn.send_bytes(b'hello')

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

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

from multiprocessing.connection import Client
from array import array

address = ('localhost', 6000)

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

multiprocessing.get_logger()

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

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

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

multiprocessing.log_to_stderr()

Эта функция выполняет вызов get_logger(), но помимо возврата логгера, созданного функцией 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.

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

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

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

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

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

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

Примечание

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

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

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

При использовании multiprocessing следует соблюдать определенные рекомендации и подходы.

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

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

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

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

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

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

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

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

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

(Никаких проблем нет, если разные процессы используют один и тот же прокси.)

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

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

Наследование лучше, чем pickle/unpickle

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

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

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

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

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

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

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

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

from multiprocessing import Process, Queue

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

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

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

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

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

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

Например

from multiprocessing import Process, Lock

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

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

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

from multiprocessing import Process, Lock

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Spec-Zone.ru

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