Deep Engineering

ЗАМЕР

bench/async-vs-sync/blocking.py

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

Цитируется в статье
/ru/python/concurrency/async-vs-sync

Запись прогона

Асинхронность против синхронности: замеры для статьи

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

Все замеры ходят в один Postgres и в одну таблицу; параметры собраны в env.py, чтобы «прочие равные» были действительно равными.

Как запустить

# 1. Поднять Postgres
/usr/lib/postgresql/16/bin/postgres -D <каталог данных> -p 5433 -k /tmp

# 2. Завести базу и таблицу
psql -h /tmp -p 5433 -U postgres -c "create database bench_async"
psql -h /tmp -p 5433 -U postgres -d bench_async -c "
  create table items (id int primary key, payload text not null);
  insert into items select g, repeat('x', 64) from generate_series(1, 100000) g;
  analyze items;"

# 3. Драйверы
python3.13 -m pip install 'psycopg[binary]' psycopg_pool asyncpg

# 4. Замеры
python3.13 pool.py       # потолок пула
python3.13 latency.py    # цена одной операции
python3.13 executor.py   # синхронный драйвер внутри async
python3.13 waiting.py    # цена одновременного ожидания
python3.13 blocking.py   # счёт рядом с запросами и хвост задержек

Адрес базы переопределяется переменной BENCH_DSN.

Скрипт Что меряет
env.py общие параметры и печать версий — не замер, а гарантия сравнимости
pool.py одна и та же работа при пуле 1…32 в трёх моделях: потоки + psycopg, async psycopg, asyncpg
latency.py время одного короткого запроса по одному соединению, без одновременности
executor.py синхронный драйвер через run_in_executor против трёх остальных моделей
waiting.py цена держать N одновременных ожиданий: потоки против корутин, время и память
blocking.py распределение задержки, когда рядом с запросами выполняется счётный код

Что показали замеры (машина замеров: 2 vCPU, Postgres 16.13 локально)

  1. Потолок ставит пул соединений, а не модель исполнения. При одном и том же размере пула три модели дают одинаковое время в пределах шума — на пуле 8 это 0,43 / 0,45 / 0,43 с при «идеале» 0,38 с.
  2. На коротком запросе без одновременности async медленнее. Запрос по первичному ключу: 59 мкс синхронно против 70 мкс на async psycopg и 68 мкс на asyncpg — то есть на 16–19 % дороже. Ускорять здесь нечего: ждать параллельно нечего, остаётся только накладной расход.
  3. run_in_executor не хуже и не лучше остальных, потому что и он упирается в тот же пул: 0,434 с против 0,428 с у чистых потоков.
  4. Где async выигрывает по-настоящему — держать много ожиданий сразу. Четыре тысячи одновременных ожиданий: 0,31 с и +0 МБ RSS у корутин против 2,03 с и +61 МБ у потоков.
  5. Счётный код рядом с запросами портит обе модели. Ожидаемого «цикл событий встал, а потоки живут» не видно: GIL делает счёт последовательным в любой модели. Один кусок на 191 мс даёт максимум 184 мс у async и 212 мс у потоков.

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

Скрипт

189 строк
"""Что делает с задержками счётный кусок кода рядом с запросами к базе.

ЗАЧЕМ. Предыдущие замеры показали: на чистом ожидании модель исполнения не
решает почти ничего. Но приложение не состоит из одного ожидания: рядом с
запросом всегда есть работа процессора — сериализация ответа, шаблон, разбор
тела, хеширование пароля. Здесь и появляется разница, из-за которой переход на
async бывает не нейтральным, а ХУДШИМ.

МЕХАНИЗМ. У корутины нет способа отобрать управление: пока функция считает,
цикл событий не крутится, и ни одна готовая операция не будет обработана — даже
та, чей ответ от базы уже пришёл. У потока способ есть: интерпретатор отпускает
GIL каждые `sys.getswitchinterval()` секунд, и операционная система переключает
поток принудительно.

КАК МЕРЯЕТСЯ, И ПОЧЕМУ ИМЕННО ТАК. Заявки приходят ПО РАСПИСАНИЮ — по одной
каждые `INTERVAL` секунд, как в сервер приходят запросы, — и задержка считается
от НАЗНАЧЕННОГО времени заявки, а не от момента, когда до неё дошли руки.

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

ЧТО ЭТОТ ЗАМЕР НЕ ГОВОРИТ. Он не говорит, что потоки быстрее вообще. Он
показывает, что портится и насколько, когда рядом с ожиданием оказывается счёт.
"""

import asyncio
import statistics
import sys
import threading
import time
from concurrent.futures import ThreadPoolExecutor

from psycopg_pool import AsyncConnectionPool, ConnectionPool

sys.path.insert(0, __file__.rsplit("/", 1)[0])
import env  # noqa: E402

QUERIES = 240
SLEEP = 0.005  # ожидание на сервере, на запрос
POOL = 8
INTERVAL = 0.002  # заявка каждые 2 мс — 500 в секунду
CHUNK_EVERY = 12  # счётный кусок после каждой N-й заявки — сценарий «часто и понемногу»
BIG_AT = QUERIES // 4  # один длинный кусок посередине — сценарий «редко и надолго»
REPEATS = 3


def burn(iterations: int) -> int:
    """Чистый счёт на Python: держит GIL, в цикл событий не возвращается."""
    total = 0
    for i in range(iterations):
        total += i * i
    return total


def calibrate(target: float = 0.020) -> tuple[int, float]:
    """Сколько итераций счёта занимают `target` секунд ИМЕННО НА ЭТОЙ машине.

    Калибровка, а не константа: число итераций, подобранное на одной машине, на
    другой означало бы другое время, и замер сравнивал бы разное.
    """
    n = 20_000
    while True:
        started = time.perf_counter()
        burn(n)
        elapsed = time.perf_counter() - started
        if elapsed > target * 0.9:
            n = int(n * target / elapsed)
            break
        n *= 2
    started = time.perf_counter()
    burn(n)
    return n, time.perf_counter() - started


def run_threads(iterations: int, once: bool = False) -> tuple[float, list[float]]:
    latencies: list[float] = []
    lock = threading.Lock()
    with ConnectionPool(env.DSN, min_size=POOL, max_size=POOL) as pool:
        pool.wait()

        def one(due: float, index: int) -> None:
            with pool.connection() as conn:
                with conn.cursor() as cur:
                    cur.execute("select pg_sleep(%s)", (SLEEP,))
                    cur.fetchone()
            # Счёт — в том же потоке, что и запрос: так и выглядит приложение,
            # где выбранное надо ещё сериализовать.
            if iterations and (index % CHUNK_EVERY == 0 if not once else index == BIG_AT):
                burn(iterations)
            with lock:
                latencies.append(time.perf_counter() - due)

        with ThreadPoolExecutor(max_workers=POOL) as ex:
            begin = time.perf_counter()
            futures = []
            for i in range(QUERIES):
                due = begin + i * INTERVAL
                now = time.perf_counter()
                if due > now:
                    time.sleep(due - now)
                futures.append(ex.submit(one, due, i))
            for f in futures:
                f.result()
            return time.perf_counter() - begin, latencies


async def _run_async(iterations: int, once: bool = False) -> tuple[float, list[float]]:
    latencies: list[float] = []
    pool = AsyncConnectionPool(env.DSN, min_size=POOL, max_size=POOL, open=False)
    await pool.open(wait=True)
    try:

        async def one(due: float, index: int) -> None:
            async with pool.connection() as conn:
                async with conn.cursor() as cur:
                    await cur.execute("select pg_sleep(%s)", (SLEEP,))
                    await cur.fetchone()
            if iterations and (index % CHUNK_EVERY == 0 if not once else index == BIG_AT):
                burn(iterations)
            latencies.append(time.perf_counter() - due)

        begin = time.perf_counter()
        tasks = []
        for i in range(QUERIES):
            due = begin + i * INTERVAL
            now = time.perf_counter()
            if due > now:
                await asyncio.sleep(due - now)
            tasks.append(asyncio.create_task(one(due, i)))
        await asyncio.gather(*tasks)
        return time.perf_counter() - begin, latencies
    finally:
        await pool.close()


def percentile(values: list[float], q: float) -> float:
    ordered = sorted(values)
    return ordered[min(int(q * len(ordered)), len(ordered) - 1)]


def report(label: str, wall: float, lat: list[float]) -> None:
    print(
        f"{label:>22} | {wall:>7.2f}с | {statistics.median(lat) * 1000:>8.1f} | "
        f"{percentile(lat, 0.95) * 1000:>7.1f} | {percentile(lat, 0.99) * 1000:>7.1f} | "
        f"{max(lat) * 1000:>7.1f}"
    )


def main() -> None:
    env.describe()
    iterations, chunk_seconds = calibrate()
    print(f"switchinterval  {sys.getswitchinterval() * 1000:.0f} мс")
    print(f"счётный кусок   {chunk_seconds * 1000:.1f} мс ({iterations} итераций)")
    print(
        f"поток заявок    {QUERIES} штук, по одной каждые {INTERVAL * 1000:.0f} мс "
        f"({1 / INTERVAL:.0f} в секунду)"
    )
    print(f"запрос          ожидание {SLEEP * 1000:.0f} мс на сервере, пул {POOL} соединений")
    print(f"задержка        от назначенного времени заявки; лучшее из {REPEATS} по общему времени\n")

    print(
        f"{'модель':>22} | {'всего':>8} | {'медиана':>8} | {'p95':>7} | {'p99':>7} | {'макс':>7}"
    )
    print("-" * 76)

    big, big_seconds = calibrate(0.200)

    scenarios = [
        (0, False, "без счёта"),
        (iterations, False, f"счёт {chunk_seconds * 1000:.0f}мс×{QUERIES // CHUNK_EVERY}"),
        (big, True, f"один кусок {big_seconds * 1000:.0f}мс"),
    ]
    for iters, once, label in scenarios:
        thr = min((run_threads(iters, once) for _ in range(REPEATS)), key=lambda r: r[0])
        asy = min(
            (asyncio.run(_run_async(iters, once)) for _ in range(REPEATS)), key=lambda r: r[0]
        )
        report(f"потоки, {label}", *thr)
        report(f"async, {label}", *asy)
        print("-" * 76)


if __name__ == "__main__":
    main()