Deep Engineering

ЗАМЕР

bench/sharding/stand.py

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

Цитируется в статье
/ru/system-design/data-scaling/sharding
Прогон
PostgreSQL 16.13, один сервер, Linux 6.18

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

Замеры: шардирование — раскладка, маршрут запроса, операции через шарды

Четыре скрипта. Первый — точное вычисление без сервера, три остальных — против живого 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.

Скрипт

292 строк
"""Стенд шардирования на 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}"')