Deep Engineering

ЗАМЕР

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

Скрипт

190 строк
"""Перешардирование средствами PostgreSQL: добавить девятый шард к восьми.

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

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

  1. Можно ли добавить девятую секцию с модулем 9 к восьми с модулем 8.
     Печатается ответ сервера.
  2. Что можно сделать вместо: разрезать один шард надвое, перейдя для него
     к модулю 16. Строки второй половины переносятся в новый, девятый шард
     через postgres_fdw — тем же путём, каким их переносил бы оператор.
     Считается, сколько строк уехало, и что показывает count(*) по всей
     таблице в промежутке между копированием и удалением.
  3. Сколько строк пришлось бы перевезти ради ровной раскладки на девять
     шардов — посчитано собственной хеш-функцией PostgreSQL, строка за
     строкой, а не моделью.

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

from __future__ import annotations

import psycopg

from stand import SHARD_TABLE, Stand, dsn_for, recreate, version


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


def attempt(conn, sql: str) -> list[str]:
    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}")
        return lines


SHARDS = 8


def main() -> None:
    print(f"{version()} | adding a ninth shard to eight")
    stand = Stand(SHARDS, 0.0, users=20_000, prefix="shr_")
    stand.build()
    total = sum(stand.rows_per_shard())
    print(f"  {SHARDS} shards behind postgres_fdw, orders partitioned by hash of user_id, {total} rows")

    # Девятый шард: пустая база с той же таблицей и внешний сервер к ней.
    new_db = "shr_8"
    recreate(new_db)
    with psycopg.connect(dsn_for(new_db), autocommit=True) as nc:
        nc.execute(SHARD_TABLE)
    with stand.connect() as conn:
        conn.execute(
            f"create server s8 foreign data wrapper postgres_fdw"
            f" options (host '{stand.host}', port '{stand.port}', dbname '{new_db}', batch_size '1000')"
        )
        conn.execute("create user mapping for current_user server s8 options (user 'postgres')")

    show("1. A NINTH PARTITION WITH MODULUS 9 NEXT TO EIGHT WITH MODULUS 8")
    with stand.connect() as conn:
        sql = (
            "create foreign table orders_8 partition of orders"
            " for values with (modulus 9, remainder 8) server s8 options (table_name 'orders')"
        )
        print("    create foreign table orders_8 partition of orders")
        print("      for values with (modulus 9, remainder 8) server s8 ...")
        for line in attempt(conn, sql):
            print(f"      {line}")
    print()
    print("  PostgreSQL routes a row to hash mod modulus, and it only accepts moduli")
    print("  that divide one another: 8 and 16 can coexist, 8 and 9 cannot. With one")
    print("  partition per shard there is no way to add one equal ninth shard; either")
    print("  one shard is split in two, or every row is placed again under a new")
    print("  modulus.")

    show("2. SPLITTING SHARD 0 IN TWO: MODULUS 8 BECOMES 16 FOR THAT SHARD ONLY")
    with stand.connect() as conn:
        # Шард 0 перестаёт быть секцией с модулем 8 и становится двумя
        # секциями с модулем 16: остаток 0 остаётся на месте, остаток 8
        # уезжает в новый шард. Старая внешняя таблица orders_0 остаётся
        # отсоединённой — через неё видны строки, ещё лежащие в шарде 0.
        conn.execute("alter table orders detach partition orders_0")
        conn.execute(
            "create foreign table orders_0_lo partition of orders"
            " for values with (modulus 16, remainder 0) server s0 options (table_name 'orders')"
        )
        conn.execute(
            "create foreign table orders_8 partition of orders"
            " for values with (modulus 16, remainder 8) server s8 options (table_name 'orders')"
        )
        # Какие строки уезжают, решается У КООРДИНАТОРА. Условие
        # satisfies_hash_partition(<oid таблицы>, ...) postgres_fdw честно
        # считает переносимым и отправляет в шард — а шард о таблице с этим
        # OID ничего не знает: «could not open relation with OID». Поэтому
        # строки шарда 0 сперва забираются целиком (materialized не даёт
        # планировщику протолкнуть условие внутрь), и решение принимается
        # здесь.
        conn.execute(
            "create temp table staged as"
            " with raw as materialized (select * from orders_0)"
            " select *, satisfies_hash_partition('orders'::regclass, 16, 8, user_id) as leaving"
            " from raw"
        )
        before_rows = conn.execute("select count(*) from staged").fetchone()[0]
        copied = conn.execute(
            "insert into orders select user_id, order_id, email, amount from staged where leaving"
        ).rowcount
        seen_during = conn.execute("select count(*) from orders").fetchone()[0]
        deleted = conn.execute(
            "delete from orders_0 o using staged s"
            " where s.leaving and o.user_id = s.user_id and o.order_id = s.order_id"
        ).rowcount
        seen_after = conn.execute("select count(*) from orders").fetchone()[0]
    per_shard = stand.rows_per_shard()
    with psycopg.connect(dsn_for(new_db), autocommit=True) as nc:
        per_shard.append(nc.execute("select count(*) from orders").fetchone()[0])
    print(f"  {'rows in shard 0 before the split':<46} {before_rows:>8}")
    print(f"  {'rows copied to the new shard 8':<46} {copied:>8}")
    print(f"  {'share of all rows that travelled':<46} {copied / total:>7.1%}")
    print(f"  {'count(*) over the table before':<46} {total:>8}")
    print(f"  {'count(*) after copying, before deleting':<46} {seen_during:>8}")
    print(f"  {'rows deleted from shard 0':<46} {deleted:>8}")
    print(f"  {'count(*) after deleting':<46} {seen_after:>8}")
    print(f"  rows per shard after, shards 0..8: {per_shard}")
    print(f"  {'largest / smallest shard':<46} {max(per_shard) / min(per_shard):>7.2f}x")
    print()
    print("  The split moves only half of one shard, and nothing else is touched.")
    print("  The price is the layout: two shards now hold half as much as the other")
    print("  seven. And for as long as the moved rows exist in both places, the")
    print("  table counts them twice - copy and delete are two steps, and a reader")
    print("  between them sees both copies.")

    show("3. WHAT AN EVEN NINE-SHARD LAYOUT WOULD MOVE, BY POSTGRESQL'S OWN HASH")
    with stand.connect() as conn:
        # Для каждой строки: её секция при модуле 8 и при модуле 9, по той же
        # функции, которой PostgreSQL раскладывает строки. Шард i при модуле
        # 8 и шард i при модуле 9 — один и тот же шард; строка остаётся на
        # месте, только если номер совпал.
        def owner(mod: int) -> str:
            cases = " ".join(
                f"when satisfies_hash_partition('orders'::regclass, {mod}, {r}, user_id) then {r}"
                for r in range(mod)
            )
            return f"case {cases} end"

        stays, rows = conn.execute(
            f"select count(*) filter (where ({owner(8)}) = ({owner(9)})), count(*) from orders"
        ).fetchone()
    moved = rows - stays
    print(f"  {'rows':<46} {rows:>8}")
    print(f"  {'rows whose shard changes, modulus 8 -> 9':<46} {moved:>8}")
    print(f"  {'share':<46} {moved / rows:>7.1%}")
    print(f"  {'the new shard fair share, 1 / 9':<46} {1 / 9:>7.1%}")
    print()
    print("  PostgreSQL's hash partitioning is modulo placement, and it pays what")
    print("  modulo pays in bench/sharding/placement.py: to spread evenly over nine,")
    print("  eight rows in nine have to change shard, where one in nine would do.")

    stand.drop()
    with psycopg.connect(stand_admin_dsn(), autocommit=True) as ac:
        ac.execute(
            "select pg_terminate_backend(pid) from pg_stat_activity where datname = %s", (new_db,)
        )
        ac.execute(f'drop database if exists "{new_db}"')


def stand_admin_dsn() -> str:
    from stand import DSN

    return DSN


if __name__ == "__main__":
    main()