multiprocessing: настоящая параллельность через процессы

В прошлом уроке мы видели что threading в Python ограничен GIL для CPU-bound задач. Решение - multiprocessing: запуск кода в отдельных процессах, каждый со своим интерпретатором и своим GIL. Это даёт настоящую параллельность на нескольких CPU. Цена - больший overhead создания и общения между процессами.

Зачем multiprocessing

CPU-bound задача в threading:

import threading
import time

def heavy():
    total = sum(i * i for i in range(10_000_000))

# 4 потока
threads = [threading.Thread(target=heavy) for _ in range(4)]
start = time.perf_counter()
for t in threads: t.start()
for t in threads: t.join()
print(f"Threading: {time.perf_counter() - start:.2f}s")
# ~примерно как 4 sequential выполнения из-за GIL

С multiprocessing:

from multiprocessing import Process

processes = [Process(target=heavy) for _ in range(4)]
start = time.perf_counter()
for p in processes: p.start()
for p in processes: p.join()
print(f"Multiprocessing: {time.perf_counter() - start:.2f}s")
# ~как 1 sequential если 4 CPU доступны

Каждый процесс работает независимо, имеет свой GIL, использует свой CPU.

Process - базовое использование

from multiprocessing import Process
import os

def worker(name):
    print(f"{name} in process {os.getpid()}")

if __name__ == "__main__":
    p = Process(target=worker, args=("A",))
    p.start()
    p.join()

API похоже на threading.Thread: target, args, start, join. Главное отличие - под капотом fork (Linux/macOS) или spawn (Windows) нового процесса.

На Windows и при `spawn` start method (Python 3.8+ default на macOS) дочерний процесс заново импортирует модуль. Без guard это вызовет рекурсивное создание процессов.

Всегда оборачивай multiprocessing-код в guard по __name__ (откуда берётся эта переменная - в уроке про импорты и пакеты):

if __name__ == "__main__":
    # код запуска процессов
    ...

На Linux с fork (legacy default) можно без guard, но лучше всегда писать - кроссплатформенно.

Передача данных между процессами

В отличие от потоков, процессы не разделяют память. Каждый имеет свой адресный пространство. Передача данных:

from multiprocessing import Process, Queue

def producer(queue):
    for i in range(5):
        queue.put(i)
    queue.put(None)   # sentinel

def consumer(queue):
    while True:
        item = queue.get()
        if item is None:
            break
        print(f"Got: {item}")

if __name__ == "__main__":
    queue = Queue()
    p1 = Process(target=producer, args=(queue,))
    p2 = Process(target=consumer, args=(queue,))
    p1.start(); p2.start()
    p1.join(); p2.join()

multiprocessing.Queue сериализует объекты (через pickle) для передачи между процессами. Поэтому объекты должны быть picklable (большинство стандартных типов да; lambdas, локальные функции - нет). pickle это бинарный родственник json-сериализации, только для произвольных Python-объектов.

Pipe - двунаправленная связь

from multiprocessing import Process, Pipe

def worker(conn):
    conn.send("hello from worker")
    response = conn.recv()
    print(f"Worker got: {response}")
    conn.close()

if __name__ == "__main__":
    parent_conn, child_conn = Pipe()
    p = Process(target=worker, args=(child_conn,))
    p.start()

    msg = parent_conn.recv()
    print(f"Parent got: {msg}")
    parent_conn.send("hi back")
    p.join()

Pipe - двунаправленный канал между двумя процессами. Быстрее Queue (нет внутренней очереди), но только для two-party communication.

ProcessPoolExecutor - удобный API

concurrent.futures.ProcessPoolExecutor похож на ThreadPoolExecutor, но создаёт процессы:

from concurrent.futures import ProcessPoolExecutor

def cpu_task(n):
    return sum(i * i for i in range(n))

if __name__ == "__main__":
    with ProcessPoolExecutor(max_workers=4) as executor:
        inputs = [1_000_000, 2_000_000, 3_000_000, 4_000_000]
        results = list(executor.map(cpu_task, inputs))
        print(results)

Этот API чаще всего используется в реальном коде - проще чем ручное управление Process. По умолчанию max_workers = количество CPU.

Pool из multiprocessing

multiprocessing.Pool - похожий API, был до ProcessPoolExecutor:

from multiprocessing import Pool

def square(n):
    return n * n

if __name__ == "__main__":
    with Pool(processes=4) as pool:
        results = pool.map(square, range(10))
        print(results)   # [0, 1, 4, 9, ..., 81]

Методы Pool:

  • map(func, iterable) - синхронное применение
  • map_async(...) - возвращает AsyncResult, не блокирует
  • apply(func, args) - синхронный вызов с аргументами
  • apply_async(func, args, callback=) - с callback на готовности
  • imap(func, iterable) - lazy map, iterator (как itertools.imap)

ProcessPoolExecutor более новый и часто предпочтительнее, но Pool активно используется в legacy и научных пакетах.

Sharing state - shared memory

Если процессам нужно работать с общими данными:

from multiprocessing import Value, Array

def increment(counter, lock):
    with lock:
        counter.value += 1

if __name__ == "__main__":
    from multiprocessing import Lock
    counter = Value("i", 0)   # shared int
    lock = Lock()

    processes = [Process(target=increment, args=(counter, lock)) for _ in range(10)]
    for p in processes: p.start()
    for p in processes: p.join()

    print(counter.value)   # 10

Value(typecode, initial) - shared скалярное значение (i=int, d=double, f=float). Array(typecode, list) - shared массив.

Используется редко - сложно отлаживать. Чаще проходить данные через Queue/Pipe.

Manager - shared объекты высокого уровня

from multiprocessing import Process, Manager

def worker(d, l):
    d["worker_pid"] = os.getpid()
    l.append(1)

if __name__ == "__main__":
    with Manager() as manager:
        shared_dict = manager.dict()
        shared_list = manager.list()

        processes = [Process(target=worker, args=(shared_dict, shared_list)) for _ in range(5)]
        for p in processes: p.start()
        for p in processes: p.join()

        print(dict(shared_dict))
        print(list(shared_list))

Manager создаёт сервер-процесс, который хостит shared объекты. Удобнее Value/Array (можно использовать dict, list, set), но медленнее (всё через сервер).

Pickle - сериализация для IPC

Передача данных между процессами требует pickle - стандартный сериализатор Python. Не всё picklable:

# OK
data = {"a": 1, "b": [1, 2, 3]}

# OK
def named_function(x): return x

# НЕ OK
lambda_func = lambda x: x          # lambdas не picklable
local_function = lambda x: ...     # локальные функции не picklable

В ProcessPoolExecutor target-функция должна быть на уровне модуля. Лямбды и closures внутри функций не передаются. Это типичная ошибка - решение через partial или вынесение в отдельные модули.

fork vs spawn

Способ создания процесса:

  • fork (Linux/macOS legacy): копирует весь процесс. Быстро, но проблемы с потоками и shared resources.
  • spawn (Windows, macOS default с Python 3.8): запускает новый интерпретатор, импортирует модуль. Медленнее, но безопаснее.
  • forkserver (Linux): промежуточный вариант.
import multiprocessing
multiprocessing.set_start_method("spawn")   # явно

Большинство нового кода работает с spawn. Если используешь fork, осторожно с pre-fork threads и open files.

Когда multiprocessing vs threading vs asyncio

ЗадачаЛучше
Чистый CPU (вычисления, ML)multiprocessing
Чистый IO (HTTP, БД, файлы)asyncio
Mixed - много sync IO + немного CPUthreading
Mixed - много IO + heavy CPU sometimesasyncio + ProcessPoolExecutor для CPU
Простой batch processingProcessPoolExecutor.map
Тысячи concurrent connectionsasyncio (потоки/процессы дорогие)

Реальные backend часто комбинируют: asyncio для основного flow + ProcessPoolExecutor для редких CPU-heavy задач (image processing, ML inference).

Производительность - overhead

Overhead процессов больше потоков:

  • Создание: процесс ~10-100ms (fork ~1ms на Linux), поток ~1ms
  • Memory: процесс ~10-50MB базы, поток ~1MB
  • IPC: процесс через pickle (медленно), поток разделяет память (быстро)

Для задач короче нескольких секунд multiprocessing невыгоден - overhead больше выигрыша. Для длительных вычислений overhead становится незначительным.

Распространённые ошибки

1. Забыли if __name__ == "__main__"

# bad.py - на Windows вызовет infinite процессы
from multiprocessing import Process

p = Process(target=lambda: print("hi"))
p.start()   # дочерний процесс импортирует bad.py и снова создаёт Process!

Всегда оборачивай top-level код в guard.

2. Lambdas как target

# Не работает с spawn
with ProcessPoolExecutor() as ex:
    ex.map(lambda x: x * 2, range(10))   # PicklingError

Используй именованные функции на module level или partial.

3. Shared state без synchronization

# Без Lock - race conditions
counter = Value("i", 0)
def bad(c):
    c.value += 1   # не атомарно

Используй Lock или multiprocessing.Manager с примитивами синхронизации.

4. Слишком много процессов

ProcessPoolExecutor(max_workers=1000)   # ОЧЕНЬ много памяти

Обычно max_workers = количество CPU достаточно. Больше не даёт прироста, ест память.

5. Передача больших данных через pickle

# 1GB list передаётся через pickle - медленно
ex.map(process, [huge_list])

Для больших данных используй shared memory (Array, Manager), общий файл, или mmap.

Сравнение с Go

В Go горутины и каналы - встроенная concurrency без отдельных процессов:

func worker(jobs <-chan int, results chan<- int) {
    for j := range jobs {
        results <- j * j
    }
}

jobs := make(chan int, 10)
results := make(chan int, 10)

for i := 0; i < 4; i++ {
    go worker(jobs, results)
}

Go использует один процесс с goroutines на всех CPU - проще и эффективнее. Python из-за GIL вынужден использовать процессы для настоящей параллельности.

Реальный пример - параллельная обработка изображений

from concurrent.futures import ProcessPoolExecutor
from PIL import Image
import os

def process_image(path):
    img = Image.open(path)
    img.thumbnail((200, 200))
    out_path = path.replace(".jpg", "_thumb.jpg")
    img.save(out_path)
    return out_path

if __name__ == "__main__":
    images = [f for f in os.listdir("photos") if f.endswith(".jpg")]
    paths = [os.path.join("photos", f) for f in images]

    with ProcessPoolExecutor(max_workers=4) as executor:
        results = list(executor.map(process_image, paths))

    print(f"Processed {len(results)} images")

Каждое изображение обрабатывается в отдельном процессе. На 4-ядерной машине - примерно 4x ускорение.

Мини-задание

  1. CPU-bound задача в нескольких процессах:
from multiprocessing import Process
import time

def heavy():
    return sum(i ** 2 for i in range(5_000_000))

if __name__ == "__main__":
    start = time.perf_counter()
    processes = [Process(target=heavy) for _ in range(4)]
    for p in processes: p.start()
    for p in processes: p.join()
    print(f"Parallel: {time.perf_counter() - start:.2f}s")

    start = time.perf_counter()
    for _ in range(4):
        heavy()
    print(f"Sequential: {time.perf_counter() - start:.2f}s")
  1. ProcessPoolExecutor для batch:
from concurrent.futures import ProcessPoolExecutor

def square(n):
    return n * n

if __name__ == "__main__":
    with ProcessPoolExecutor() as ex:
        results = list(ex.map(square, range(20)))
        print(results)
  1. Передача данных через Queue:
from multiprocessing import Process, Queue
import time

def worker(queue, name):
    for i in range(3):
        queue.put(f"{name}-{i}")
        time.sleep(0.1)

if __name__ == "__main__":
    queue = Queue()
    processes = [Process(target=worker, args=(queue, f"P{i}")) for i in range(3)]
    for p in processes: p.start()
    for p in processes: p.join()

    while not queue.empty():
        print(queue.get())

Что дальше

Модуль 7 завершён. Освоили iterators, generators, asyncio, threading с GIL, multiprocessing. У тебя теперь полное понимание моделей concurrency в Python. В следующем модуле перейдём к packaging и stdlib: imports, упаковка проекта, тур по полезным модулям стандартной библиотеки.

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