MEASUREMENT
bench/sharding/routing.py
The script that produced the numbers in the article, and the record of the run. The file is read from the repository at build time — this is the code that was run, not a copy of it.
- Cited in
- /en/system-design/data-scaling/sharding
- Run on
- PostgreSQL 16.13, один сервер, Linux 6.18
The run below is recorded in Russian. It is a lab record, kept in the language it was written in; the numbers, the tables and the code read the same either way.
Record of the run
Замеры: шардирование — раскладка, маршрут запроса, операции через шарды
Четыре скрипта. Первый — точное вычисление без сервера, три остальных — против живого PostgreSQL, на котором шарды собраны штатными средствами: секционированием по хешу и postgres_fdw.
| скрипт | что показывает |
|---|---|
placement.py |
одно изменение состава — восемь узлов становятся девятью — при пяти стратегиях: сколько ключей переезжает и насколько неровной становится раскладка; неравномерный доступ (закон Ципфа) и популярный ключ; возрастающий ключ и диапазоны; блок 5 — ещё два правила без таблицы слотов, rendezvous (HRW) и jump consistent hash |
routing.py |
план запроса с ключом шардирования и без него; что координатор отправляет шарду, разобранное прямо из протокола; цена запроса без ключа по очереди и одновременно; сколько кругов оплачивается строго друг за другом |
crossshard.py |
уникальность по столбцу вне ключа шардирования; транзакция в два шарда, один из которых отказывает на фиксации; двухфазная фиксация по умолчанию |
resharding.py |
можно ли добавить девятый шард к восьми; что остаётся — разрезать один; что видит читатель, пока строки едут; сколько строк сдвинула бы ровная раскладка по хешу самого PostgreSQL |
stand.py — не скрипт, а общая часть трёх последних: собирает координатор и
шарды и ставит между ними канал с задержкой.
Запуск (для трёх последних нужен живой сервер; адрес — в DE_BENCH_DSN):
python3 bench/sharding/placement.py
DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \
python3 bench/sharding/routing.py
DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \
python3 bench/sharding/crossshard.py
DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \
python3 bench/sharding/resharding.py
Нужны psycopg 3 и расширение postgres_fdw (оно входит в стандартную поставку
PostgreSQL). Скрипты создают и удаляют собственные базы с префиксами shd_,
shx_ и shr_; чужие базы не трогаются.
Почему шарды — базы на одном сервере
Статья проверяет не производительность железа, а свойства механизма: куда уходит запрос, сколько раз он ходит к шарду, что происходит с транзакцией, задевшей два шарда, что сервер разрешает при смене состава. Ни одно из этих свойств не зависит от того, лежат шарды на одной машине или на разных: шард здесь — отдельная база с отдельными соединениями и отдельными транзакциями.
От задержки сети зависит всё, поэтому она задана явно. Координатор ходит к
каждому шарду через ретранслятор DelayedLink из
bench/stretched-cache/link.py — тот же, что в статье про растянутый кэш. На
петлевом интерфейсе круг занимает десятки микросекунд, и цена похода во все
шарды на нём была бы обманчиво маленькой.
Почему модель, а не стенд, для стратегий раскладки
Доля переехавших ключей — не величина, которую нужно мерить: она вычисляется
точно по раскладке. Стенд здесь добавил бы только шум. Поэтому placement.py
считает, а стенд проверяет то, что модель проверить не может: разрешает ли
сервер нужное изменение и что видно снаружи, пока оно идёт. Сверка между ними
есть — блок 3 resharding.py считает долю переехавших строк собственной
хеш-функцией PostgreSQL и получает то же, что модель для деления по модулю.
placement.py берёт ключи, их число и хеш из bench/hashring/ring.py импортом:
числа кольца здесь и в уроке про консистентное хеширование обязаны совпадать.
Протокол замера времени
routing.py меряет каждую конфигурацию девятью кругами по 30 запросов, печатает
медиану круга, лучший и худший круг и разброс; порядок конфигураций внутри круга
чередуется. Первый запрос в каждой конфигурации выбрасывается — в нём
устанавливаются соединения к шардам.
Задержка — параметр: 1 мс в каждую сторону как основная и 3 мс как вторая точка. По двум точкам считается число кругов, оплачиваемых строго друг за другом, — оно от задержки не зависит, и выводы сформулированы в нём.
Ошибки этих замеров, оставленные в истории
Первая редакция routing.py заканчивалась выводом, написанным до прогона:
при одновременном походе в шарды цена запроса без ключа «растёт с числом шардов
гораздо медленнее». Прогон показал, что она растёт почти так же, как при походе
по очереди. Разбор протокола объяснил почему: из пяти обменов с шардом
async_capable совмещает только выборку, а parallel_commit — только фиксацию;
открытие удалённой транзакции, объявление курсора и его закрытие по-прежнему
идут шард за шардом. Отсюда блоки 2 и 5 — трассировка и счёт последовательных
кругов.
Ретранслятор из статьи про растянутый кэш не рвал соединение до конца.
Закрытие через close() сокета, на котором другой поток сидит в recv(), в
Linux не отправляет FIN, и обслуживающие процессы шардов оставались висеть.
Там соединения жили весь прогон, и это не проявлялось; здесь каждый круг
открывает новые, и первый прогон упёрся в max_connections. stand.py
переопределяет закрытие: сначала shutdown(), потом close().
Условие переноса строк нельзя отдавать шарду. satisfies_hash_partition(<oid>, ...)
postgres_fdw считает переносимым и отправляет на шард, а шард о таблице с этим
OID не знает ничего. resharding.py решает, какие строки уезжают, у
координатора.
Вычитка нашла выводы шире вычисленного. Первая редакция статьи и
placement.py говорили, что шард с горячим ключом несёт 2,14 средней нагрузки
«при любой стратегии», и называли деление по модулю на восемь «самой ровной
раскладкой из блока 1». Ни того ни другого вычисление не показывало: 2,14
посчитано для одной раскладки, а в блоке 1 восьми узлов нет вовсе. Теперь блок
2 печатает ровность самой раскладки (12,3–12,7 %) и долю одного ключа против
средней доли шарда — 1,57 при s = 1,2. Эта доля и есть нижняя граница для
любой стратегии, а остальной перекос зависит от того, куда легли другие ключи.
Там же resharding.py утверждал, что равноправный девятый шард добавить
нельзя вообще. Это верно только при одной секции на шард, и оговорка теперь
стоит в выводе. Оба прогона повторены, и числа не изменились.
Аудит 02.10.2026: «только слоты» было выводом из пяти строк. Статья говорила, что мало перевозить и оставаться ровными одновременно умеют только фиксированные слоты. Блок 1 сравнивал пять стратегий, а не все существующие, и вывод был шире вычисленного. Теперь блок 5 считает на том же изменении состава ещё два правила без таблицы: rendezvous (HRW) — 11,3 % перевезено, перекос 1,03, — и jump consistent hash — 10,9 % и 1,04. Оба укладываются около минимума и ровны, как слоты; платят они другим — перебором всех узлов на каждый поиск и нумерацией узлов подряд. Блоки 1–4 при повторном прогоне не изменились.
Что получилось (PostgreSQL 16.13, один сервер, Linux 6.18)
Записи прогонов — в runs/: placement.txt, routing.txt, crossshard.txt,
resharding.txt.
Script
250 lines"""Куда уходит запрос к шардированной таблице и сколько он стоит.
ЗАЧЕМ ЭТОТ СКРИПТ. Про запрос без ключа шардирования говорят «он идёт во все
шарды». Это верно и недоговорено: во сколько раз это дороже, зависит не от
числа шардов, а от того, ходит ли координатор к ним по очереди или
одновременно, — и в PostgreSQL это одна настройка, выключенная по умолчанию.
ЧТО ИМЕННО МЕРЯЕТСЯ.
1. План: запрос с ключом шардирования и без него. Это не замер, а то, что
печатает сам PostgreSQL, — отсечение секций видно в плане буквально.
2. Сколько порций уходит к шардам за один запрос. Порция — единица, к
которой ретранслятор применяет задержку, то есть это число обменов.
2. Что координатор отправляет шарду: ретранслятор разбирает протокол и
печатает сообщения. Отсюда число обменов на шард.
3. Время запроса на канале с задержкой: с ключом, без ключа по очереди,
без ключа одновременно (async_capable), и то же с параллельной
фиксацией (parallel_commit). Для 4 и 8 шардов.
4. Сколько кругов оплачивается строго друг за другом. Каждая конфигурация
меряется при двух задержках; прирост времени, делённый на прирост
круга, — это число последовательных кругов. Так видно, какие из пяти
обменов async_capable и parallel_commit действительно совмещают.
ПРОТОКОЛ ЗАМЕРА. Один запрос — отдельная транзакция, как в приложении без
явных BEGIN. Соединения к шардам прогреты: первый запрос в каждой конфигурации
выбрасывается, иначе в замер попала бы установка соединения. Девять кругов по
30 запросов, ключи из фиксированного зерна; печатается медиана круга, лучший
и худший круг и разброс между ними. Порядок конфигураций внутри круга
чередуется, чтобы дрейф машины не ложился на одну из них.
ЗАДЕРЖКА — ПАРАМЕТР, А НЕ СВОЙСТВО МИРА. Основная — 1 мс в каждую сторону,
круг 2 мс: порядок задержки между стойками или зонами доступности одного
региона. Вторая, 3 мс, нужна только для того, чтобы посчитать число
последовательных кругов. Выводы сформулированы в кругах, а не в
миллисекундах, и переносятся на любую задержку умножением.
ОШИБКА ЭТОГО СКРИПТА, ОСТАВЛЕННАЯ В ИСТОРИИ. Первая редакция заканчивалась
фразой, написанной до прогона: при одновременном походе в шарды цена
запроса без ключа «растёт с числом шардов гораздо медленнее». Прогон её
опроверг: при переходе с 4 на 8 шардов она росла почти так же, как у похода
по очереди. Отсюда трассировка протокола (блок 2) и счёт
последовательных кругов (блок 6): они показывают, что async_capable
совмещает только выборку, а parallel_commit — только фиксацию.
ЗАПУСК:
DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \\
python3 bench/sharding/routing.py
Вывод: runs/routing.txt
"""
from __future__ import annotations
import random
import statistics
import sys
import time
from stand import Stand, version
ROUNDS = 9
PER_ROUND = 30
SEED = 20260923
def show(title: str) -> None:
print()
print(title)
print("-" * len(title))
KEY_Q = "select order_id, amount from orders where user_id = %s"
NONKEY_Q = "select order_id, amount from orders where email = %s"
def plan(conn, sql: str) -> list[str]:
return [r[0] for r in conn.execute("explain (costs off) " + sql).fetchall()]
def time_query(conn, sql: str, arg) -> float:
t = time.perf_counter()
conn.execute(sql, (arg,)).fetchall()
return (time.perf_counter() - t) * 1000
CONFIGS = [
("with the shard key", KEY_Q, {"async_capable": "false", "parallel_commit": "false"}),
("no key, shards in turn", NONKEY_Q, {"async_capable": "false", "parallel_commit": "false"}),
("no key, async_capable", NONKEY_Q, {"async_capable": "true", "parallel_commit": "false"}),
("no key, async + parallel_commit", NONKEY_Q, {"async_capable": "true", "parallel_commit": "true"}),
]
def measure(stand: Stand, users: int) -> dict[str, list[float]]:
rng = random.Random(SEED)
medians: dict[str, list[float]] = {label: [] for label, _, _ in CONFIGS}
for rnd in range(ROUNDS):
order = CONFIGS if rnd % 2 == 0 else list(reversed(CONFIGS))
for label, sql, opts in order:
for k, v in opts.items():
stand.set_option(k, v)
with stand.connect() as conn:
arg0 = 1 if sql is KEY_Q else "user1@example.com"
time_query(conn, sql, arg0) # прогрев соединений к шардам
times = []
for _ in range(PER_ROUND):
u = rng.randint(1, users)
arg = u if sql is KEY_Q else f"user{u}@example.com"
times.append(time_query(conn, sql, arg))
medians[label].append(statistics.median(times))
return medians
def main() -> None:
print(f"{version()} | routing a query to shards through postgres_fdw")
print(" shards are databases on one server; each sits behind its own delaying relay")
print(f" {ROUNDS} rounds x {PER_ROUND} queries, median of a round; seed {SEED}")
users = 20_000
stand = Stand(4, 0.0, users=users)
stand.build()
show("1. WHERE THE QUERY GOES: THE PLAN WITH AND WITHOUT THE SHARD KEY")
with stand.connect() as conn:
print(" where user_id = 42 (the table is partitioned by hash of user_id)")
for line in plan(conn, "select * from orders where user_id = 42"):
print(f" {line}")
print(" where email = 'user42@example.com'")
for line in plan(conn, "select * from orders where email = 'user42@example.com'"):
print(f" {line}")
stand.set_option("async_capable", "true")
print(" the same, with async_capable 'true' on the foreign servers")
for line in plan(conn, "select * from orders where email = 'user42@example.com'"):
print(f" {line}")
stand.set_option("async_capable", "false")
show("2. WHAT THE COORDINATOR SENDS TO A SHARD FOR ONE QUERY")
print(" decoded from the wire by the relay; connections already open")
with stand.connect() as conn:
conn.execute(NONKEY_Q, ("user41@example.com",)).fetchall() # прогрев
for link in stand.links:
link.trace = []
conn.execute(NONKEY_Q, ("user42@example.com",)).fetchall()
traces = [link.trace for link in stand.links]
for link in stand.links:
link.trace = None
print(" query without the shard key, shard 0:")
for n, message in enumerate(traces[0], 1):
first = message.splitlines()[0]
print(f" {n}. {first}")
per_shard = [len(t) for t in traces]
print(f" messages per shard, all shards: {per_shard}")
conn.execute(KEY_Q, (41,)).fetchall()
for link in stand.links:
link.trace = []
conn.execute(KEY_Q, (42,)).fetchall()
per_shard_key = [len(link.trace) for link in stand.links]
for link in stand.links:
link.trace = None
print(f" query with the shard key, messages per shard: {per_shard_key}")
print()
print(" Five exchanges per shard touched: open a remote transaction, declare a")
print(" cursor, fetch, close it, commit. A query without the key pays all five")
print(" on every shard.")
stand.drop()
delays = (1.0, 3.0)
results: dict[tuple[int, float], dict[str, list[float]]] = {}
for n in (4, 8):
stand = Stand(n, 0.0, users=users)
stand.build()
for d in delays:
stand.attach_links(d)
results[(n, d)] = measure(stand, users)
stand.drop()
block = 3
for n in (4, 8):
d = delays[0]
show(f"{block}. TIME PER QUERY, {n} SHARDS, {d:g} MS EACH WAY")
block += 1
r = results[(n, d)]
print(f" {'configuration':<34} {'median':>8} {'best round':>11} {'worst round':>12} {'spread':>7}")
for label, _, _ in CONFIGS:
m = r[label]
med = statistics.median(m)
print(
f" {label:<34} {med:>6.1f} ms {min(m):>8.1f} ms {max(m):>9.1f} ms"
f" {(max(m) - min(m)) / med:>6.0%}"
)
key = statistics.median(r["with the shard key"])
seq = statistics.median(r["no key, shards in turn"])
both = statistics.median(r["no key, async + parallel_commit"])
wins = sum(
1 for a, b in zip(r["no key, async + parallel_commit"], r["no key, shards in turn"]) if a < b
)
print()
print(f" no key in turn / with the key: {seq / key:.1f}x")
print(f" no key async + parallel_commit / with key: {both / key:.1f}x")
print(f" async + parallel_commit faster than in turn in {wins} of {ROUNDS} rounds")
show(f"{block}. HOW MANY ROUND TRIPS ARE PAID ONE AFTER ANOTHER")
block += 1
lo, hi = delays
print(f" each configuration timed at {lo:g} ms and at {hi:g} ms each way;")
print(f" (time at {hi:g} - time at {lo:g}) / ({2 * hi:g} - {2 * lo:g} ms) = round trips in sequence")
# Сколько кругов ДОЛЖНО получиться, если из пяти обменов блока 2
# совмещается только то, что обещает документация: выборка при
# async_capable и фиксация при parallel_commit. Совмещённый обмен стоит
# один круг на все шарды вместо одного на каждый, то есть экономит N - 1.
expected = {
"with the shard key": lambda n: 5,
"no key, shards in turn": lambda n: 5 * n,
"no key, async_capable": lambda n: 5 * n - (n - 1),
"no key, async + parallel_commit": lambda n: 5 * n - 2 * (n - 1),
}
print(f" {'configuration':<34} {'4 shards':>9} {'expected':>9} {'8 shards':>9} {'expected':>9}")
for label, _, _ in CONFIGS:
cells = []
for n in (4, 8):
t_lo = statistics.median(results[(n, lo)][label])
t_hi = statistics.median(results[(n, hi)][label])
cells.append((t_hi - t_lo) / (2 * (hi - lo)))
e4, e8 = expected[label](4), expected[label](8)
print(f" {label:<34} {cells[0]:>9.1f} {e4:>9} {cells[1]:>9.1f} {e8:>9}")
print()
print(" expected: of the five exchanges, fetch overlaps under async_capable and")
print(" commit under parallel_commit; an overlapped exchange costs one round trip")
print(" for all shards instead of one per shard, saving N - 1")
print()
print(" With the key the count is five - the five exchanges of block 2 - and it")
print(" does not depend on the number of shards. Without the key, in turn, it is")
print(" five per shard. async_capable overlaps the fetches and parallel_commit")
print(" the commits; opening the remote transaction, declaring the cursor and")
print(" closing it still happen shard after shard, so the count keeps growing")
print(" with the number of shards.")
show(f"{block}. WHAT GROWS WITH THE NUMBER OF SHARDS, {lo:g} MS EACH WAY")
print(f" {'configuration':<34} {'4 shards':>9} {'8 shards':>9} {'8 / 4':>7}")
for label, _, _ in CONFIGS:
a = statistics.median(results[(4, lo)][label])
b = statistics.median(results[(8, lo)][label])
print(f" {label:<34} {a:>6.1f} ms {b:>6.1f} ms {b / a:>6.2f}x")
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
sys.exit(1)