Конкурентность: дерево решений

Потоки, asyncio, процессы, субинтерпретаторы — не вкусовщина, а четыре разных ответа на два вопроса: где нагрузка проводит время и сколько данных придётся передавать. С числами.

L1 · 9 минL2 · 9 мин📱 телефонсверено · CPython 3.11 + numpy 2.4, linux x86-64

1Предскажи

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

Фрагмент A
import time
import numpy as np
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor

M = np.random.rand(1600, 1600)

def job(_):
    return float((M @ M).sum())

def run(Pool, label):
    with Pool(max_workers=4) as ex:
        list(ex.map(job, range(4)))          # прогрев: пул и сам код
        t0 = time.perf_counter()
        list(ex.map(job, range(4)))
        print(f"{label}: {time.perf_counter() - t0:.2f} c")

if __name__ == "__main__":                   # обязательно: ProcessPool импортирует модуль заново
    run(ThreadPoolExecutor, "потоки  ")
    run(ProcessPoolExecutor, "процессы")
потоки  : 0.27 c
процессы: 0.48 c
Процессы дают настоящий параллелизм — и проигрывают втрое.
Фрагмент B
import pickle, numpy as np
a = np.random.rand(5_000_000)   # 40 МБ
# сколько стоит pickle туда-обратно?
0.358 c
Только сериализация. Полезной работы ещё не было.
Фрагмент C
import asyncio, time

async def call():
    await asyncio.sleep(0.010)               # вместо сетевого запроса на 10 мс

async def sequential():
    for _ in range(50):
        await call()

async def concurrent():
    async with asyncio.TaskGroup() as tg:
        for _ in range(50):
            tg.create_task(call())

async def main():
    t0 = time.perf_counter(); await sequential()
    print(f"последовательно: {time.perf_counter() - t0:.3f} c")
    t0 = time.perf_counter(); await concurrent()
    print(f"одновременно:    {time.perf_counter() - t0:.3f} c")

asyncio.run(main())
последовательно: 0.510 c
одновременно:    0.012 c
Один поток, одно ядро, в одиннадцать раз быстрее.
Фрагмент D
import time, threading

def work(n=8_000_000):
    t = 0
    for i in range(n):
        t += i
    return t

t0 = time.perf_counter(); work(); work()
seq = time.perf_counter() - t0

ths = [threading.Thread(target=work) for _ in range(2)]
t0 = time.perf_counter()
for t in ths: t.start()
for t in ths: t.join()
par = time.perf_counter() - t0

print(f"последовательно: {seq:.2f} c")
print(f"два потока:      {par:.2f} c")
последовательно: 0.54 c
два потока:      0.50 c
Нулевой выигрыш — и это не баг.

2Механизм

Все четыре механизма разобраны в предыдущих уроках. Здесь они сводятся в одно решение, и решается оно двумя вопросами, а не перебором вариантов.

Вопрос первый: где нагрузка проводит время? В ожидании, в C-коде, отпустившем лок, или в байткоде.

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

Следствие 1 · ожидание — это asyncio, и разница на порядок

Фрагмент C. Пятьдесят сетевых вызовов последовательно — сумма ожиданий; одновременно — время самого долгого. Один поток, одно ядро, одиннадцатикратная разница.

Потоки решают ту же задачу, но каждый стоит стека ОС и переключения ядром: тысячи одновременных соединений на потоках — уже проблема, на корутинах — норма (урок 14). Порог примерно там, где счёт идёт на сотни.

Ограничение известно: нужны асинхронные библиотеки насквозь. Один синхронный вызов кладёт весь цикл.

Следствие 2 · счёт в C — это потоки, и они дешевле процессов

Фрагмент A — самый контринтуитивный. Процессы дают настоящий параллелизм без всякого GIL, а проигрывают потокам втрое. Причина в фрагменте B: сорок мегабайт данных едут через pickle, и одна только сериализация стоит 0.358 с — больше, чем вся полезная работа.

четыре задачи над массивом 40 МБвремя
ThreadPoolExecutor0.46 с
ProcessPoolExecutor1.24 с

Правило, выводимое отсюда: если работа уходит в C, начинай с потоков. Они делят память бесплатно, а параллелизм получают за счёт отпускания лока (уроки 15 и 19).

Следствие 3 · счёт в байткоде — процессы, и данные надо считать

Фрагмент D. Чистый Python держит лок всё время: потоки не дадут ничего, asyncio тем более. Остаются процессы.

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

Обратный случай — «переслать много, посчитать мало» — почти всегда проигрышный. Обходят его тем, что данные не пересылают: shared_memory, memory-mapped файлы, или процессы читают вход независимо по путям.

Следствие 4 · субинтерпретаторы — четвёртый вариант, пока узкий

С 3.12 каждый субинтерпретатор может иметь собственный лок (PEP 684): параллелизм в одном процессе без отдельных адресных пространств. Старт дешевле процесса, изоляция сильнее потока.

Но объекты между ними не разделяются, обмен идёт через каналы с сериализацией, а C-расширения должны быть к этому готовы — многие не готовы. То есть выигрыш по сравнению с процессами пока в основном в стоимости старта, а не в передаче данных.

Дерево целиком

где времячто братьпочему
ожидание сети или диска, сотни задачasyncioдешёвое переключение, нет стеков ОС
ожидание, но библиотека синхроннаяпотокилок отпускается на I/O
счёт внутри C-расширенияпотокилок отпущен, память общая
счёт в байткоде, данных малопроцессыединственный настоящий параллелизм
счёт в байткоде, данных многопроцессы + shared_memory
или переписать горячее в C
сериализация съедает выигрыш

Обрати внимание, что первый вопрос почти всегда решается замером, а не рассуждением: тест из урока 19 — прогнать последовательно и в двух потоках — отвечает на него за минуту.

3Границы модели

Где сказанное перестаёт держать

Смешивать можно и часто нужно. Асинхронный сервис с пулом процессов под тяжёлые задачи — нормальная архитектура: loop.run_in_executor связывает их. Плохо не смешение, а выбор наугад.

fork и потоки несовместимы 🕳. Форк из процесса с потоками оставляет в потомке заблокированные локи, взятые несуществующими потоками. Отсюда зависания при multiprocessing в сервисе с потоками и смена умолчания стартового метода в новых версиях.

Числа зависят от размера задачи. Все замеры здесь — про конкретные объёмы. На мелких задачах накладные расходы любого механизма перевешивают, и правильный ответ «не распараллеливать» тоже входит в дерево.

4Корень

Корень R3: счётчик ссылок и GIL. Всё дерево — прямое следствие того, что байткод исполняется под глобальным локом, а память в процессе общая.

Порядок появления объясняет, почему вариантов четыре, а не один. Потоки были с самого начала, вместе с локом. multiprocessing добавили в 2008-м именно как обход GIL — и вместе с ним пришла цена сериализации. asyncio появился в 2014-м, решая другую задачу: не параллелизм, а масштаб по числу ожиданий. Субинтерпретаторы с отдельным локом — 2023-й, попытка получить параллелизм без отдельных адресных пространств.

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

5Аналогия

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

Поэтому знакомая из C интуиция работает наполовину: «потоки дешевле в обмене» остаётся верным и остаётся главным аргументом, а «потоки дают параллелизм» верно только за пределами интерпретатора. Отсюда и контринтуитивный результат замера, где процессы проигрывают втрое: изоляция стоит сериализации, и на больших массивах эта цена больше выигрыша от лишних ядер.

Asyncio в эту аналогию тоже ложится — это epoll с одним потоком, то есть событийная модель, знакомая по сетевым серверам на C.

6Вопросы пытливого ума

Вопросы, которые возникают сами, если читать внимательно. Ответ — под вопросом.

Потоки, процессы, asyncio — как выбрать за десять секунд?

Два вопроса подряд: где нагрузка проводит время и сколько данных надо передать.

где времяинструментпочему
ждём сеть или диск, задач сотниasyncioОжидание не требует ни GIL, ни потока ОС
ждём, но задач десятки и код синхронныйпотокиНе нужно переписывать под async
считаем внутри C-расширенияпотокиGIL отпущен, память общая — 1.9× на двух ядрах
считаем на чистом PythonпроцессыЕдинственный способ занять несколько ядер
считаем на Python, данных многопроцессы + shared_memoryИначе всё съест сериализация

Второй вопрос важнее, чем кажется. Процессы не разделяют память, и всё, что уходит воркеру, проходит через pickle. Замер из этого урока: круг в 40 МБ через сериализацию — 0.358 с. Если сама работа занимает 0.1 с, процессы проигрывают однопоточному варианту, сколько ядер ни дай.

потоки:   0.27 c
процессы: 0.45 c        та же работа в numpy, четыре задачи

Здесь потоки быстрее именно потому, что numpy отпускает GIL, а процессам пришлось передавать данные.

Отсюда практическое правило: процессы окупаются, когда на единицу переданных данных приходится много вычислений. Считаешь долго над маленьким входом — процессы. Гоняешь массивы туда-обратно ради дешёвой операции — не окупятся никогда.

Почему fork опасен, если в процессе есть потоки?

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

Поток, державший в момент fork внутреннюю блокировку аллокатора, логгера или пула соединений, в потомке не появится. Блокировка при этом скопировалась в состоянии «занята», и снять её теперь некому. Первая же попытка её взять — вечное ожидание.

Коварство в том, что падает не сразу и не всегда: зависит от того, что именно делал другой поток в момент форка. Симптом — процесс, который иногда молча зависает на старте воркера.

Именно поэтому в 3.8 macOS перевели на spawn по умолчанию, а в 3.12 добавили DeprecationWarning при форке многопоточного процесса; в 3.14 умолчанием на Linux стал forkserver.

import multiprocessing as mp
print(mp.get_start_method())          # что используется сейчас
print(mp.get_all_start_methods())     # что доступно на этой платформе

mp.set_start_method("spawn", force=True)   # явно и предсказуемо

Три метода и когда какой. fork — быстрый старт и копирование состояния родителя, годится только для процесса без потоков. spawn — чистый интерпретатор с нуля: безопасно, но старт дороже и всё передаваемое обязано быть сериализуемым. forkserver — компромисс: заранее форкается один чистый служебный процесс, от него уже форкаются воркеры.

Практический вывод: если в приложении есть хоть один фоновый поток — а он есть почти всегда, потому что его заводят логгер, метрики и http-клиент, — fork использовать нельзя.

Субинтерпретаторы заменят процессы?

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

Идея (PEP 554 и PEP 684): несколько интерпретаторов в одном процессе, у каждого свои модули, свои глобальные и — с 3.12 — свой GIL. Значит настоящий параллелизм по ядрам без отдельных процессов.

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

Что не выигрывается: объекты Python по-прежнему нельзя передавать между интерпретаторами — у каждого своя куча и свои счётчики ссылок. Обмен идёт через каналы с сериализацией или через разделяемую память. То есть главная боль процессов — цена передачи данных — остаётся.

Что мешает прямо сейчас: большинство C-расширений не поддерживают несколько интерпретаторов в процессе. Модуль с глобальным состоянием на уровне C (а таких много, включая исторически numpy) при загрузке во второй интерпретатор ведёт себя непредсказуемо. Требуется переход на multi-phase init (PEP 489), и это работа на стороне каждой библиотеки.

В 3.13 появился модуль interpreters в стандартной библиотеке, но экосистема догоняет медленно.

Практический вывод на сегодня: процессы остаются рабочим ответом для CPU-bound Python. Субинтерпретаторы стоит держать в виду как то, что через несколько релизов может сделать ProcessPoolExecutor с его сериализацией ненужным — но принимать решения по ним сейчас рано.

Итог. Выбор решается двумя вопросами, а не перебором: где нагрузка проводит время и сколько данных придётся передать. Ожидание с сотнями задач — asyncio: пятьдесят вызовов по десять миллисекунд дают 0.509 с последовательно против 0.011 с одновременно, в одном потоке. Счёт внутри C-расширения — потоки: лок отпущен, память общая, и они втрое обгоняют процессы на массиве в сорок мегабайт, потому что pickle туда-обратно стоит 0.358 с сам по себе. Счёт в байткоде — только процессы, и там надо считать объём передачи: выигрывают задачи, где работы много, а данных мало, иначе сериализация съедает всё. Субинтерпретаторы — четвёртый вариант с дешёвым стартом и слабой поддержкой расширений. И первый вопрос почти всегда решается замером за минуту, а не рассуждением.
Дальше — по желанию
контрфактуалC++, Go и JavaScriptОдин механизм вместо четырёх — и что за это платят

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

Go сделал горутины единственным механизмом: вытесняющий планировщик, настоящий параллелизм по ядрам, каналы для обмена. Выбирать не из чего, и в этом сила. Цена — сложный рантайм и то, что логические гонки остаются (детектор -race существует не зря).

JavaScript — противоположный полюс: один поток и цикл событий, точка. Проблема выбора отсутствует, но и счётную нагрузку деть некуда, кроме воркеров с обменом через сериализацию — ровно как процессы в Python, со всей той же ценой передачи.

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

🐛 багПроцессы там, где нужны были потокиПерешли на ProcessPoolExecutor ради «настоящего параллелизма» и получили втрое медленнее

Обработка массивов ускорялась потоками недостаточно, взяли процессы. Стало хуже, и объяснить никто не может.

from concurrent.futures import ProcessPoolExecutor

def process_batch(arr):          # arr — numpy-массив на 40 МБ
    return (arr * arr).sum()

with ProcessPoolExecutor(4) as ex:
    results = list(ex.map(process_batch, batches))
четыре задачивремя
потоки0.46 с
процессы1.24 с

Рассуждение было верным в общем и неверным в частности: GIL действительно мешает параллелить байткод, но здесь байткода почти нет — вся работа внутри numpy, который лок отпускает. Зато появилась сериализация: сорок мегабайт на задачу, туда и обратно, 0.358 с на один цикл — больше, чем сам расчёт.

Ревью пропускает, потому что замена одного исполнителя на другой — однострочный дифф с убедительным обоснованием.

Лечение — вернуться к потокам. А если данные всё же должны идти в процессы, не пересылать их:

from multiprocessing import shared_memory
import numpy as np

shm = shared_memory.SharedMemory(create=True, size=arr.nbytes)
view = np.ndarray(arr.shape, dtype=arr.dtype, buffer=shm.buf)
view[:] = arr[:]                  # один раз
# дочерние процессы открывают shm по имени — копирования нет

Работает это ровно на buffer protocol из урока 15: numpy делает массив поверх чужой памяти. Правило: прежде чем брать процессы, посчитай объём передачи и сравни с объёмом работы.

⚡ приёмОтветить на первый вопрос замеромТри прогона по минуте дают ответ надёжнее любого рассуждения о GIL

Вместо спора о том, поможет ли параллелизм, — короткий тест, который сразу указывает ветку дерева.

import time, threading
from concurrent.futures import ProcessPoolExecutor

def probe(fn, arg, n=4):
    t = time.perf_counter()
    for _ in range(n): fn(arg)
    seq = time.perf_counter() - t

    th = [threading.Thread(target=fn, args=(arg,)) for _ in range(n)]
    t = time.perf_counter()
    for x in th: x.start()
    for x in th: x.join()
    thr = time.perf_counter() - t

    t = time.perf_counter()
    with ProcessPoolExecutor(n) as ex: list(ex.map(fn, [arg]*n))
    prc = time.perf_counter() - t

    print(f"последовательно {seq:.2f} | потоки {thr:.2f} ({seq/thr:.1f}×) "
          f"| процессы {prc:.2f} ({seq/prc:.1f}×)")

Как читать результат. Потоки ускорили — берём потоки, вопрос закрыт. Потоки не ускорили, процессы ускорили — работа в байткоде, берём процессы и проверяем объём передачи. Ничего не ускорило — либо задача слишком мелкая для накладных расходов, либо упирается не в процессор.

Что обязательно проверить перед выводами: задача должна быть достаточно крупной, иначе меряешь накладные расходы; и она должна быть той же, что в проде, — на синтетике numpy отпускает лок, а на твоей комбинации операций может и нет.

Отдельная ветка, которую тест не покажет: если нагрузка — это сотни одновременных ожиданий, ни потоки, ни процессы не масштабируются по числу задач, и правильный ответ asyncio. Этот случай виден не по замеру, а по постановке.

💻 терминалПроверь в терминалеЧетыре фрагмента плюс fork с потоками, shared_memory и субинтерпретаторы

Прогони фрагменты, сверяясь с предсказанием. Затем:

import multiprocessing as mp
print(mp.get_start_method())
print(mp.get_all_start_methods())
# Какой умолчательный на твоей системе? Чем spawn отличается от fork по цене старта?

import os, threading
def t(): 
    import time; time.sleep(5)
threading.Thread(target=t).start()
pid = os.fork()
# Почему так делать нельзя? Что происходит с локами в потомке?

from concurrent.futures import ThreadPoolExecutor
import time
def io_task(_): time.sleep(0.1)
# Прогони 100 задач в пуле на 4 и на 50 воркеров. Где предел полезности?

import sys
print(sys.version_info >= (3, 12))
# На 3.12+ посмотри модуль _xxsubinterpreters — сколько стоит старт субинтерпретатора?

Третий вопрос практичнее прочих: на I/O число воркеров можно поднимать далеко за число ядер, и понимание где именно перестаёт помогать — половина настройки любого пула.

исходникиconcurrent.futures и multiprocessingОдин API поверх двух моделей — и где видно, что модели разные

Lib/concurrent/futures/thread.py и Lib/concurrent/futures/process.py. Стоит прочитать оба подряд: интерфейс одинаковый, внутренности принципиально разные. Поток берёт задачу из очереди в общей памяти; процесс — сериализует аргументы, отправляет в канал, получает обратно результат. Вся цена из следствия 2 видна в этой разнице.

Заметь в процессной версии обработку случая, когда аргумент не пиклится: сообщение об ошибке приходит не из дочернего процесса, а из родительского при попытке отправить. Отсюда загадочные PicklingError на лямбдах и локальных классах (урок 06).

Lib/multiprocessing/shared_memory.py. Тонкая обёртка над системным вызовом, экспортирующая буфер. Всё остальное делает buffer protocol — тот же механизм, что связывает numpy с torch.

Lib/asyncio/base_events.py, метод run_in_executor. Мостик между моделями: он позволяет асинхронному сервису отдать блокирующую работу в пул потоков или процессов. Три строки, снимающие ограничение из урока 14.

L2Цена каждого механизма в цифрахСтоимость старта, переключения и передачи данных

Из L1 ты знаешь дерево. Здесь — порядки величин, по которым оно считается.

механизмстарт исполнителяпереключениепередача данных
корутинамикросекундыдесятки наносекундбесплатно, общая память
потокдесятки микросекундмикросекунды, через ядробесплатно, общая память
процессдесятки миллисекундмикросекундысериализация: 0.358 с на 40 МБ
субинтерпретатормиллисекундымикросекундысериализация через каналы

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

Про способ старта процесса. fork дешёв, потому что копирует адресное пространство лениво, — но несовместим с потоками и наследует всё состояние родителя, включая открытые соединения. spawn запускает интерпретатор заново и импортирует главный модуль (урок 17, отсюда нужда в if __name__), поэтому дороже, но предсказуем. Умолчание менялось между версиями и платформами — проверяй get_start_method, а не память.

import multiprocessing as mp
print((mp.get_start_method(), mp.get_all_start_methods()))
# → на Linux обычно 'fork' и список из трёх; на macOS 'spawn'
L3Где это в исходниках CPythonfutures, multiprocessing, PEP 684

Тег v3.11.15.

  • Lib/concurrent/futures/process.py — весь конвейер: сериализация, очереди, управление воркерами. Видно, где именно платится цена из следствия 2.
  • Lib/multiprocessing/context.py — способы старта и их различия.
  • Modules/_threadmodule.c — низкоуровневые потоки и примитивы синхронизации.
  • PEP 684 — отдельный лок на субинтерпретатор; PEP 554 — их API.
связиКуда это ведётL14, L15, L19 · контраст с горутинами Go и однопоточным JS
корень R3 · Подсчёт ссылок как модель памяти
← основа L19 · GIL — почему потоки не дают параллелизма на байткоде
← основа L15 · Граница с C — почему дают на numpy, и что такое shared_memory
← основа L14 · async/await — механизм для ветки «много ожиданий»
↔ контраст Go: один механизм вместо четырёх, выбирать не из чего
↔ контраст JavaScript: одна модель и воркеры с сериализацией, как процессы здесь
→ дальше L21 · Производительность — как понять, что вообще упёрлось в процессор
источникиЧто почитатьconcurrent.futures 🟧 · PEP 684 🟦 · process.py 🟥