Pandora ResearchPandora
Research
RUEN
Python

Процессы в Python

Pandora ResearchPandora Research
9 июня 20269 мин чтения

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

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

Базовый пример

Рассмотрим базовый пример:

Python
import time
import os

from multiprocessing import Process, current_process


def foo(text: str) -> None:
    time.sleep(3)

    print('[foo] process name:', current_process().name)
    print('[foo] process pid:', current_process().pid)
    print('msg:', text)


if __name__ == '__main__':
    print('Main process pid:', os.getpid())

    p = Process(target=foo, args=('Hello world', ), name='foo-process')
    process_pid = p.pid
    process_is_alive = p.is_alive()
    print(f'Process with pid: {process_pid} is alive: {process_is_alive}')

    print('Process started')
    p.start()
    process_pid = p.pid
    process_is_alive = p.is_alive()
    print(f'Process with pid: {process_pid} is alive: {process_is_alive}')

    print('Waiting for the process ends...')
    p.join()
    process_pid = p.pid
    process_is_alive = p.is_alive()
    print(f'Process with pid: {process_pid} is alive: {process_is_alive}')

    print('End program')

Вывод в консоль:

Terminal
Main process pid: 24569
Process with pid: None is alive: False
Process started
Process with pid: 24570 is alive: True
Waiting for the process ends...
[foo] process name: foo-process
[foo] process pid: 24570
msg: Hello world
Process with pid: 24570 is alive: False
End program

Также, пока выполняется скрипт, во втором окне терминала выполните команду ps aux | grep <pid_number> для каждого полученного pid. В результате будет примерно следующее:

Terminal
~> ps aux | grep 24569
darksto+   24569  0.6  0.0  28316 11192 ?        S    14:58   0:00 ~/PycharmProjects/edu/async-python-sprint-1/.venv/bin/python ~/.config/JetBrains/PyCharm2022.3/scratches/scratch.py
darksto+   24578  0.0  0.0  17868  2280 pts/0    S+   14:58   0:00 grep --color=auto 24569

~> ps aux | grep 24570
darksto+   24570  0.0  0.0  28316  8968 ?        S    14:58   0:00 ~/PycharmProjects/edu/async-python-sprint-1/.venv/bin/python ~/.config/JetBrains/PyCharm2022.3/scratches/scratch.py
darksto+   24590  0.0  0.0  17868  2312 pts/0    S+   14:58   0:00 grep --color=auto 24570

Видим, что в один момент времени существуют два процесса, с соответствующими pid.

Разберём пример кода по шагам. Вначале добавляем необходимые импорты и функцию foo. Внутри функции добавляем вывод названия и pid текущего процесса, в котором выполняется функция.

Начинаем основной блок программы. С помощью модуля os выводим pid основного процесса.

Создаём экземпляр класса Process, передав в него через параметр target callable-объект (функцию foo), её аргументы и название процесса. Метод is_alive() позволяет определить, запущен процесс или нет. С помощью метода start() запускаем процесс, а добавление метода join() в основной код позволяет дождаться завершения дочернего процесса.

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

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

Взаимодействие между процессами

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

Для обеспечения этой потребности существуют различные модели обмена данными. Но мы остановимся на одной из них — очереди.

Очередь — это структура данных с типом FIFO (First In, First Out), «первым пришёл — первым ушёл». Очереди позволяют обмениваться сообщениями между группами процессов.

Задача «producer-consumer» описывает два процесса: один является поставщиком данных или задач, то есть их производителем, а второй получает их и обрабатывает, то есть потребляет и удаляет их из очереди. Совместно они используют общий буфер для обмена сообщениями — очередь.

Рассмотрим на примере:

Python
from multiprocessing import Process, get_context
from multiprocessing.queues import Queue
import time


class Producer(Process):
    def __init__(self, queue: Queue):
        super().__init__()
        self.__queue = queue

    def run(self):
        for msg in range(5):
            self.__queue.put(msg)
            time.sleep(0.5)


class Consumer(Process):
    def __init__(self, queue: Queue, queue_result: Queue):
        super().__init__()
        self.__queue = queue
        self.__queue_result = queue_result

    def run(self):
        while True:
            if self.__queue.empty():
                print('Queue is empty. Exit.')
                break
            else:
                item = self.__queue.get()
                res = item ** 2
                print(res)
                self.__queue_result.put(res)
                time.sleep(1)


if __name__ == '__main__':
    data = []
    context = get_context('spawn')
    queue = Queue(ctx=context)
    queue_result = Queue(ctx=context)

    producer_process = Producer(queue)
    consumer_process = Consumer(queue, queue_result)
    producer_process.start()
    consumer_process.start()
    producer_process.join()
    consumer_process.join()

    while not queue_result.empty():
        data.append(queue_result.get())
    print(data)

Сначала импортируем необходимые библиотеки. Затем создаём класс Producer, унаследованный от Process. В методе __init__ сохраняем очередь и переопределяем метод run. По умолчанию этот метод вызывается в дочернем процессе после инициализации. Если он не определён, то вызывается функция, переданная в параметре target. В текущей реализации мы просто генерируем числа от 0 до 4 и кладём их в очередь queue.

Далее реализован класс Consumer, также от родительского класса Process. В конструктор класса передаём уже две очереди. Одна, из которой будем вычитывать данные, и вторая, в которую будем помещать результат расчётов. В методе run производим возведение в квадрат, перемещение результата в очередь результатов, а также проверку — пуста ли очередь. Если пуста, то завершаем работу консьюмера. Если этого не сделать, то выход из программы не произойдёт и она будет работать постоянно.

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

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

Результат работы программы представлен ниже.

Terminal
0
1
4
9
16
Queue is empty. Exit.
[0, 1, 4, 9, 16]

Также разберём отдельно функцию get_context. Для каждой очереди необходимо указать контекст. Всего доступно три варианта:

spawn (используется по умолчанию) — основной процесс запускает новый процесс интерпретатора Python. Наиболее медленный метод из трёх, но доступен на всех ОС.

fork — использует системную команду fork для создания дочерних процессов. Быстрее, чем spawn, но недоступен на Windows.

forkserver — задействуется специальный процесс-сервер, к которому обращается главный процесс программы, чтобы создать новый процесс. Ненужные ресурсы не копируются. Доступен на всех unix-based системах.

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

Для упрощения работы с несколькими параллельными процессами можно использовать класс Pool из библиотеки multiprocessing.

Основными методами экземпляра класса Pool являются:

apply — эквивалентен обычному вызову функции. Блокируется до момента получения результата и выполняется только в одном рабочем процессе пула. Синтаксис: apply(func, *args).

map — синтаксис: map(func, iter, [chunksize]); разделяет все итеративные данные iter на указанное число фрагментов chunksize, которые поставляются пулу процессов в виде отдельных задач. Блокирует пул до тех пор, пока не получен конечный результат.

imap — lazy-версия map. Позволяет получать результаты по мере их завершения.

apply_async — неблокирующий вызов apply. Синтаксис: apply_async(func, *args, [callback]). Результаты выполнения могут быть переданы в callback-функцию для их дальнейшей обработки, либо можно получить результат с помощью метода get().

Рассмотрим несколько примеров.

Python
import os
from multiprocessing import Pool, current_process
import time
import sys

import logging

logger = logging.getLogger(__name__)
logger.setLevel(logging.DEBUG)
handler = logging.StreamHandler(stream=sys.stdout)
handler.setFormatter(logging.Formatter(fmt='%(asctime)s: %(message)s'))
logger.addHandler(handler)


def calc(value: int) -> int:
    logger.info(f'calc started in process {current_process().name} with pid {current_process().pid}')
    time.sleep(1)
    return value**10


if __name__ == '__main__':
    logger.info(f'main pid: {os.getpid()}')
    with Pool(processes=2) as pool:
        result1 = pool.apply(calc, (10,))
        logger.info('result1: %s', result1)
        result2 = pool.apply(calc, (10,))
        logger.info('result2: %s', result2)
        result3 = pool.apply(calc, (10,))
        logger.info('result3: %s', result3)
        result4 = pool.apply(calc, (10,))
        logger.info('result4: %s', result4)


    print('End program')

Пример вывода:

Terminal
2023-01-25 18:21:40,027: main pid: 34098
2023-01-25 18:21:40,041: calc started in process ForkPoolWorker-1 with pid 34099
2023-01-25 18:21:41,044: result1: 10000000000
2023-01-25 18:21:41,047: calc started in process ForkPoolWorker-2 with pid 34100
2023-01-25 18:21:42,050: result2: 10000000000
2023-01-25 18:21:42,050: calc started in process ForkPoolWorker-1 with pid 34099
2023-01-25 18:21:43,052: result3: 10000000000
2023-01-25 18:21:43,053: calc started in process ForkPoolWorker-2 with pid 34100
2023-01-25 18:21:44,056: result4: 10000000000
End program

Рассмотрим пример с apply. В начале скрипта происходит добавление необходимых импортов и базовая настройка логгера для отображения времени. Определяем базовую функцию calc, которая будет возводить в 10-ю степень полученное число.

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

В консоли можно увидеть, что участвуют два дочерних процесса с pid 34100 и 34099, но при этом параллельно они не работают: задержка между каждым выводом результата составляет одну секунду. Никакого выигрыша в производительности в данном случае нет.

Рассмотрим пример с функцией map:

Python
import os
from multiprocessing import Pool, current_process
import time
import sys

import logging

logger = logging.getLogger(__name__)
logger.setLevel(logging.DEBUG)
handler = logging.StreamHandler(stream=sys.stdout)
handler.setFormatter(logging.Formatter(fmt='%(asctime)s: %(message)s'))
logger.addHandler(handler)


def calc(value: int) -> int:
    logger.info(f'calc started in process {current_process().name} with pid {current_process().pid}')
    time.sleep(1)
    return value**10

if __name__ == '__main__':
    logger.info(f'main pid: {os.getpid()}')

    l = [1, 2, 5, 10]
    with Pool(processes=2) as pool:
        results = pool.map(calc, l)
        print(results)

    print('End program')

Вывод в консоль:

Terminal
2023-01-25 18:28:27,135: main pid: 34311
2023-01-25 18:28:27,148: calc started in process ForkPoolWorker-2 with pid 34313
2023-01-25 18:28:27,148: calc started in process ForkPoolWorker-1 with pid 34312
2023-01-25 18:28:28,150: calc started in process ForkPoolWorker-1 with pid 34312
2023-01-25 18:28:28,150: calc started in process ForkPoolWorker-2 with pid 34313
[1, 1024, 9765625, 10000000000]
End program

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

В консоли видим, что два дочерних процесса стартовали одновременно, затем ещё два. Результат представлен в виде списка выходных данных функции calc. Также следует обратить внимание на то, что список упорядочен в соответствии со входными данными в списке l.

Попробуем вариант с использованием apply_async:

Python
import os
from multiprocessing import Pool, current_process
import time
import sys

import logging

logger = logging.getLogger(__name__)
logger.setLevel(logging.DEBUG)
handler = logging.StreamHandler(stream=sys.stdout)
handler.setFormatter(logging.Formatter(fmt='%(asctime)s: %(message)s'))
logger.addHandler(handler)


def calc(value: int) -> int:
    logger.info(f'calc started in process {current_process().name} with pid {current_process().pid}')
    time.sleep(1)
    return value**10


def callback(result):
    logger.info(f'Callback result: {result}, process pid: {current_process().pid}')


if __name__ == '__main__':
    logger.info(f'main pid: {os.getpid()}')

    with Pool(processes=2) as pool:
        result1 = pool.apply_async(calc, (1,), callback=callback)
        logger.info('result1: %s', result1)
        result2 = pool.apply_async(calc, (2,), callback=callback)
        logger.info('result2: %s', result2)
        result3 = pool.apply_async(calc, (3,), callback=callback)
        logger.info('result3: %s', result3)
        result4 = pool.apply_async(calc, (4,), callback=callback)
        logger.info('result4: %s', result4)

        pool.close()
        pool.join()

    print('End program')

По сравнению с предыдущими примерами добавлена функция callback, которая будет вызываться при завершении процесса. В основном коде происходит последовательная передача чисел 1, 2, 3, 4 в функцию calc, передаётся также callback для каждого вызова. В конце контекстного менеджера добавлены две инструкции — close() и join(). close() предотвращает отправку новых задач в пул. Как только все задачи будут выполнены, рабочие процессы завершатся, а join() ожидает завершения рабочих процессов. Без указания этих команд основная программа просто завершит свою работу, не дождавшись выполнения процессов и вызова callback-функции.

Вывод в консоль:

Terminal
2023-01-25 18:33:03,614: main pid: 34515
2023-01-25 18:33:03,625: result1: <multiprocessing.pool.ApplyResult object at 0x7fa55b754100>
2023-01-25 18:33:03,626: result2: <multiprocessing.pool.ApplyResult object at 0x7fa55b754130>
2023-01-25 18:33:03,626: result3: <multiprocessing.pool.ApplyResult object at 0x7fa55b754340>
2023-01-25 18:33:03,626: result4: <multiprocessing.pool.ApplyResult object at 0x7fa55b754460>
2023-01-25 18:33:03,626: calc started in process ForkPoolWorker-1 with pid 34516
2023-01-25 18:33:03,626: calc started in process ForkPoolWorker-2 with pid 34517
2023-01-25 18:33:04,628: calc started in process ForkPoolWorker-2 with pid 34517
2023-01-25 18:33:04,628: calc started in process ForkPoolWorker-1 with pid 34516
2023-01-25 18:33:04,628: Callback result: 1024, process pid: 34515
2023-01-25 18:33:04,628: Callback result: 1, process pid: 34515
2023-01-25 18:33:05,630: Callback result: 1048576, process pid: 34515
2023-01-25 18:33:05,630: Callback result: 59049, process pid: 34515
End program

Здесь мы видим, что все задачи для процессов были созданы одновременно. Функция apply_async вернула объект класса multiprocessing.pool.ApplyResult. Затем стартовали первые два процесса с pid 34516 и 34517, а затем вторые два.

После этого происходил вызов callback-функции в основном процессе с pid 34515, но результаты перемешались. Сначала вывод для инструкции, которая была второй, затем вывод результата для первой инструкции и т.д. apply_async не гарантирует корректной последовательности получения результата, в отличие от функции map.

Заключение

Рассмотрим ещё раз основные достоинства и недостатки процессов в Python.

Достоинства: возможность использовать все доступные ядра процессора; GIL не накладывает своих ограничений; у каждого процесса своя область памяти; есть возможность прервать процесс; подходит для тяжёлых вычислений.

Недостатки: более сложный интерфейс и большие накладные расходы; большое потребление памяти.

Python

Есть задача для нашей команды?

Расскажите о проекте — оценим и предложим решение в течение одного рабочего дня.

Обсудить проект