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) нового процесса.
Всегда оборачивай 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 + немного CPU | threading |
| Mixed - много IO + heavy CPU sometimes | asyncio + ProcessPoolExecutor для CPU |
| Простой batch processing | ProcessPoolExecutor.map |
| Тысячи concurrent connections | asyncio (потоки/процессы дорогие) |
Реальные 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 ускорение.
Мини-задание
- 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")
- 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)
- Передача данных через 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, упаковка проекта, тур по полезным модулям стандартной библиотеки.