Python multiprocessing, потоки и asyncio: что выбрать

Python Автор: Среда и версия: CPython 3.14.5 с GIL; Fedora; Intel Core i7-12700KF; JIT выключен

Правило одно: процессы дают настоящую параллельность на CPU, потоки и asyncio ускоряют только ожидание. Считаешь числа в чистом Python, бери multiprocessing. Ждёшь сеть, соединений десятки, библиотека синхронная: потоки. Соединений тысячи: asyncio. Дальше эта развилка разобрана на одном примере и подтверждена замерами.

Сквозной пример: ночная джоба обходит 500 магазинов. По каждому надо забрать выгрузку из внутреннего API и пересчитать агрегат.

import time

def fetch(shop_id):
    time.sleep(0.2)            # вместо реального HTTP-запроса
    return shop_id

def recalc(shop_id):
    total = 0
    for amount in range(20_000_000):
        total += (amount * shop_id) % 97
    return total

fetch просто ждёт, процессор в это время свободен. recalc крутит чистый Python и занимает ядро целиком: 0.64 c на магазин.

Все числа ниже я померил сам на Intel i7-12700KF (12 ядер, 20 потоков), Fedora, CPython 3.14.5, обычная сборка с GIL, JIT выключен. На другой машине абсолютные значения будут другими, соотношения сохранятся.

Что выбрать: развилка в одну таблицу

Чем занята задачаЧем распараллеливатьЧто портит картину
ждёт сеть, соединений сотни и тысячиasyncioодин блокирующий вызов останавливает весь цикл событий
ждёт сеть или диск, соединений десяткипотоки10 000 потоков стоят 226 МБ, корутины столько не стоят
считает в чистом Pythonпроцессыстарт пула и сериализация аргументов не бесплатны
считает внутри NumPy, hashlib, кодековсначала потокиэти C-библиотеки отпускают GIL, замерь оба варианта
задач мало, каждая короче миллисекундыничегонакладные расходы съедят весь выигрыш

Вторая строка на практике решает чаще первой. Если в проекте requests, синхронный драйвер БД или чужой SDK без асинхронной версии, asyncio не даст ничего до полного переписывания стека, а потоки заработают сразу.

Почему потоки не ускоряют счёт

Возьмём восемь магазинов вместо пятисот, чтобы не ждать пять минут, и прогоним recalc тремя способами.

from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

if __name__ == "__main__":
    totals = [recalc(s) for s in range(1, 9)]                 # последовательно

    with ThreadPoolExecutor(max_workers=8) as pool:           # и то же с ProcessPoolExecutor
        totals = list(pool.map(recalc, range(1, 9)))

Охрана if __name__ == "__main__" тут обязательна: с потоками код проживёт и без неё, а с ProcessPoolExecutor упадёт. Зачем она нужна, разобрано ниже.

Лучшее из трёх прогонов:

последовательно  5.16 c
8 потоков        5.43 c
8 процессов      0.89 c

Потоки не дали ничего и отняли ещё четверть секунды сверху. Байт-код Python в каждый момент времени исполняет ровно один поток, остальные стоят в очереди за глобальной блокировкой; восемь потоков поделили ту же работу на восемь кусков и доплатили за переключения. Механика разобрана в статье про GIL, API потоков — в статье про threading.

Процессы дали 5.8x: у каждого свой интерпретатор и свой GIL, конкурировать им не за что.

Исключение, которое ломает правило: C-расширения отпускают GIL на время своей работы. Матричные операции NumPy, hashlib, сжатие, кодеки изображений — внутри этих вызовов остальные потоки Python бегут свободно. Если счёт идёт там, потоки могут выиграть и обойдутся дешевле процессов. Это проверяют замером, а не рассуждением.

Как выглядит multiprocessing: Process и Pool

Process запускает одну функцию в отдельном процессе, Pool раздаёт задачи пачке воркеров.

import multiprocessing as mp
import os

def report(shop_id):
    print(f"магазин {shop_id}: pid {os.getpid()}")

if __name__ == "__main__":
    p = mp.Process(target=report, args=(1,))
    p.start()
    p.join()
    print("exitcode:", p.exitcode)

    with mp.Pool(processes=8) as pool:
        totals = pool.map(recalc, range(1, 9))
    print("итогов:", len(totals))
магазин 1: pid 1832737
exitcode: 0
итогов: 8

join() ждёт завершения, exitcode показывает исход: 0 — успех, отрицательное число — номер сигнала, которым процесс убили (SIGKILL даёт -9). У concurrent.futures.ProcessPoolExecutor своя реализация, не Pool, но идея та же: процессы плюс pickle через трубу. Зато интерфейс общий с потоками: submit возвращает Future, а перевести код с процессов на потоки можно заменой одного класса. У Pool взамен есть imap, imap_unordered и starmap, и дефолтный chunksize у них разный, см. «Типичные ошибки».

Аргументы и результаты едут между процессами через pickle, поэтому не проедет ничего, что им не сериализуется. Лямбда даёт PicklingError: Can't pickle <function <lambda> ...>: it's not found as __main__.<lambda>, открытый файл — TypeError: cannot pickle 'TextIOWrapper' instances. Так же не переживают дорогу сокет, соединение с БД, курсор и замыкание. Отсюда два требования: функцию объявляй на верхнем уровне импортируемого модуля, а в аргументы клади данные, а не ресурсы. Передавай путь и строку подключения, открывай файл и соединение уже внутри воркера.

Зачем в multiprocessing нужен if __name__ == "__main__"

Без него код падает ещё до первой полезной задачи: RuntimeError: An attempt has been made to start a new process before the current process has finished its bootstrapping phase.

Новый процесс не наследует память родителя, поэтому он импортирует твой __main__ заново, чтобы найти там recalc. Если создание пула лежит на верхнем уровне модуля, воркер выполнит его при импорте и запустит собственных воркеров, те своих. CPython не даёт этой рекурсии начаться: spawn._check_not_importing_main() видит, что процесс прямо сейчас импортирует __main__, и сразу бросает RuntimeError. Охрана оставляет в модуле только определения, а запуск прячет от импорта.

На Linux этот код годами работал и без охраны, потому что старый дефолт fork ничего не импортировал. В 3.14 дефолт сменили.

fork, spawn и forkserver: что стоит по умолчанию

Проверено на месте, CPython 3.14.5, Linux:

>>> mp.get_start_method()
'forkserver'
>>> mp.get_all_start_methods()
['forkserver', 'fork', 'spawn']

Раньше на Linux по умолчанию стоял fork. В 3.14 дефолтом стал forkserver; в multiprocessing/context.py рядом с выбором лежит комментарий про gh-84559 и «thread safeish». На macOS с 3.8 по умолчанию spawn, потому что после fork() там ненадёжно выполнять произвольный код. На Windows spawn вообще единственный вариант.

Цена запуска пула из 4 воркеров в модуле с тяжёлыми импортами, минимум из четырёх прогонов:

МетодКак создаёт процессПул из 4Ограничения
forkкопирует текущий процесс целиком3.3 мстолько POSIX, ломается рядом с потоками
forkserverфоркает чистый служебный процесс8.0 мстолько POSIX, модуль импортируется один раз
spawnподнимает новый интерпретатор с нуля51.7 мсработает везде, каждый воркер импортирует сам

forkserver выбран дефолтом ровно из-за этой середины: безопасность почти как у spawn, цена ближе к fork.

Смена дефолта ломает один конкретный класс кода: тот, что настраивал глобальное состояние в родителе и рассчитывал увидеть его в воркере.

CONFIG = {}                                # модульный уровень

def read_config(_):
    return CONFIG.get("rate", "пусто")

if __name__ == "__main__":
    CONFIG["rate"] = 1.07                  # прочитали конфиг на старте
    for method in ("fork", "forkserver", "spawn"):
        with mp.get_context(method).Pool(2) as pool:
            print(f"{method:<11} ->", pool.map(read_config, [0, 1]))
fork        -> [1.07, 1.07]
forkserver  -> ['пусто', 'пусто']
spawn       -> ['пусто', 'пусто']

Воркер импортировал модуль заново и получил CONFIG таким, каким тот был в момент импорта. Лечится передачей конфига аргументом или через initializer пула.

Почему fork опасен рядом с потоками

fork() копирует только тот поток, который его вызвал. Остальные в ребёнке не появляются, а вот их мьютексы копируются как есть, вместе с состоянием «занято». Отпустить такой мьютекс некому.

import threading

lock = threading.Lock()

def holder():
    with lock:
        time.sleep(5)                          # держим лок в момент fork

def child():
    print("ребёнок стартовал", flush=True)
    with lock:                                 # лок скопирован занятым
        print("ребёнок взял лок", flush=True)

if __name__ == "__main__":
    threading.Thread(target=holder, daemon=True).start()
    p = mp.get_context("fork").Process(target=child)
    p.start()
    p.join(timeout=3)
    print("жив через 3 c:", p.is_alive())
    p.kill()
    p.join()
ребёнок стартовал
жив через 3 c: True

Вторая строка про лок не появится никогда. Убить ребёнка обязательно: без p.kill() не завершится и родитель, потому что на выходе multiprocessing джойнит живых детей.

Интерпретатор умеет предупредить, но по умолчанию молчит: DeprecationWarning фильтруется везде, кроме __main__, а срабатывает он не в твоём коде, а внутри multiprocessing/popen_fork.py. Запусти с python3 -W default или -X dev:

/usr/lib64/python3.14/multiprocessing/popen_fork.py:76: DeprecationWarning: This process (pid=2126372) is multi-threaded, use of fork() may lead to deadlocks in the child.
  self.pid = os.fork()

В обычном прогоне сигнала не будет. Практический вывод жёстче, чем кажется по игрушечному примеру. Потоки под капотом держит почти всё: пул соединений к БД, gRPC-клиент, Sentry SDK, цикл asyncio со своим исполнителем. Форкать живое приложение — лотерея, где выигрыш незаметен, а проигрыш выглядит как зависший воркер без стектрейса.

Сколько стоят процессы

Замер вхолостую, никакой полезной работы:

пул из 8 потоков              0.3 мс
пул из 8 процессов           20.0 мс
sum() по 200k чисел локально  0.9 мс
то же через 2 процесса       20.3 мс

Двадцать миллисекунд на пул, и это уже прогретый forkserver; самый первый пул в процессе стоит около 70 мс. Список на 200 000 чисел, отправленный через границу процесса, превращает миллисекундную работу в двадцатимиллисекундную: его надо упаковать, протолкнуть через трубу и распаковать обратно.

Память тоже своя у каждого воркера. Под spawn модуль исполняется заново в каждом, так что тяжёлые импорты оплачиваются столько раз, сколько воркеров. Под forkserver модуль импортируется один раз, в служебном процессе, и воркеры форкаются уже с готовыми импортами: ровно поэтому пул из четырёх стоит 8 мс, а не 52. Процессы окупаются, когда на воркера приходятся хотя бы десятые доли секунды счёта, а через границу едут небольшие данные.

Где выигрывает asyncio, а где потоки

Возвращаемся к fetch и берём все 500 магазинов:

последовательно      100.03 c
500 потоков            0.25 c
asyncio                0.20 c

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

10000 потоков:       старт 0.56 c, RSS 19 -> 226 МБ
10000 задач asyncio: старт 0.21 c, RSS 19 -> 31 МБ

Поток — объект ядра со своим стеком, корутина живёт в куче. Отсюда 226 МБ против 31. Пока соединений десятки, бери то, что проще написать; когда их тысячи, потоки упираются в память и планировщик. Как устроен цикл событий и почему один синхронный вызов вешает всю программу, разобрано в статье про asyncio.

Как совместить asyncio с потоками и процессами

Схем две, и обе нужны там, где в одной джобе есть и ожидание, и счёт. asyncio.to_thread(fn, arg) уносит блокирующий вызов в поток: синхронная библиотека перестаёт останавливать цикл событий. Второй рычаг тяжелее. loop.run_in_executor(pool, fn, arg) с ProcessPoolExecutor отдаёт счёт отдельным процессам, пока корутины держат соединения. Наш пример укладывается в обе сразу:

raw = await asyncio.gather(*(asyncio.to_thread(fetch, s) for s in range(1, 9)))

loop = asyncio.get_running_loop()
with ProcessPoolExecutor(max_workers=8) as pool:
    totals = await asyncio.gather(*(loop.run_in_executor(pool, recalc, s) for s in raw))

Выгрузка восьми магазинов через to_thread заняла 0.20 c, пересчёт в пуле процессов 0.85 c. Обе части работают так же, как поодиночке, только управляет ими один цикл событий.

Типичные ошибки

Мелкие задачи без chunksize. Двести тысяч вызовов x + 1 на восьми воркерах, одна и та же работа четырьмя вызовами:

Pool.map (умолчание)        0.02 c
Pool.imap (умолчание)       9.17 c
ProcessPoolExecutor.map     8.47 c
  он же с chunksize=5000    0.04 c

Разница в 400 раз, и вся она в упаковке и пересылке. chunksize считают сами только методы Pool, идущие через _map_async: map, starmap, map_async, starmap_async. Длина делится на число воркеров, умноженное на четыре, с округлением вверх. imap, imap_unordered и ProcessPoolExecutor.map по умолчанию гонят задачи по одной, и chunksize там надо ставить руками.

set_start_method после get_start_method. Чтение через get_start_method() фиксирует контекст, а get_start_method(allow_none=True) нет, и следующая попытка сменить метод даёт RuntimeError: context has already been set. Ставь метод первой строкой в __main__ либо не трогай глобальный дефолт и работай через mp.get_context("spawn").

Общая переменная вместо общего состояния. Счётчик, который увеличивают четыре задачи на двух воркерах под fork, даёт [1, 1, 2, 2] в детях и ровный 0 в родителе. Copy-on-write работает на чтение, любая запись создаёт приватную страницу. Для обмена есть Queue, Pipe, Value, Manager, и все они платят сериализацией.

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

Своё исключение с нестандартным __init__. Pickle восстанавливает исключение вызовом Класс(*args), поэтому конструктор с двумя обязательными аргументами не соберётся. У Pool на этом умирает поток, читающий результаты, и вместо ошибки получаешь вечное ожидание: TypeError: BadError.__init__() missing 1 required positional argument улетает в stderr, а map не возвращается никогда. Наследуйся от стандартных исключений и клади в них строку.

Частые вопросы

Сколько воркеров запускать

Для CPU-задач отправная точка — число ядер, дальше подбирается замером; воркеров больше, чем ядер, смысла не имеет, они начнут вытеснять друг друга. Для I/O-потоков ориентир другой: столько, сколько одновременных запросов выдержит тот, к кому ты ходишь. os.process_cpu_count() учитывает маску affinity, os.cpu_count() её игнорирует, и под taskset или --cpuset-cpus второй наврёт в разы. А вот квоту cgroup не видит ни один из них: под docker run --cpus=2 на этой двадцатипоточной машине оба вернули 20 при cpu.max = 200000 100000. Там число воркеров бери из /sys/fs/cgroup/cpu.max или из явной переменной окружения.

Работает ли multiprocessing из REPL и Jupyter

Под fork да, под forkserver и spawn нет. Воркеру нужно импортировать модуль с твоей функцией, а у интерактивной сессии файла нет. В stderr прилетит FileNotFoundError: [Errno 2] No such file or directory: '/твоя/директория/<stdin>' от служебного процесса, а в самом скрипте ConnectionResetError: [Errno 104] Connection reset by peer. Выноси функции в отдельный .py и импортируй их в ноутбук.

Отменяет ли сборка без GIL всю эту развилку

Не целиком. В сборке free-threading потоки действительно считают параллельно, и recalc в восьми потоках там ускорится. Свою сборку проверь вызовом sys._is_gil_enabled(): обычные дистрибутивные сборки, включая ту, на которой сделаны замеры выше, возвращают True. Процессы останутся нужны там, где важна изоляция по памяти и по падению.

Что учить дальше

Развилка помогает, только когда знаешь, что стоит за каждой веткой. Начни с той, которую выбрал: потоки и threading или asyncio и цикл событий. Если непонятно, откуда вообще берётся ограничение на CPU, читай разбор GIL.

Источники