Spec-Zone.ru › Python 3.9

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

По умолчанию не передаются аргументы в target.

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

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

run()

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

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

start()

Запуск активности процесса.

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

join([timeout])

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

Процесс можно присоединить многократно.

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

name

Имя процесса. Имя — это строка, используемая только для идентификации. Оно не имеет семантики. Несколько процессов могут иметь одинаковое имя.

Начальное имя задаётся в конструкторе. Если явное имя не предоставлено конструктору, создаётся имя вида ‘Process-N1:N2:…:Nk’, где каждый Nk — N-й дочерний процесс своего родительского процесса.

is_alive()

Возвращает, жив ли процесс.

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

daemon

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

Начальное значение унаследовано от создающего процесса.

При выходе процесса он пытается завершить все свои демонические дочерние процессы.

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

В дополнение к API threading.Thread, объекты Process также поддерживают следующие атрибуты и методы:

pid

Возвращает идентификатор процесса. До создания процесса это будет None.

exitcode

Код выхода дочернего процесса. Будет None если процесс ещё не завершился.

Если метод run() дочернего процесса завершился успешно, код выхода будет 0. Если он завершился с помощью sys.exit() с целочисленным аргументом N, код выхода будет N.

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

authkey

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

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

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

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

sentinel

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

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

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

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

terminate()

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

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

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

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

kill()

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

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

close()

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

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

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

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

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

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

exception multiprocessing.BufferTooShort

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

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

exception multiprocessing.AuthenticationError

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

exception multiprocessing.TimeoutError

Возникает в методах с таймаутом, когда таймаут истекает.

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

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

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

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

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

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

Примечание

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

Примечание

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

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

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

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

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

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

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

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

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

multiprocessing.Pipe([duplex])

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

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

class multiprocessing.Queue([maxsize])

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

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

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

qsize()

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

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

empty()

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

full()

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

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

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

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

put_nowait(obj)

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

get([block[, timeout]])

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

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

get_nowait()

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

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

close()

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

join_thread()

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

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

cancel_join_thread()

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

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

Примечание

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

class multiprocessing.SimpleQueue

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

close()

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

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

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

empty()

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

get()

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

put(item)

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

class multiprocessing.JoinableQueue([maxsize])

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

task_done()

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

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

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

join()

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

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

Разное

multiprocessing.active_children()

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

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

multiprocessing.cpu_count()

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

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

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

См. также

os.cpu_count()

multiprocessing.current_process()

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

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

multiprocessing.parent_process()

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

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

multiprocessing.freeze_support()

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

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

from multiprocessing import Process, freeze_support

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

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

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

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

multiprocessing.get_all_start_methods()

Возвращает список поддерживаемых методов запуска, первый из которых является по умолчанию. Возможные методы запуска — 'fork', '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(executable)

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

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

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

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

multiprocessing.set_start_method(method)

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

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

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

Примечание

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

Объекты подключения

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

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

class multiprocessing.connection.Connection
send(obj)

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

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

recv()

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

fileno()

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

close()

Закрыть подключение.

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

poll([timeout])

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

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

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

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

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

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

recv_bytes([maxlength])

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

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

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

recv_bytes_into(buffer[, offset])

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

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

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

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

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

Например:

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

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

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

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

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

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

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

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

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

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

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

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

class multiprocessing.BoundedSemaphore([value])

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

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

Примечание

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

class multiprocessing.Condition([lock])

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

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

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

class multiprocessing.Event

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

class multiprocessing.Lock

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

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

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

acquire(block=True, timeout=None)

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

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

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

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

release()

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

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

class multiprocessing.RLock

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

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

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

acquire(block=True, timeout=None)

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

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

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

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

release()

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

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

class multiprocessing.Semaphore([value])

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

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

Примечание

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

Примечание

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

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

Примечание

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

Общие объекты ctypes

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

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

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

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

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

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

counter.value += 1

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

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

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

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

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

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

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

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

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

Модуль multiprocessing.sharedctypes

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

Примечание

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

multiprocessing.sharedctypes.RawArray(typecode_or_type, size_or_initializer)

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

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

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

multiprocessing.sharedctypes.RawValue(typecode_or_type, *args)

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

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

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

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

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

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

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

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

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

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

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

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

multiprocessing.sharedctypes.copy(obj)

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

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

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

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

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

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

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

ctypes

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

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

c_double(2.4)

RawValue(c_double, 2.4)

RawValue(‘d’, 2.4)

MyStruct(4, 6)

RawValue(MyStruct, 4, 6)

(c_short * 7)()

RawArray(c_short, 7)

RawArray(‘h’, 7)

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

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

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

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

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

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

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

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

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

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

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

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

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

Менеджеры

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

multiprocessing.Manager()

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

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

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

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

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

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

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

start([initializer[, initargs]])

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

get_server()

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

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

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

connect()

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

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

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

Это можно вызывать многократно.

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

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

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

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

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

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

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

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

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

address

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

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

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

class multiprocessing.managers.SyncManager

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

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

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

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

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

BoundedSemaphore([value])

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

Condition([lock])

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

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

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

Event()

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

Lock()

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

Namespace()

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

Queue([maxsize])

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

RLock()

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

Semaphore([value])

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

Array(typecode, sequence)

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

Value(typecode, value)

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

dict()
dict(mapping)
dict(sequence)

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

list()
list(sequence)

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

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

class multiprocessing.managers.Namespace

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

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

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

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

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

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

from multiprocessing.managers import BaseManager

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

class MyManager(BaseManager):
    pass

MyManager.register('Maths', MathsClass)

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

class multiprocessing.pool.AsyncResult

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

get([timeout])

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

wait([timeout])

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

ready()

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

successful()

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

Изменено в версии 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, он будет заблокирован на неопределённый срок. Отрицательный таймаут эквивалентен нулевому таймауту.

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

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

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

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

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

В версии 3.3.

Примеры

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

from multiprocessing.connection import Listener
from array import array

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

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

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

        conn.send_bytes(b'hello')

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

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

from multiprocessing.connection import Client
from array import array

address = ('localhost', 6000)

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

multiprocessing.get_logger()

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

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

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

multiprocessing.log_to_stderr(level=None)

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

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

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

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

Модуль multiprocessing.dummy

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

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

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

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

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

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

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

Примечание

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

from multiprocessing import Process, Queue

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

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

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

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

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

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

Например

from multiprocessing import Process, Lock

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

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

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

from multiprocessing import Process, Lock

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

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

Осторожно с заменой sys.stdin на «объект типа файла»

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

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

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

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

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

Если вы создаёте объект типа файла и реализуете собственный кэширование, вы можете сделать его безопасным при использовании 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.9/library/multiprocessing.html

Spec-Zone.ru

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