MEASUREMENT
bench/sharding/stand.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
292 lines"""Стенд шардирования на PostgreSQL: координатор, шарды и канал между ними.
ЭТО НЕ СКРИПТ, А ОБЩАЯ ЧАСТЬ трёх скриптов рядом: `routing.py`,
`crossshard.py`, `resharding.py`. Сам по себе он ничего не печатает.
КАК УСТРОЕН. Один сервер PostgreSQL, на нём N+1 баз:
shd_coord — координатор. Здесь живёт секционированная по хешу таблица
`orders`, и каждая её секция — внешняя таблица postgres_fdw,
указывающая в свой шард.
shd_0..N-1 — шарды. В каждом — обычная таблица `orders` со своей долей
строк.
Это штатный способ шардирования средствами самого PostgreSQL: секционирование
решает, КУДА идёт строка, postgres_fdw — КАК туда дойти. Специализированные
расширения делают ту же работу своим кодом; здесь нужен именно встроенный
механизм, потому что его поведение описано в документации дословно и его
можно сверить.
ПОЧЕМУ ВСЕ ШАРДЫ НА ОДНОМ СЕРВЕРЕ. Статья проверяет не производительность
железа, а три свойства: куда уходит запрос, сколько раз он ходит по сети и
что происходит с транзакцией, задевшей два шарда. Ни одно из них не зависит
от того, на одной машине шарды или на разных, — шард здесь отдельная база с
отдельными соединениями и отдельными транзакциями. А вот от задержки сети
зависит всё, поэтому её стенд даёт явно.
КАНАЛ. Координатор ходит к каждому шарду через свой ретранслятор с
задержкой — `DelayedLink` из bench/stretched-cache/link.py, тот же, что в
статье про растянутый кэш. На петлевом интерфейсе круг занимает десятки
микросекунд, и цена похода во все шарды на нём вышла бы обманчиво
маленькой. Ретранслятор применяет задержку в каждую сторону и считает
порции — это прибор, которым меряется число кругов.
ЗАГРУЗКА ИДЁТ МИМО КАНАЛА. Строки заливаются, пока внешние серверы указывают
прямо на порт PostgreSQL, и только потом серверы переключаются на
ретрансляторы. Загрузка — не предмет замера, и платить за неё задержкой
незачем.
"""
from __future__ import annotations
import os
import sys
import psycopg
HERE = os.path.dirname(os.path.abspath(__file__))
sys.path.insert(0, os.path.join(HERE, "..", "stretched-cache"))
from link import DelayedLink # noqa: E402
def decode(chunk: bytes) -> str:
"""Порция от координатора к шарду — словами.
Разбираются сообщения клиентского протокола PostgreSQL: тип (один байт),
длина (четыре), тело. Из простого запроса (Q) и из разбора (P) берётся
текст SQL, у остальных — только буква. Шифрования на петле нет: сервер
стенда поднят без ssl, и postgres_fdw при sslmode=prefer идёт открытым
текстом, иначе разбирать было бы нечего.
"""
parts, texts, i = [], [], 0
while i + 5 <= len(chunk):
kind = chr(chunk[i])
size = int.from_bytes(chunk[i + 1 : i + 5], "big")
body = chunk[i + 5 : i + 1 + size]
parts.append(kind)
if kind == "Q":
texts.append(body.rstrip(b"\0").decode(errors="replace"))
elif kind == "P":
fields = body.split(b"\0")
if len(fields) > 1:
texts.append(fields[1].decode(errors="replace"))
i += 1 + size
head = "+".join(parts) if parts else "?"
return f"{head}: {texts[0]}" if texts else head
class ClosingLink(DelayedLink):
"""DelayedLink, который действительно рвёт соединение, когда клиент ушёл.
В link.py конец потока закрывает оба сокета через close(). Но на
противоположном сокете в этот момент сидит в recv() читающий поток, и в
Linux close() такого сокета не отправляет FIN: ядро держит его открытым,
пока recv() не вернётся, а он не вернётся никогда. Шард не узнаёт, что
координатор ушёл, и его обслуживающий процесс остаётся висеть.
В статье про растянутый кэш этого не было видно: там соединения живут весь
прогон. Здесь каждый круг открывает новые, и за двадцать соединений
координатора висящих шардовых процессов набиралось двадцать на каждый
шард — прогон упирался в max_connections. shutdown() перед close()
отправляет FIN сразу и будит recv() на другой стороне.
"""
# Запись того, что координатор отправляет в шард: по порции — список
# сообщений протокола PostgreSQL. Включается на время одного запроса.
trace: list | None = None
def _reader(self, src, q, upstream: bool) -> None:
import time as _time
delay = self.delay_ms / 1000
try:
while True:
data = src.recv(65536)
if not data:
break
self.stats.add(upstream, len(data))
if upstream and self.trace is not None:
self.trace.append(decode(data))
q.put((_time.monotonic() + delay, data))
except OSError:
pass
finally:
q.put((_time.monotonic() + delay, self._EOF))
@staticmethod
def _shutdown(*socks) -> None:
import socket
for s in socks:
try:
s.shutdown(socket.SHUT_RDWR)
except OSError:
pass
try:
s.close()
except OSError:
pass
DSN = os.environ.get("DE_BENCH_DSN", "postgres://postgres@localhost:5433/postgres")
PREFIX = "shd_"
def dsn_for(db: str) -> str:
"""Тот же адрес, другая база."""
base = DSN.rsplit("/", 1)[0]
return f"{base}/{db}"
def server_addr() -> tuple[str, int]:
info = psycopg.conninfo.conninfo_to_dict(DSN)
return info.get("host", "localhost"), int(info.get("port", 5432))
def admin() -> psycopg.Connection:
return psycopg.connect(DSN, autocommit=True)
def version() -> str:
with admin() as conn:
return conn.execute("select version()").fetchone()[0].split(" on ")[0]
def recreate(db: str) -> None:
with admin() as conn:
conn.execute(
"select pg_terminate_backend(pid) from pg_stat_activity "
"where datname = %s and pid <> pg_backend_pid()",
(db,),
)
conn.execute(f'drop database if exists "{db}"')
conn.execute(f'create database "{db}"')
SHARD_TABLE = """
create table orders (
user_id bigint not null,
order_id bigint not null,
email text not null,
amount numeric(10,2) not null
);
create index on orders (user_id);
create index on orders (email);
"""
class Stand:
"""N шардов за координатором; секционирование по хешу user_id."""
def __init__(
self,
shards: int,
delay_ms: float = 0.0,
users: int = 20_000,
per_user: int = 10,
prefix: str = PREFIX,
):
self.n = shards
self.delay_ms = delay_ms
self.users = users
self.per_user = per_user
# Свой префикс у каждого скрипта: иначе два прогона, запущенные
# рядом, пересоздавали бы базы друг друга посреди замера.
self.coord = f"{prefix}coord"
self.shards = [f"{prefix}{i}" for i in range(shards)]
self.links: list[ClosingLink] = []
host, port = server_addr()
self.host, self.port = host, port
# ------------------------------------------------------------ сборка
def build(self) -> None:
for db in self.shards:
recreate(db)
with psycopg.connect(dsn_for(db), autocommit=True) as conn:
conn.execute(SHARD_TABLE)
recreate(self.coord)
with self.connect() as conn:
conn.execute("create extension postgres_fdw")
conn.execute(
"create table orders (user_id bigint not null, order_id bigint not null,"
" email text not null, amount numeric(10,2) not null)"
" partition by hash (user_id)"
)
for i, db in enumerate(self.shards):
conn.execute(
f"create server s{i} foreign data wrapper postgres_fdw"
f" options (host '{self.host}', port '{self.port}', dbname '{db}',"
" batch_size '1000')"
)
conn.execute(f"create user mapping for current_user server s{i} options (user 'postgres')")
conn.execute(
f"create foreign table orders_{i} partition of orders"
f" for values with (modulus {self.n}, remainder {i})"
f" server s{i} options (table_name 'orders')"
)
# Строки раскладывает сам PostgreSQL: вставка в секционированную
# таблицу уходит в ту секцию, которую выбирает его хеш-функция.
# Раскладывать строки вручную значило бы повторить эту функцию —
# и получить стенд, который проверяет сам себя.
conn.execute(
"insert into orders"
" select u, u * 100 + k, 'user' || u || '@example.com', (u % 97) + k"
f" from generate_series(1, {self.users}) u, generate_series(1, {self.per_user}) k"
)
for db in self.shards:
with psycopg.connect(dsn_for(db), autocommit=True) as conn:
conn.execute("analyze orders")
self.attach_links(self.delay_ms)
def attach_links(self, delay_ms: float) -> None:
"""Переключить внешние серверы на ретрансляторы с задержкой."""
self.delay_ms = delay_ms
self.links = [ClosingLink(self.port, delay_ms) for _ in self.shards]
with self.connect() as conn:
for i, link in enumerate(self.links):
conn.execute(f"alter server s{i} options (set port '{link.port}')")
def set_option(self, name: str, value: str) -> None:
"""Выставить параметр postgres_fdw на всех внешних серверах."""
with self.connect() as conn:
for i in range(self.n):
present = conn.execute(
"select 1 from pg_foreign_server, unnest(srvoptions) o"
" where srvname = %s and o like %s",
(f"s{i}", f"{name}=%"),
).fetchone()
verb = "set" if present else "add"
conn.execute(f"alter server s{i} options ({verb} {name} '{value}')")
# ------------------------------------------------------------ доступ
def connect(self) -> psycopg.Connection:
return psycopg.connect(dsn_for(self.coord), autocommit=True)
def shard_conn(self, i: int) -> psycopg.Connection:
return psycopg.connect(dsn_for(self.shards[i]), autocommit=True)
def rows_per_shard(self) -> list[int]:
out = []
for i in range(self.n):
with self.shard_conn(i) as conn:
out.append(conn.execute("select count(*) from orders").fetchone()[0])
return out
def portions(self) -> int:
"""Сколько порций ушло к шардам через все каналы с последнего сброса."""
return sum(link.stats.chunks_up for link in self.links)
def reset_portions(self) -> None:
for link in self.links:
link.reset_stats()
def drop(self) -> None:
for db in [self.coord, *self.shards]:
with admin() as conn:
conn.execute(
"select pg_terminate_backend(pid) from pg_stat_activity "
"where datname = %s and pid <> pg_backend_pid()",
(db,),
)
conn.execute(f'drop database if exists "{db}"')