Deep Engineering

MEASUREMENT

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

172 lines
"""Операции, которые задевают больше одного шарда: уникальность и фиксация.

ЗАЧЕМ ЭТОТ СКРИПТ. Пока запрос идёт в один шард, шардированная таблица ведёт
себя как обычная. Два обещания обычной таблицы при этом молча перестают
действовать, стоит операции задеть второй шард: уникальность по столбцу, не
входящему в ключ шардирования, и атомарность транзакции. Здесь проверяется,
что именно PostgreSQL говорит в каждом из двух случаев.

ЧТО ИМЕННО ПРОВЕРЯЕТСЯ.

  1. Уникальный email на таблице, шардированной по user_id: на таблице с
     внешними секциями (шарды) и на обычной секционированной. Печатается
     ответ сервера дословно.
  2. Транзакция, вставившая по строке в два шарда, когда один из шардов
     отказывает при фиксации. Отказ устроен отложенным триггером, который
     срабатывает ровно на COMMIT, — то есть ПОСЛЕ того, как обе вставки
     прошли. Проверяется, что осталось в каждом шарде, при отказе каждого
     из двух и с parallel_commit и без него.
  3. Двухфазная фиксация на шарде: разрешена ли она по умолчанию.

ПОЧЕМУ ОТКАЗ ИМЕННО НА ФИКСАЦИИ. Ошибка во время вставки откатывает всё —
это неинтересно и честно. Опасный случай — когда всё прошло, и отказал один
участник на последнем шаге: postgres_fdw фиксирует удалённые транзакции по
очереди (или одновременно, с parallel_commit), и к моменту отказа второго
первый уже мог зафиксироваться.

ЗАПУСК:
    DE_BENCH_DSN="postgres://postgres@localhost:5433/postgres" \\
      python3 bench/sharding/crossshard.py
Вывод: runs/crossshard.txt
"""

from __future__ import annotations

import psycopg

from stand import Stand, version


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


def attempt(conn, sql: str) -> list[str]:
    """Выполнить и вернуть ответ сервера строками: OK или текст ошибки."""
    try:
        conn.execute(sql)
        return ["OK"]
    except psycopg.Error as e:
        lines = [f"ERROR:  {e.diag.message_primary}"]
        if e.diag.message_detail:
            lines.append(f"DETAIL:  {e.diag.message_detail}")
        if e.diag.message_hint:
            lines.append(f"HINT:  {e.diag.message_hint}")
        return lines


def main() -> None:
    print(f"{version()} | operations that touch more than one shard")
    print("  4 shards behind postgres_fdw, orders partitioned by hash of user_id")

    stand = Stand(4, 0.0, users=200, prefix="shx_")
    stand.build()

    show("1. A UNIQUE EMAIL ON A TABLE SHARDED BY THE USER ID")
    with stand.connect() as conn:
        print("  the sharded table: partitions are foreign tables on the shards")
        for sql in (
            "alter table orders add unique (email)",
            "alter table orders add unique (user_id, email)",
        ):
            print(f"    {sql}")
            for line in attempt(conn, sql):
                print(f"      {line}")
        conn.execute(
            "create table accounts (user_id bigint, email text) partition by hash (user_id)"
        )
        for i in range(4):
            conn.execute(
                f"create table accounts_{i} partition of accounts"
                f" for values with (modulus 4, remainder {i})"
            )
        print("  the same table with ordinary local partitions")
        for sql in (
            "alter table accounts add unique (email)",
            "alter table accounts add unique (user_id, email)",
        ):
            print(f"    {sql}")
            for line in attempt(conn, sql):
                print(f"      {line}")
    print()
    print("  Across shards there is no unique constraint at all. Even locally, a")
    print("  partitioned table can only enforce uniqueness that includes the")
    print("  partition key: each partition checks its own rows, and a duplicate")
    print("  email under another user_id lives in another partition.")

    # Пользователь, чьи строки лежат в шарде 0, и пользователь из шарда 1.
    owners = {}
    for i in (0, 1):
        with stand.shard_conn(i) as sc:
            owners[i] = sc.execute("select min(user_id) from orders").fetchone()[0]

    show("2. ONE TRANSACTION, TWO SHARDS, ONE SHARD REFUSES AT COMMIT")
    print(f"  insert one row for user {owners[0]} (shard 0) and one for user {owners[1]} (shard 1),")
    print("  then COMMIT; the refusing shard has a deferred trigger that raises at commit")
    print(f"  {'parallel_commit':<16} {'refusing shard':>15} {'client sees':>13} {'row on shard 0':>15} {'row on shard 1':>15}")
    outcomes = []
    for pc in ("false", "true"):
        stand.set_option("parallel_commit", pc)
        for refusing in (0, 1):
            for i in range(4):
                with stand.shard_conn(i) as sc:
                    sc.execute("drop trigger if exists refuse on orders")
                    sc.execute("drop function if exists refuse()")
                    sc.execute("delete from orders where order_id < 0")
            with stand.shard_conn(refusing) as sc:
                sc.execute(
                    "create function refuse() returns trigger language plpgsql as"
                    " $$ begin raise exception 'shard refuses at commit'; end $$"
                )
                sc.execute(
                    "create constraint trigger refuse after insert on orders"
                    " deferrable initially deferred for each row execute function refuse()"
                )
            seen = "committed"
            with stand.connect() as conn:
                try:
                    with conn.transaction():
                        conn.execute("insert into orders values (%s, -1, 'a@example.com', 1)", (owners[0],))
                        conn.execute("insert into orders values (%s, -2, 'b@example.com', 1)", (owners[1],))
                except psycopg.Error:
                    seen = "an error"
            present = []
            for i in (0, 1):
                with stand.shard_conn(i) as sc:
                    n = sc.execute("select count(*) from orders where order_id < 0").fetchone()[0]
                    present.append("kept" if n else "none")
            outcomes.append((pc, refusing, present))
            print(f"  {pc:<16} {refusing:>15} {seen:>13} {present[0]:>15} {present[1]:>15}")
    partial = sum(1 for _, _, p in outcomes if "kept" in p)
    print()
    print(f"  the client got an error every time; a row survived anyway in {partial} of {len(outcomes)} cases")
    print()
    print("  postgres_fdw commits each remote transaction in turn at the local commit.")
    print("  When the first shard to commit is the one that refuses, the rest are")
    print("  rolled back and nothing survives. When it is the other one, that one has")
    print("  already committed and stays committed. With parallel_commit every COMMIT")
    print("  goes out at once, so the shard that did not refuse has committed whichever")
    print("  shard refuses.")

    show("3. TWO-PHASE COMMIT ON A SHARD, DEFAULT SETTINGS")
    with stand.shard_conn(0) as sc:
        setting = sc.execute("show max_prepared_transactions").fetchone()[0]
        print(f"  max_prepared_transactions = {setting}")
        sc.execute("begin")
        sc.execute("insert into orders values (1, -9, 'c@example.com', 1)")
        print("    prepare transaction 'order-42'")
        for line in attempt(sc, "prepare transaction 'order-42'"):
            print(f"      {line}")
    print()
    print("  Two-phase commit is what would make block 2 atomic: every shard first")
    print("  promises to commit, and only then is anyone told to. PostgreSQL has it,")
    print("  turned off by default, and postgres_fdw does not use it.")

    stand.drop()


if __name__ == "__main__":
    main()