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

Последнее обновление: 13.09.2026

Обмен данными между процессами через Queue (Очередь)

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

multiprocessing.Queue([maxsize])

В качестве необязательного параметра конструктор Queue принимает размер очереди. В кратце рассмотрим основные методы класса:

  • empty(): возвращает True, если очередь пуста, и False в противном случае. Из-за особенностей многопоточности и многопроцессности результат не является абсолютно надежным.

  • full(): возвращает True, если очередь заполнена, и False в противном случае. Из-за особенностей многопоточности и многопроцессности результат не является абсолютно надежным.

  • put(obj[, block[, timeout]]): помещает объект obj в очередь. Если необязательный аргумент block равен True (значение по умолчанию), а timeout равен None (значение по умолчанию), то при необходимости выполнение блокируется до тех пор, пока не освободится место.

  • get([block[, timeout]]): извлекает и возвращает элемент из очереди. Если необязательный аргумент block равен True (значение по умолчанию), а timeout равен None (значение по умолчанию), то при необходимости выполнение блокируется до появления элемента.

Рассмотрим небольшой пример:

from multiprocessing import Process, Queue

# возводим числа в куб
def cubes(numbers, queue):
    for num in numbers:
        queue.put(num ** 3) # Отправляем данные в очередь
        
if __name__ == "__main__":
    numbers_list = [1, 2, 3]
    data_queue = Queue()    

    # Передаем очередь как аргумент в процесс
    proc = Process(target=cubes, args=(numbers_list, data_queue))
    proc.start()
    proc.join()

    # Извлекаем данные из очереди в главном процессе
    while not data_queue.empty():
        print(f"Получено из процесса cubes: {data_queue.get()}")

Здесь основное действие запускаемого процесса представляет функцию cubes, которая принимает набор чисел numbers и очередь queue. Перебирает набор и отправляет их кубы в очередь (queue.put(num ** 3)

В блоке __main__ получаем данные из очереди с помощью вызова data_queue.get(). Консольный вывод программы:

Получено из процесса cubes: 1
Получено из процесса cubes: 8
Получено из процесса cubes: 27

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

from multiprocessing import Process, Queue

# возводим числа в куб
def cubes(inputs: Queue, outputs: Queue):
    # Извлекаем данные из входной очереди inputs
    while True:
        n = inputs.get()
        if not n: break
        outputs.put(n ** 3) # Отправляем данные в очередь
        
if __name__ == "__main__":

    inputs = Queue()
    outputs = Queue()

    input_values = [1, 2, 3, 4, 5, 0]

    # передаем данные из списка во входную очередь - inputs
    for i in input_values:
        inputs.put(i)
    # Передаем очередь как аргумент в процесс
    proc = Process(target=cubes, args=(inputs, outputs))
    proc.start()
    proc.join()

    # Извлекаем данные из очереди в главном процессе
    for n in input_values:
        if not n: break
        print(f"{n}: {outputs.get()}")
   

В главном процессе создаются две очереди: inputs (для входящих чисел) и outputs (для результатов). В очередь inputs последовательно добавляются числа от 1 до 5. Последнее число - 0 служит признаком конца очереди.

Затем запускается отдельный процесс proc, который выполняет функцию cubes. В качестве аргументов ему передаются обе очереди. Функция cubes в цикле извлекает числа из inputs, возводит их в куб и отправляет результат в очередь outputs. Цикл завершается, когда входная очередь пустеет - когда дойдет до 0.

Главный процесс вызывает proc.join(), останавливая свою работу до тех пор, пока процесс proc полностью не завершится. После этого главный процесс поочередно достает значения из outputs и выводит их на экран.

1: 1
2: 8
3: 27
4: 64
5: 125

Передача сложных объектов между процессами

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

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

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

Общий процесс работы выглядит следующим образом:

  1. BaseManager запускает отдельный управляющий процесс (сервер).

  2. В этом процессе создается реальный экземпляр нашего сложного объекта.

  3. Другие процессы вызывают методы объекта через так называемые прокси-объекты (или клиенты).

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

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

Сначала создадим обычный класс, который естественно ничего не знает о процессах и никак с ними не связан:

class Inventory:
    def __init__(self):
        self.items = {}

    def add_item(self, name: str, quantity: int):
        if name in self.items:
            self.items[name] += quantity
        else:
            self.items[name] = quantity
        print(f"[Склад] Добавлено: {name}, количество {quantity}")

    def get_stock(self) -> dict:
        return self.items

Класс довольно небольшой. Тем не менее разберем его. Прежде всего внутри конструктора __init__() создается пустой словарь (self.items = {}). В нем будут храниться пары вида "название_товара": количество.

Метод add_item() будет использоваться для добавления товара. Этот метод принимает два аргумента: name (название товара, тип str) и quantity (количество, тип int). Двоеточия и типы данных - это аннотации типов, которые помогают программисту понять, какие данные сюда нужно передавать.

Внутри метода с помощью выражения if name in self.items проверяем, есть ли уже такой товар в нашем словаре. И если товар уже есть, код self.items[name] += quantity прибавляет новое количество к старому. А если такого товара еще нет, товар создается в словаре с нуля (self.items[name] = quantity) и ему присваивается указанное количество.

И последний метод - получение остатков (get_stock) просто возвращает текущее содержимое словаря со всеми товарами. Причем метод возвращает данные в формате словаря.

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

from multiprocessing import Process 
from multiprocessing.managers import BaseManager
import time

class Inventory:
    def __init__(self):
        self.items = {}

    def add_item(self, name: str, quantity: int):
        if name in self.items:
            self.items[name] += quantity
        else:
            self.items[name] = quantity
        print(f"[Склад] Добавлено: {name}, количество {quantity}")

    def get_stock(self) -> dict:
        return self.items


# Создаем кастомный менеджер
class CustomManager(BaseManager):
    pass

# Функция, которую будут выполнять рабочие процессы
def worker_task(shared_inventory, item_name, quantity):
    print(f"[Процесс {item_name}] Начинает работу...")
    time.sleep(1)  # Имитация работы
    # Вызываем метод через прокси-объект так же, как у обычного класса
    shared_inventory.add_item(item_name, quantity)
    print(f"[Процесс {item_name}] Завершил добавление.")
    
if __name__ == "__main__":
    # Регистрируем наш класс в менеджере под именем "Inventory"
    CustomManager.register("Inventory", Inventory)

    # Запускаем менеджер
    with CustomManager() as manager:
        # Создаем общий объект через менеджер
        # Внутри менеджера создается реальный Inventory, а нам возвращается прокси
        shared_inventory = manager.Inventory()

        # Для демонстрации создаем два процесса и передаем им наш прокси-объект
        p1 = Process(target=worker_task, args=(shared_inventory, "Яблоки", 10))
        p2 = Process(target=worker_task, args=(shared_inventory, "Бананы", 5))

        p1.start()
        p2.start()

        # Ждем завершения процессов
        p1.join()
        p2.join()

        # Проверяем финальное состояние объекта в главном процессе
        print("\n--- Финальный отчет о складе ---")
        print(shared_inventory.get_stock())

Пройдем по основным моментам. Сначала создаем кастомный класс менеджера:

class CustomManager(BaseManager):

В блоке __main__ регистрируем наш класс в менеджере под именем "Inventory":

CustomManager.register("Inventory", Inventory)

Запускаем менеджер:

with CustomManager() as manager:

Создаем общий объект через менеджер:

shared_inventory = manager.Inventory()

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

И при запуске программы мы увидим в консоли:

[Процесс Яблоки] Начинает работу...
[Процесс Бананы] Начинает работу...
[Склад] Добавлено: Яблоки, количество 10
[Процесс Яблоки] Завершил добавление.
[Склад] Добавлено: Бананы, количество 5
[Процесс Бананы] Завершил добавление.

--- Финальный отчет о складе ---
{"Яблоки": 10, "Бананы": 5}

При работе с BaseManager следует учитывать ряд аспектов:

  • Менеджеры не обеспечивают потокобезопасность по умолчанию. Если два процесса одновременно вызовут метод, который меняет одно и то же поле (например, self.items["Яблоки"] += 1), может возникнуть состояние гонки (race condition). Если же методы класса делают сложные вычисления или перезаписывают данные, лучше использовать multiprocessing.Lock внутри методов класса для синхронизации.

  • Прокси-объекты возвращают копии данных, а не ссылки. Когда мы вызываем shared_inventory.get_stock(), мы получаем копию словаря items. Если же мы изменим полученный словарь (stock["Груши"] = 1), на реальном складе внутри менеджера ничего не изменится. Изменения нужно делать только через вызовы методов самого прокси.

  • Регистрация должна происходить до старта. Всегда регистрируйте классы (CustomManager.register...) до того, как вызовете manager.start() или войдете в контекст with CustomManager().

Помощь сайту
Юмани:
410011174743222
Номер карты:
4048415020898850