Deep Engineering

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)