ЗАМЕР
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 локально)
- Потолок ставит пул соединений, а не модель исполнения. При одном и том же размере пула три модели дают одинаковое время в пределах шума — на пуле 8 это 0,43 / 0,45 / 0,43 с при «идеале» 0,38 с.
- На коротком запросе без одновременности async медленнее. Запрос по первичному ключу: 59 мкс синхронно против 70 мкс на async psycopg и 68 мкс на asyncpg — то есть на 16–19 % дороже. Ускорять здесь нечего: ждать параллельно нечего, остаётся только накладной расход.
run_in_executorне хуже и не лучше остальных, потому что и он упирается в тот же пул: 0,434 с против 0,428 с у чистых потоков.- Где async выигрывает по-настоящему — держать много ожиданий сразу. Четыре тысячи одновременных ожиданий: 0,31 с и +0 МБ RSS у корутин против 2,03 с и +61 МБ у потоков.
- Счётный код рядом с запросами портит обе модели. Ожидаемого «цикл событий встал, а потоки живут» не видно: 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()