Deep Engineering

MEASUREMENT

bench/sharding/placement.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

327 lines
"""Шардирование: что переезжает при смене состава и где скапливается нагрузка.

ЗАЧЕМ ЭТОТ ФАЙЛ. Стратегии раскладки данных по шардам обычно сравнивают
словами: «диапазоны удобны для выборок», «хеш распределяет равномерно»,
«кольцо переносит мало». Все три фразы верны и все три недоговорены: каждая
стратегия платит за своё достоинство чем-то другим, и цену видно только в
числах. Здесь эти числа считаются.

ЧТО ЗДЕСЬ СЧИТАЕТСЯ. Это вычисление, а не симуляция: раскладка ключей по узлам
получается точно, и доли — тоже. Случайны только сами ключи, и они
фиксированы зерном.

  1. Одно и то же изменение состава — восемь узлов становятся девятью — при
     пяти стратегиях: сколько ключей сменило владельца и насколько неровной
     стала раскладка после.
  2. Ровная раскладка КЛЮЧЕЙ против ровной НАГРУЗКИ: доступ к ключам
     неравномерен (закон Ципфа), и шард с популярным ключом получает больше
     остальных. Доля самого ключа — нижняя граница для любой стратегии;
     сверх неё перекос зависит от того, куда легли остальные ключи.
  3. Возрастающий ключ и диапазоны: куда уходят новые записи.
  5. Ещё два правила без таблицы — rendezvous (HRW) и jump consistent hash —
     на том же изменении состава, что в блоке 1. Добавлены 02.10.2026 после
     внешнего аудита: вывод «мало перевозить и остаться ровными умеют только
     слоты» был сделан по пяти стратегиям из блока 1, а не по всем
     существующим, и эти две его опровергают.

ТЕ ЖЕ КЛЮЧИ И ТОТ ЖЕ ХЕШ, ЧТО В bench/hashring/ring.py. Ключи, их число, зерно
и функция хеша берутся оттуда импортом, а не копией. Иначе числа кольца здесь
и в уроке про консистентное хеширование разошлись бы при первой правке одного
из файлов, и ни одна проверка этого бы не заметила.

ЗАПУСК: python3 bench/sharding/placement.py
Вывод: runs/placement.txt
"""

import os
import sys

HERE = os.path.dirname(os.path.abspath(__file__))
sys.path.insert(0, os.path.join(HERE, "..", "hashring"))

from ring import KEYS, NODES, Ring, digest, keys  # noqa: E402

SPACE = 1 << 64  # хеш — первые восемь байт SHA-1, то есть число из [0, 2^64)


def show(title: str) -> None:
    print()
    print(title)
    print("-" * len(title))


def moved(before: list, after: list) -> float:
    return sum(1 for a, b in zip(before, after) if a != b) / len(before)


def spread(owners: list, nodes: int) -> tuple[float, float, float]:
    """Наименьшая доля, наибольшая доля и их отношение."""
    counts = [0] * nodes
    for owner in owners:
        counts[owner] += 1
    total = len(owners)
    lo, hi = min(counts) / total, max(counts) / total
    return lo, hi, hi / lo


# ------------------------------------------------------------ стратегии
def by_modulo(h: int, nodes: int) -> int:
    return h % nodes


def by_even_ranges(h: int, nodes: int) -> int:
    """Хеш-пространство поделено на nodes равных отрезков."""
    return h * nodes // SPACE


def by_split_range(h: int) -> int:
    """Восемь равных отрезков, и новый узел забирает верхнюю половину нулевого.

    Так добавляют узел в схему с диапазонами, когда перестраивать все границы
    не хотят: режут один отрезок пополам, остальные не трогают.
    """
    owner = h * NODES // SPACE
    if owner == 0 and h >= SPACE // (2 * NODES):
        return NODES  # новый, девятый узел
    return owner


def slot_table(slots: int, nodes: int) -> list[int]:
    """Слоты раздаются узлам сплошными кусками: узел i держит i-й кусок."""
    return [s * nodes // slots for s in range(slots)]


def grow_slot_table(table: list[int], old: int) -> list[int]:
    """Новый узел забирает у каждого старого поровну, остальное не трогается.

    Ровно так переносят слоты в кластере Valkey: целиком, по одному, и только
    те, что уходят новому узлу.
    """
    slots = len(table)
    target = slots // (old + 1)  # сколько слотов должен получить новый узел
    new = list(table)
    per_node = [[s for s in range(slots) if table[s] == n] for n in range(old)]
    taken = 0
    i = 0
    while taken < target:
        node = i % old
        own = per_node[node]
        # берётся последний ещё не отданный слот этого узла
        slot = own.pop()
        new[slot] = old
        taken += 1
        i += 1
    return new


def by_slot(h: int, table: list[int]) -> int:
    return table[h % len(table)]


def by_rendezvous(key: str, nodes: int) -> int:
    """Rendezvous, или HRW: у каждого узла своя оценка ключа, владелец — максимум.

    Оценка — тот же SHA-1, что везде в файле, от пары «ключ, узел». Добавленный
    узел забирает ровно те ключи, где его оценка выше всех прежних, и больше
    ничего не переезжает. Цена — поиск владельца перебирает все узлы.
    """
    return max(range(nodes), key=lambda n: digest(f"{key}|node-{n}"))


def by_jump(h: int, nodes: int) -> int:
    """Jump consistent hash (Lamping, Veach, 2014), дословно по статье.

    Таблицы нет, памяти не нужно, но узлы обязаны быть пронумерованы подряд:
    добавить можно только узел с номером N, убрать — только последний.
    """
    b, j = -1, 0
    key = h
    while j < nodes:
        b = j
        key = (key * 2862933555777941757 + 1) & 0xFFFFFFFFFFFFFFFF
        j = int((b + 1) * (float(1 << 31) / float((key >> 33) + 1)))
    return b


def main() -> None:
    print(f"Python {sys.version.split()[0]} · Linux {os.uname().release}")
    print("exact computation, not a simulation: only the keys are random")
    print(
        f"{KEYS} keys, {NODES} nodes, hash is the first 8 bytes of SHA-1"
        " (keys and hash imported from bench/hashring/ring.py)"
    )

    hashes = [digest(k) for k in keys()]
    names = [f"node-{i}" for i in range(NODES)]
    minimum = 1 / (NODES + 1)

    # ------------------------------------------------------------ блок 1
    show("1. THE SAME CHANGE UNDER FIVE STRATEGIES: 8 NODES BECOME 9")
    print(f"  {'strategy':<38} {'keys moved':>11} {'min share':>10} {'max share':>10} {'max/min':>8}")

    def report(label: str, before: list, after: list, nodes_after: int) -> None:
        lo, hi, ratio = spread(after, nodes_after)
        print(f"  {label:<38} {moved(before, after):>10.1%} {lo:>9.1%} {hi:>9.1%} {ratio:>7.2f}x")

    report(
        "modulo: hash mod N",
        [by_modulo(h, NODES) for h in hashes],
        [by_modulo(h, NODES + 1) for h in hashes],
        NODES + 1,
    )
    report(
        "ranges, all boundaries redrawn evenly",
        [by_even_ranges(h, NODES) for h in hashes],
        [by_even_ranges(h, NODES + 1) for h in hashes],
        NODES + 1,
    )
    report(
        "ranges, one range split in half",
        [by_even_ranges(h, NODES) for h in hashes],
        [by_split_range(h) for h in hashes],
        NODES + 1,
    )
    index = {name: i for i, name in enumerate(names + ["node-8"])}
    ring_before = Ring(names, 128)
    ring_after = Ring(names + ["node-8"], 128)
    report(
        "ring, 128 virtual nodes per node",
        [index[ring_before.owner_of_hash(h)] for h in hashes],
        [index[ring_after.owner_of_hash(h)] for h in hashes],
        NODES + 1,
    )
    for slots in (1024, 16384):
        table = slot_table(slots, NODES)
        grown = grow_slot_table(table, NODES)
        report(
            f"slots: {slots}, handed over whole",
            [by_slot(h, table) for h in hashes],
            [by_slot(h, grown) for h in hashes],
            NODES + 1,
        )
    print()
    print(f"  For the new node to get its fair share, at least {minimum:.1%} has to move;")
    print("  a scheme that moves less leaves the new node short of it.")
    print("  Modulo and evenly redrawn ranges stay even and move most of the data.")
    print("  Splitting one range moves the least and leaves two nodes with half the")
    print("  load of the rest. The ring moves about the minimum, but its arcs are")
    print("  still random: even with 128 points per node the busiest node holds")
    print("  noticeably more than the least busy one. Fixed slots move about the")
    print("  minimum and stay even, because the hand-over is decided slot by slot")
    print("  rather than left to where points fall; what they cost is a slot table")
    print("  every router has to hold and keep current.")

    # ------------------------------------------------------------ блок 2
    show("2. EVEN KEYS ARE NOT EVEN LOAD: ACCESS FOLLOWS A ZIPF LAW")
    owners = [by_modulo(h, NODES) for h in hashes]
    average = 1 / NODES
    counts = [owners.count(i) / KEYS for i in range(NODES)]
    print(f"  keys placed by hash mod 8: key shares per shard {min(counts):.1%} to {max(counts):.1%},")
    print(f"  max/min {max(counts) / min(counts):.2f}x - an even placement of keys")
    print("  key of rank r gets a share of requests proportional to 1 / r^s;")
    print("  'top key alone' is the top key's share over the average shard's")
    print(f"  {'s':>5} {'top key, share of all requests':>32} {'hottest shard':>14} {'vs average':>11} {'top key alone':>14}")
    hot_rows = {}
    for s in (0.8, 1.0, 1.2):
        weights = [1 / (r ** s) for r in range(1, KEYS + 1)]
        total = sum(weights)
        load = [0.0] * NODES
        for owner, w in zip(owners, weights):
            load[owner] += w / total
        top = weights[0] / total
        hottest = max(load)
        hot_rows[s] = (weights, total, load)
        print(f"  {s:>5} {top:>31.1%} {hottest:>13.1%} {hottest / average:>10.2f}x {top / average:>13.2f}x")
    print()
    print("  The shard that holds the most popular key carries that key's traffic")
    print("  whole: a key cannot be split between shards by any placement rule,")
    print("  because the rule maps a key to exactly one owner. The last column is")
    print("  therefore a floor for every rule; how far above it the hottest shard")
    print("  goes depends on where the other keys land, and that part is specific")
    print("  to this placement.")

    show("3. SPLITTING ONE HOT KEY INTO SUB-KEYS")
    print("  s = 1.2; the top key is stored as K copies key#0 .. key#K-1,")
    print("  each copy placed by its own hash, requests spread evenly over copies")
    print(f"  {'K':>5} {'shards the copies hit':>22} {'hottest shard':>15} {'vs average':>11}")
    weights, total, _ = hot_rows[1.2]
    top_key = keys()[0]
    rest = [0.0] * NODES
    for owner, w in list(zip(owners, weights))[1:]:
        rest[owner] += w / total
    for k in (1, 2, 4, 8, 16):
        load = list(rest)
        if k == 1:
            hit = {owners[0]}
            load[owners[0]] += weights[0] / total
        else:
            hit = set()
            for c in range(k):
                shard = by_modulo(digest(f"{top_key}#{c}"), NODES)
                hit.add(shard)
                load[shard] += weights[0] / total / k
        hottest = max(load)
        print(f"  {k:>5} {len(hit):>22} {hottest:>14.1%} {hottest / average:>10.2f}x")
    print()
    print("  More copies is not always better: the copies are placed by hash too, and")
    print("  a copy that lands on a shard already carrying other heavy keys can raise")
    print("  the peak.")
    print("  Splitting helps only as far as the copies land on different shards,")
    print("  and every read of the key now has to pick a copy - or read all of them,")
    print("  if the value is a counter that has to be summed.")

    # ------------------------------------------------------------ блок 4
    show("4. AN INCREASING KEY AND RANGES: WHERE NEW WRITES GO")
    existing = 100_000
    new = 10_000
    width = existing // NODES
    print(f"  ids 1..{existing} already stored, ranges of {width} ids each, 8 shards;")
    print(f"  the next {new} inserts get ids {existing + 1}..{existing + new}")
    print(f"  {'placement':<28} {'share of new writes on the busiest shard':>42}")

    def range_owner(i: int) -> int:
        return min((i - 1) // width, NODES - 1)  # последний отрезок открыт сверху

    new_ids = range(existing + 1, existing + new + 1)
    for label, owner in (
        ("ranges of the id", range_owner),
        ("hash of the id, mod 8", lambda i: by_modulo(digest(str(i)), NODES)),
    ):
        counts = [0] * NODES
        for i in new_ids:
            counts[owner(i)] += 1
        print(f"  {label:<28} {max(counts) / new:>41.1%}")
    print()
    print("  With ranges every new id is larger than all stored ones and lands in")
    print("  the last range: one shard takes every insert while seven take none.")
    print("  Hashing the id spreads inserts evenly - and gives up the ordered range")
    print("  scan that was the reason to use ranges in the first place.")

    # ------------------------------------------------------------ блок 5
    show("5. TWO MORE RULES WITHOUT A SLOT TABLE: RENDEZVOUS AND JUMP HASH")
    print("  the same change as in block 1: 8 nodes become 9, the same keys")
    print(f"  {'strategy':<38} {'keys moved':>11} {'min share':>10} {'max share':>10} {'max/min':>8}")
    names_list = keys()
    report(
        "rendezvous (HRW), score per node",
        [by_rendezvous(k, NODES) for k in names_list],
        [by_rendezvous(k, NODES + 1) for k in names_list],
        NODES + 1,
    )
    report(
        "jump consistent hash",
        [by_jump(h, NODES) for h in hashes],
        [by_jump(h, NODES + 1) for h in hashes],
        NODES + 1,
    )
    print()
    print("  Both move about the minimum and stay about as even as fixed slots,")
    print("  without a slot table. What they pay instead is different: rendezvous")
    print("  scores every node on every lookup, so a lookup costs N hashes; jump")
    print("  hash needs no memory at all, but the nodes must be numbered 0..N-1,")
    print("  so only the last node can be removed and an arbitrary one cannot.")


if __name__ == "__main__":
    main()