Шардирование: ключ решает, куда идёт каждый запрос
Разложить данные по узлам — половина дела; вторая половина — то, что после этого перестаёт работать. Здесь посчитано и измерено на живом PostgreSQL: при переходе с восьми шардов на девять деление по модулю перевозит 88,9 % ключей, а минимум — 11,1 %; на стенде с postgres_fdw запрос без ключа шардирования платит пять обменов с каждым шардом, и две настройки совмещают из них только два; транзакция через postgres_fdw, задевшая два шарда, оставила строку после ошибки у клиента в трёх случаях из четырёх.
Полное техническое изложение
TL;DR
- Шардирование — это когда одна большая база разложена по нескольким машинам, и каждая строка принадлежит ровно одной из них. Какая строка где лежит, решает одно правило и один столбец — ключ.
- Всё хорошее и всё плохое следует из ключа. Запрос, в котором ключ есть, идёт в одну машину. Запрос без ключа — во все, если нет отдельного справочника «где что лежит», и платит каждой.
- Добавить машину дорого по-разному. Простейшее правило — остаток от деления — при переходе с восьми машин на девять перевозит 88,9 % данных, хотя хватило бы 11,1 %. Правила поумнее укладываются около минимума.
- Операция, задевшая две машины, теряет гарантии обычной базы, если не добавить для неё отдельный механизм. Уникальность по столбцу вне ключа сама не проверяется, а транзакция на две машины через postgres_fdw может зафиксироваться наполовину — и при этом сообщить об ошибке.
Зачем делить базу на части
Представьте картотеку, которая перестала помещаться в один шкаф. Можно поставить рядом копии шкафа — это реплики: удобно, когда карточки часто читают, и спасает, если один шкаф сгорит. Но каждая копия всё равно хранит все карточки и получает все новые, поэтому места не прибавляется. А можно поставить несколько разных шкафов и разложить карточки по ним. Это и есть шардирование: каждая карточка принадлежит ровно одному шкафу, и шкаф хранит только свою часть. Одно другому не мешает: у каждого из разных шкафов могут быть свои копии.
Шкафы здесь — отдельные машины с базой, их называют шардами. Чтобы найти карточку, нужно знать, в каком она шкафу. Для этого у картотеки есть правило: например, «фамилии на А–Г — в первом шкафу, на Д–К — во втором». Столбец, по которому работает правило, называется ключом шардирования.
Ключ решает, в какой шкаф
Правило раскладки смотрит только на ключ. Если карточки разложены по фамилии, то вопрос «где карточка Иванова» решается мгновенно: открываем один шкаф. А вопрос «где карточка с таким-то телефоном» правилом не решить — оно про телефоны ничего не знает. Если нет отдельного справочника «телефон → шкаф», придётся открыть все шкафы подряд. Справочник можно завести, но его самого придётся держать в порядке.
Отсюда главное, что надо помнить про шардирование. Ключ выбирают не по тому, какой столбец «главный», а по тому, какие вопросы к картотеке задают чаще всего. Если почти все вопросы — про конкретного человека, ключ «человек» хорош. Если половина вопросов — поиск по телефону, этот ключ будет стоить дорого.
Когда шкафов становится больше
Картотека растёт, и восьми шкафов перестаёт хватать. Девятый шкаф должен получить свою девятую часть карточек — меньше перенести нельзя. Сколько переносится на самом деле, зависит от правила:
Самое простое правило — «номер шкафа = остаток от деления номера карточки на число шкафов». Оно раскладывает ровно, но при смене числа шкафов остаток меняется почти у всех карточек сразу. Это легко проверить самим:
import hashlib
def owner(key: str, shards: int) -> int:
h = int.from_bytes(hashlib.sha1(key.encode()).digest()[:8], "big")
return h % shards
keys = [f"user-{i}" for i in range(100_000)]
moved = sum(owner(k, 8) != owner(k, 9) for k in keys)
print(f"{moved / len(keys):.1%}")Программа напечатает, что сменили шкаф почти девять карточек из десяти. Правила поумнее — кольцо, заранее нарезанные «слоты» и ещё пара хитрых формул — переносят около минимума: новый шкаф забирает понемногу у каждого старого, а остальные карточки остаются на месте. Как устроено кольцо, подробно разобрано в уроке про консистентное хеширование.
Есть и беда, которую не лечит ни одно правило. Карточки можно разложить поровну, а обращаются к ним неравномерно: к одной знаменитой карточке могут приходить чаще, чем ко всем остальным в её шкафу вместе. Такая карточка лежит в одном шкафу, и весь поток к ней идёт туда. В модели, где на самую популярную карточку приходится 19,6 % всех обращений, её шкаф работает в 2,14 раза больше среднего. Одна эта карточка даёт своему шкафу в 1,57 раза больше средней нагрузки — и это не исправит никакое правило раскладки.
Запрос с ключом и без
Вернёмся к вопросу «где карточка с этим телефоном». На живом PostgreSQL это можно измерить. Главная база — координатор, в нашем стенде это PostgreSQL с расширением postgres_fdw, — на каждый шкаф, который ей пришлось открыть, тратит пять коротких разговоров по сети: начать, найти, забрать, закрыть, подтвердить. Если сеть отвечает за 2 мс, запрос с ключом занимает около 13 мс, сколько бы шкафов ни было. Запрос без ключа платит пять разговоров каждому шкафу:
У PostgreSQL есть две настройки, которые разрешают разговаривать со шкафами одновременно. Они помогают: на восьми шкафах запрос без ключа дешевеет с 104,0 до 67,9 мс. Но одновременно ведутся только два разговора из пяти — «забрать» и «подтвердить», — а три по-прежнему идут шкаф за шкафом. Поэтому чем больше шкафов, тем дороже запрос без ключа, как его ни настраивай.
Когда операция задевает два шкафа
Пока операция касается одного шкафа, всё работает как в обычной базе. Как только двух шкафов — две привычные гарантии сами больше не держатся. Вернуть их можно, но отдельным механизмом, и за него платят скоростью и сложностью.
Уникальность. Представьте правило «один телефон — одна карточка», а карточки разложены по фамилии. Каждый шкаф может проверить, что в нём нет двух карточек с одним телефоном, но не может заглянуть в соседний. PostgreSQL в такой ситуации честно отказывается создавать уникальный индекс вообще.
Всё или ничего. Операция кладёт по карточке в два шкафа, и один шкаф в последний момент отказывается. В обычной базе отказ означал бы, что не положено ничего. В нашем стенде шкафы подтверждают по очереди, и если первым подтвердил не тот, что отказался, его карточка останется:
Хуже всего то, что программа при этом получает ошибку. Обычная реакция на ошибку — повторить, и повтор кладёт в уцелевший шкаф вторую такую же карточку.
Что выбирать
- Сначала вопросы, потом ключ. Выберите ключ так, чтобы самые частые запросы и все важные операции касались одного шкафа.
- Шкафов — с запасом. Добавлять их по одному дорого, поэтому заранее делят данные на много мелких частей и при росте переносят части целиком.
- Переезд — с планом. Пока карточки переносят в новый шкаф и ещё не убрали из старого, они лежат в двух местах сразу: на стенде таблица насчитала 212 670 строк вместо 200 000. Заранее решите, кто отвечает за карточку в пути.
- Популярные записи — отдельной заботой. Их можно разбить на несколько копий, но за это платят чтением.
- Операции на два шкафа — без надежды на «всё или ничего». Уникальность вне ключа держат отдельной таблицей, а запись в два шкафа делают двумя шагами, так, чтобы повтор не создавал дублей.
Чем измерено
bench/sharding/placement.py— сколько данных переезжает при каждом правиле и что делает неравномерный доступ.bench/sharding/routing.py— запрос с ключом и без на живом PostgreSQL.bench/sharding/crossshard.py— уникальность и транзакция на два шарда.bench/sharding/resharding.py— добавление девятого шарда к восьми.
TL;DR
- Шардирование — это выбор ключа, а не нарезка таблицы. Ключ решает, в какой шард идёт строка, и тем самым делит все операции на два сорта: те, у которых ключ есть, идут в один шард, остальные — во все, если у маршрутизатора нет отдельных сведений о том, где искать. Всё остальное в статье — следствия этого деления.
- Добавить шард дорого по-разному. Восемь шардов становятся девятью: деление хеша по модулю перевозит 88,9 % ключей при минимуме 11,1 %, перерисованные диапазоны — 50,0 %. Кольцо, фиксированные слоты, rendezvous и jump hash укладываются около минимума, но кольцо оставляет перекос в 1,41 раза.
- Ровная раскладка ключей — не ровная нагрузка. При доступе по закону Ципфа с показателем 1,2 самый популярный ключ — 19,6 % всех запросов, и его шард несёт в 2,14 раза больше средней нагрузки. Один этот ключ — в 1,57 раза больше средней доли шарда, и ниже этого его шард не опустит никакая стратегия: ключ целиком живёт в одном шарде.
- Запрос без ключа платит каждому шарду. На стенде PostgreSQL с postgres_fdw координатор обменивается с каждым задетым шардом пять раз. При круге по сети 2 мс запрос с ключом стоит около 13 мс при любом числе шардов — 13,2 при четырёх и 13,4 при восьми, — без ключа 52,1 и 104,0.
- Настройки, которые обещают параллельность, совмещают не всё.
async_capableиparallel_commitвыключены по умолчанию; вместе они совмещают выборку и фиксацию, а три обмена из пяти по-прежнему идут шард за шардом. Запрос без ключа всё равно дорожает с каждым шардом: 67,9 мс на восьми. - Два шарда в одной транзакции через postgres_fdw — не атомарно. postgres_fdw
не использует двухфазную фиксацию. Когда один шард отказал на
COMMIT, клиент получил ошибку во всех четырёх опытах, а строка в другом шарде осталась в трёх из них. Сparallel_commit— при отказе любого из двух. Другая распределённая база может давать другой договор, но атомарность через шарды всегда требует отдельного механизма и оплачивается задержкой и сложностью. - Уникальность вне ключа шардирования сам PostgreSQL не обеспечивает. На таблице, чьи секции — шарды, PostgreSQL не создаёт уникального индекса вовсе.
- Перешардирование — это перенос данных на ходу. Пока строки копировались в
новый шард и ещё не были удалены из старого,
count(*)показал 212 670 строк вместо 200 000. Простое «скопировать, потом удалить» оставляет окно, в котором данные видны дважды.
База: шард, ключ и тот, кто знает, куда идти
Одна машина с базой упирается в потолок: в объём диска, в число записей в секунду, которое выдерживает журнал, в память под рабочий набор. Реплики этот потолок не поднимают: каждая держит все данные и принимает все записи. Они дают другое — масштабируют чтение, переживают отказ узла, сохраняют данные при потере машины, — но не объём и не поток записи. Шардирование делит сами данные: каждая строка принадлежит ровно одному шарду, и шард держит только свою долю.
Шардирование и репликация отвечают на разные вопросы и обычно используются вместе: у каждого шарда — свои реплики.
ключ → шард 3
├ основной узел: принимает записи
├ реплика
└ реплика
Четыре слова, без которых дальше не обойтись:
- Шард — часть строк на своём узле, со своими соединениями и своими транзакциями. Для приложения шард — отдельная база.
- Ключ шардирования — столбец (или несколько), по значению которого
решается, в каком шарде лежит строка. Здесь это
user_id. - Стратегия раскладки — правило «значение ключа → шард»: остаток от деления хеша, диапазоны, кольцо, таблица слотов.
- Маршрутизатор — тот, кто применяет правило к каждому запросу. Это может быть библиотека в приложении, отдельный прокси или, как на стенде этой статьи, координатор PostgreSQL с внешними таблицами.
Отсюда главное свойство, из которого следует всё остальное. Маршрутизатор
может отправить запрос в один шард, только если может вывести шард из
запроса. Проще всего — по ключу в условии: запрос по user_id идёт туда, где
лежат строки этого пользователя. Запрос по email правило раскладки адресовать
не умеет — оно про email не знает ничего. Если других сведений у маршрутизатора
нет, запрос идёт во все шарды. Такие сведения можно завести — каталог
email → шард, глобальный вторичный индекс, — но это отдельная структура со
своей ценой, и к ней статья ещё вернётся.
То же с операциями: уникальность, транзакция, соединение таблиц работают как обычно, пока задевают один шард. Задев два, они не исчезают, а перестают быть локальными: те же гарантии через несколько шардов требуют отдельного механизма — глобального индекса, распределённой фиксации, координатора соединений — и оплачиваются задержкой, доступностью и сложностью.
Поэтому ключ выбирают не по тому, какой столбец «естественный», а по тому,
какие запросы и какие транзакции должны остаться внутри одного шарда. Ключ
user_id хорош для системы, где почти всё происходит внутри одного
пользователя, и плох для системы, где половина запросов — поиск по email.
Как ключ становится адресом: пять стратегий и одна цена
Пока состав не меняется, любая разумная стратегия раскладывает ключи ровно. Различаются они тем, что происходит при смене состава. Самое частое изменение — добавить узел: восемь шардов становятся девятью. Чтобы новый узел получил свою долю — девятую часть, — должно переехать не меньше 11,1 % ключей. Это минимум при трёх условиях: узлы равноправны, у ключа один владелец, и раскладка до перехода была ровной. Всё сверх него перевозится зря, а меньше значит, что новый узел свою долю недополучил.
1. THE SAME CHANGE UNDER FIVE STRATEGIES: 8 NODES BECOME 9
----------------------------------------------------------
strategy keys moved min share max share max/min
modulo: hash mod N 88.9% 10.9% 11.3% 1.04x
ranges, all boundaries redrawn evenly 50.0% 10.9% 11.3% 1.03x
ranges, one range split in half 6.4% 6.3% 12.6% 1.99x
ring, 128 virtual nodes per node 10.7% 9.4% 13.3% 1.41x
slots: 1024, handed over whole 11.0% 11.0% 11.3% 1.03x
slots: 16384, handed over whole 11.2% 11.0% 11.2% 1.02x
Вывод скриптов здесь и дальше приводится дословно, по-английски: при сборке он сверяется с записью прогона символ в символ. В столбцах: keys moved — доля ключей, сменивших владельца; min share и max share — наименьшая и наибольшая доля ключей у одного узла после перехода; max/min — их отношение.
Остаток от деления — hash(key) mod N — раскладывает ровно и стоит одну
операцию. Его беда в том, что владелец ключа оказывается свойством не ключа, а
текущего числа узлов. Меняется делитель — меняется ответ почти для всех
ключей сразу: 88,9 % при минимуме 11,1 %, в восемь раз больше необходимого.
Диапазоны хранят ключи упорядоченными, и выборка «все заказы за неделю» или «пользователи с 1000 по 2000» идёт в один-два шарда, а не во все. При смене состава у них два пути, и оба неприятны. Перерисовать все границы ровно — и переедет половина ключей: сдвиг каждой границы задевает каждый отрезок. Разрезать один отрезок пополам — и переедет меньше минимума, 6,4 %, но новый узел получит не свою долю, а половину чужой: два узла держат вдвое меньше остальных семи.
Кольцо привязывает владельца к самому ключу: ключ идёт к ближайшей точке на окружности, и новый узел забирает только дуги перед своими точками. Переезжает 10,7 % — чуть меньше минимума, и это не нарушение: всё перевезённое досталось новому узлу, и он получил чуть меньше девятой части. Цена — случайность дуг: даже со 128 точками на узел самый загруженный держит в 1,41 раза больше самого свободного. Почему точек нужно много и чем они оплачиваются, разобрано в уроке про консистентное хеширование; числа кольца здесь и там — из одного и того же вычисления.
Фиксированные слоты разделяют два решения, которые остальные стратегии смешивают. Ключ отображается в один из заранее заданного числа слотов — это решение не меняется никогда. Слот принадлежит узлу — это решение меняется при смене состава, но слотами целиком: новый узел получает по нескольку слотов от каждого старого, и больше не переезжает ничего. Отсюда 11,0 % при 1024 слотах и ровная раскладка. Так устроен кластер Valkey:
HASH_SLOT = CRC16(key) mod 16384
ХЕШ_СЛОТ = CRC16(ключ) по модулю 16384
Слотов берут с запасом, во много раз больше узлов: слот — единица переноса, и раскладка не может быть ровнее, чем позволяет его размер. Это и есть схема «много логических шардов на немного физических узлов»: число слотов выбирают один раз и надолго, а узлы передают их друг другу целиком.
Из пяти стратегий таблицы мало перевозить и оставаться ровными одновременно умеют только слоты. Платят они таблицей «слот → узел», которую каждый маршрутизатор должен держать и вовремя обновлять. Но слоты — не единственный способ: ещё два правила обходятся без таблицы. Они посчитаны на том же изменении состава и тех же ключах:
5. TWO MORE RULES WITHOUT A SLOT TABLE: RENDEZVOUS AND JUMP HASH
----------------------------------------------------------------
the same change as in block 1: 8 nodes become 9, the same keys
strategy keys moved min share max share max/min
rendezvous (HRW), score per node 11.3% 10.9% 11.3% 1.03x
jump consistent hash 10.9% 10.9% 11.3% 1.04x
Столбцы — как в первом блоке.
Rendezvous, или HRW, даёт каждому узлу свою оценку ключа — хеш пары «ключ, узел» — и отдаёт ключ узлу с наибольшей. Новый узел забирает ровно те ключи, где его оценка выше всех прежних. Jump consistent hash (Lamping, Veach, 2014) вычисляет номер узла короткой арифметикой от хеша ключа и не хранит ничего. Оба уложились около минимума и ровны, как слоты. Платят они другим: rendezvous перебирает все узлы на каждый поиск, а jump hash требует, чтобы узлы были пронумерованы подряд, — убрать можно только последний. Есть и третий вопрос, помимо перевозки и ровности: какие запросы остаются внутри одного шарда. По нему выигрывают только диапазоны — выборка по отрезку ключей идёт в один-два шарда, а при хеше во все.
Ровные ключи — не ровная нагрузка
Таблица выше считает ключи. Нагрузку создают запросы, а к ключам они
обращаются неравномерно: популярный товар, знаменитость в социальной сети,
общий счётчик. Классическая модель такой неравномерности — закон Ципфа: доля
запросов к ключу с номером популярности r пропорциональна 1 / r^s.
Ключи ниже разложены делением хеша по модулю на восемь шардов, и разложены ровно: доли шардов — от 12,3 до 12,7 % ключей.
2. EVEN KEYS ARE NOT EVEN LOAD: ACCESS FOLLOWS A ZIPF LAW
---------------------------------------------------------
keys placed by hash mod 8: key shares per shard 12.3% to 12.7%,
max/min 1.03x - an even placement of keys
key of rank r gets a share of requests proportional to 1 / r^s;
'top key alone' is the top key's share over the average shard's
s top key, share of all requests hottest shard vs average top key alone
0.8 2.2% 13.9% 1.11x 0.18x
1.0 8.3% 18.3% 1.46x 0.66x
1.2 19.6% 26.8% 2.14x 1.57x
s — показатель закона Ципфа. top key, share of all requests — доля запросов к самому популярному ключу; hottest shard — доля нагрузки на самом загруженном шарде; vs average — во сколько раз она больше средней; top key alone — то же для одного популярного ключа, без остальных.
При s = 1,2 на самый популярный ключ приходится 19,6 % всех запросов, и шард,
в котором он лежит, несёт в 2,14 раза больше средней нагрузки. Часть перекоса
зависит от того, куда легли остальные ключи, но не весь: один этот ключ — в 1,57
раза больше средней доли шарда, и ни одна стратегия раскладки этого не
исправит. Правило отображает ключ ровно в одного владельца, и весь поток к
ключу идёт туда.
Единственный способ разнести горячий ключ — перестать хранить его одним
ключом: разбить на K копий key#0 .. key#K-1 и раскладывать копии по их
собственному хешу. Это помогает ровно настолько, насколько копии попадают в
разные шарды:
3. SPLITTING ONE HOT KEY INTO SUB-KEYS
--------------------------------------
s = 1.2; the top key is stored as K copies key#0 .. key#K-1,
each copy placed by its own hash, requests spread evenly over copies
K shards the copies hit hottest shard vs average
1 1 26.8% 2.14x
2 2 20.9% 1.67x
4 3 22.5% 1.80x
8 5 20.0% 1.60x
16 7 18.8% 1.50x
K — число копий ключа; shards the copies hit — в сколько разных шардов они легли; hottest shard и vs average — как в предыдущем блоке.
Четыре копии оказались хуже двух: они легли в три шарда, и одна из них — рядом с другими тяжёлыми ключами. Копии раскладываются тем же хешем, что и всё остальное, и выбрать им шарды нельзя. А платит за разбиение чтение: каждое обращение к ключу теперь выбирает копию — или читает все, если значение счётчик, который надо сложить.
Отдельный случай неравномерности — возрастающий ключ при диапазонах. Автоинкремент, время, упорядоченный по времени идентификатор: каждое новое значение больше всех прежних и попадает в последний диапазон.
4. AN INCREASING KEY AND RANGES: WHERE NEW WRITES GO
----------------------------------------------------
ids 1..100000 already stored, ranges of 12500 ids each, 8 shards;
the next 10000 inserts get ids 100001..110000
placement share of new writes on the busiest shard
ranges of the id 100.0%
hash of the id, mod 8 13.1%
Столбец справа — какая доля из десяти тысяч новых вставок пришлась на самый загруженный шард: при диапазонах по идентификатору и при хеше от него.
Все новые записи идут в один шард из восьми, остальные семь не получают ни одной. Хеш от идентификатора разносит вставки ровно — и отдаёт упорядоченную выборку по диапазону, ради которой диапазоны и брали.
Куда идёт запрос
Стенд собран штатными средствами PostgreSQL. Таблица orders на координаторе
секционирована по хешу user_id, и каждая её секция — внешняя таблица
postgres_fdw, указывающая в свой шард. Секционирование решает, куда идёт
строка, postgres_fdw — как туда дойти. Специализированные расширения делают ту
же работу своим кодом; здесь нужен встроенный механизм, потому что его
поведение описано в документации и его можно сверить.
Что делает маршрутизатор, видно прямо в плане запроса:
1. WHERE THE QUERY GOES: THE PLAN WITH AND WITHOUT THE SHARD KEY
----------------------------------------------------------------
where user_id = 42 (the table is partitioned by hash of user_id)
Foreign Scan on orders_2 orders
where email = 'user42@example.com'
Append
-> Foreign Scan on orders_0 orders_1
-> Foreign Scan on orders_1 orders_2
-> Foreign Scan on orders_2 orders_3
-> Foreign Scan on orders_3 orders_4
С ключом в условии — один шард: планировщик отсёк остальные секции ещё до выполнения. Без ключа — все четыре. Сколько стоит каждый задетый шард, видно, если разобрать, что координатор отправляет по сети. Ретранслятор стенда понимает протокол PostgreSQL и печатает сообщения:
2. WHAT THE COORDINATOR SENDS TO A SHARD FOR ONE QUERY
------------------------------------------------------
decoded from the wire by the relay; connections already open
query without the shard key, shard 0:
1. Q: START TRANSACTION ISOLATION LEVEL REPEATABLE READ
2. P+B+D+E+S: DECLARE c1 CURSOR FOR
3. Q: FETCH 100 FROM c1
4. Q: CLOSE c1
5. Q: COMMIT TRANSACTION
messages per shard, all shards: [5, 5, 5, 5]
query with the shard key, messages per shard: [0, 0, 5, 0]
Пять обменов на каждый задетый шард: открыть удалённую транзакцию, объявить
курсор, выбрать строки, закрыть курсор, зафиксировать. Каждый обмен — круг по
сети. Запрос с ключом платит пять кругов, запрос без ключа — пять на каждый
шард. Пять — пока шард возвращает не больше ста строк: FETCH 100 — размер
пачки postgres_fdw по умолчанию, и строки сверх него стоят новых обменов.
По умолчанию postgres_fdw ходит в шарды по очереди. Две настройки внешнего
сервера обещают это изменить, и обе выключены: allows foreign tables to be scanned concurrently for asynchronous execution
(разрешает сканировать внешние таблицы одновременно, в асинхронном исполнении)
у async_capable и commits, in parallel, remote transactions opened on a foreign server
(фиксирует параллельно удалённые транзакции, открытые на внешнем сервере)
у parallel_commit. Время одного запроса при задержке 1 мс в каждую сторону:
3. TIME PER QUERY, 4 SHARDS, 1 MS EACH WAY
------------------------------------------
configuration median best round worst round spread
with the shard key 13.2 ms 13.0 ms 13.6 ms 4%
no key, shards in turn 52.1 ms 51.3 ms 53.3 ms 4%
no key, async_capable 44.4 ms 43.9 ms 47.3 ms 8%
no key, async + parallel_commit 36.6 ms 36.4 ms 37.5 ms 3%
Каждая конфигурация замерена девятью сериями по тридцать запросов. median — медиана по сериям, best round и worst round — лучшая и худшая серия, spread — разница между ними в долях медианы.
Обе настройки помогают, и помогают устойчиво: вместе быстрее похода по очереди во всех девяти сериях замеров. Но насколько — видно только по росту с числом шардов, и здесь замер разошёлся с ожиданием. Одновременный поход, казалось бы, должен сделать цену почти плоской: все шарды опрашиваются разом, значит, и платить надо как за один. Этого не произошло:
6. WHAT GROWS WITH THE NUMBER OF SHARDS, 1 MS EACH WAY
------------------------------------------------------
configuration 4 shards 8 shards 8 / 4
with the shard key 13.2 ms 13.4 ms 1.01x
no key, shards in turn 52.1 ms 104.0 ms 2.00x
no key, async_capable 44.4 ms 85.6 ms 1.93x
no key, async + parallel_commit 36.6 ms 67.9 ms 1.86x
8 / 4 — во сколько раз запрос на восьми шардах дороже, чем на четырёх.
Вдвое больше шардов — и даже с обеими настройками почти вдвое дороже. Объяснение даёт счёт кругов, оплачиваемых строго друг за другом. Каждая конфигурация измерена при двух задержках, 1 и 3 мс; прирост времени, делённый на прирост круга, — это и есть число последовательных кругов. Постоянные расходы при вычитании сокращаются.
5. HOW MANY ROUND TRIPS ARE PAID ONE AFTER ANOTHER
--------------------------------------------------
each configuration timed at 1 ms and at 3 ms each way;
(time at 3 - time at 1) / (6 - 2 ms) = round trips in sequence
configuration 4 shards expected 8 shards expected
with the shard key 5.3 5 5.2 5
no key, shards in turn 20.8 20 41.1 40
no key, async_capable 17.6 17 34.2 33
no key, async + parallel_commit 14.6 14 27.0 26
Для каждого числа шардов два столбца: слева измеренное число последовательных кругов, справа expected — то, что следует из разбора протокола.
Колонка expected — не подгонка, а следствие разбора протокола: из пяти
обменов async_capable совмещает только выборку, parallel_commit — только
фиксацию, и совмещённый обмен стоит один круг на все шарды вместо одного на
каждый. Разности между строками это подтверждают независимо от остатка:
«по очереди» минус async_capable — 3,2 круга при четырёх шардах и 6,9 при
восьми, async_capable минус async + parallel_commit — 3,0 и 7,2. То есть
N − 1, три и семь, с точностью до 0,2 круга. Измеренное число везде выше разобранного на
0,2–1,2 круга; откуда этот остаток, разбор протокола не объясняет.
Открыть удалённую транзакцию, объявить курсор и закрыть его postgres_fdw по-прежнему делает шард за шардом. Три круга на каждый шард остаются при любых настройках, и в postgres_fdw запрос без ключа дорожает линейно с числом шардов, как его ни настраивай. Настройки стоит включать — на восьми шардах они снимают треть цены, 104,0 мс против 67,9, — но исправляет положение только ключ в условии.
Пять обменов на шард — свойство этого маршрутизатора, а не шардирования вообще: другой координатор может обмениваться с шардом иначе. Общее правило шире стенда: цена запроса без ключа растёт с числом задетых шардов настолько, насколько работу с ними не удаётся вести одновременно.
Что перестаёт работать само, когда операция задевает два шарда
Уникальность
Уникальный email на таблице, шардированной по user_id, — обычное требование:
двух аккаунтов с одним адресом быть не должно. Сервер отвечает на него так:
1. A UNIQUE EMAIL ON A TABLE SHARDED BY THE USER ID
---------------------------------------------------
the sharded table: partitions are foreign tables on the shards
alter table orders add unique (email)
ERROR: unique constraint on partitioned table must include all partitioning columns
DETAIL: UNIQUE constraint on table "orders" lacks column "user_id" which is part of the partition key.
alter table orders add unique (user_id, email)
ERROR: cannot create unique index on partitioned table "orders"
DETAIL: Table "orders" contains partitions that are foreign tables.
the same table with ordinary local partitions
alter table accounts add unique (email)
ERROR: unique constraint on partitioned table must include all partitioning columns
DETAIL: UNIQUE constraint on table "accounts" lacks column "user_id" which is part of the partition key.
alter table accounts add unique (user_id, email)
OK
На шардированной таблице уникального индекса нет никакого — даже с ключом шардирования в составе. На обычной секционированной таблице уникальность возможна, но только включающая ключ, и документация объясняет почему:
the constraint's columns must include all of the partition key columns. This limitation exists because the individual indexes making up the constraint can only directly enforce uniqueness within their own partitions
столбцы ограничения должны включать все столбцы ключа секционирования. Это ограничение существует потому, что отдельные индексы, из которых состоит ограничение, могут непосредственно обеспечивать уникальность только внутри своих секций
Каждый шард проверяет только свои строки, а дубль email под другим user_id
живёт в другом шарде. Обычный выход — вторая таблица, шардированная по
email: email → user_id. Уникальность тогда держит первичный ключ на каждом
шарде этой таблицы — на координаторе его не создать, как видно выше. Держит
потому, что маршрутизатор отправляет все строки одного адреса в один шард, и
проверяет это правило только сам маршрутизатор. А регистрация превращается в
две записи в два разных шарда — то есть в следующую проблему.
Такая таблица — уже не вспомогательная, а вторичный индекс, от которого зависит правильность. Держать его согласованным приходится на каждой операции: регистрация — запись в обе таблицы; удаление аккаунта — удаление из обеих; смена email — удаление старой строки и вставка новой, то есть снова два шарда; повтор после ошибки не должен создавать второй записи; а после частичной записи — строка в одной таблице есть, в другой нет — нужна процедура, которая найдёт и доведёт или откатит такие пары. Вторая таблица решает уникальность, но платят за неё согласованностью.
Атомарность
Транзакция вставляет по строке в шард 0 и в шард 1 и фиксируется. Один из
шардов отказывает ровно на COMMIT — после того как обе вставки прошли. Это
не экзотика: так срабатывают отложенные ограничения, и так же может выглядеть
обрыв соединения с шардом в неудачный момент.
2. ONE TRANSACTION, TWO SHARDS, ONE SHARD REFUSES AT COMMIT
-----------------------------------------------------------
insert one row for user 1 (shard 0) and one for user 3 (shard 1),
then COMMIT; the refusing shard has a deferred trigger that raises at commit
parallel_commit refusing shard client sees row on shard 0 row on shard 1
false 0 an error none kept
false 1 an error none none
true 0 an error none kept
true 1 an error kept none
the client got an error every time; a row survived anyway in 3 of 4 cases
refusing shard — какой шард отказал; client sees — что получил клиент, an error — ошибку; row on shard 0 и row on shard 1 — осталась ли в шарде вставленная строка: kept — да, none — нет.
Клиент получил ошибку все четыре раза, а строка пережила её в трёх. Механизм
описан документацией одной фразой: The remote transaction is committed or aborted when the local transaction commits or aborts
(Удалённая транзакция фиксируется или отменяется, когда фиксируется или отменяется локальная).
При фиксации локальной транзакции postgres_fdw фиксирует удалённые по очереди.
Если первым идёт отказавший шард, остальные откатываются, и ничего не
остаётся. Если первым идёт другой — он уже зафиксирован, и отменить
зафиксированное нечем. Какой шард идёт первым, документация не задаёт; в этом
прогоне первым шёл шард 1.
parallel_commit, который ускоряет запросы из предыдущего раздела, в этом
опыте с отказом делает хуже: все COMMIT уходят разом, и шард, который не
отказывал, зафиксировался в обоих опытах — какой бы из двух ни отказал.
Везение порядка пропадает вместе с порядком. Это не довод против настройки
вообще: обычную фиксацию она ускоряет, а меняет исход только при отказе.
Самое опасное в этом — не сама частичная фиксация, а то, что клиент видит ошибку. Обычная реакция на ошибку — повторить, и повтор создаёт в уцелевшем шарде дубль.
Двухфазная фиксация
Атомарность через несколько участников даёт двухфазная фиксация: сначала
каждый шард обещает, что зафиксирует, и только когда пообещали все, им велят
фиксировать. PostgreSQL её умеет, но postgres_fdw ею не пользуется:
Note that it is currently not supported by postgres_fdw to prepare the remote transaction for two-phase commit
(Заметим, что подготовка удалённой транзакции к двухфазной фиксации в postgres_fdw сейчас не поддерживается).
А на самих шардах она выключена по умолчанию:
3. TWO-PHASE COMMIT ON A SHARD, DEFAULT SETTINGS
------------------------------------------------
max_prepared_transactions = 0
prepare transaction 'order-42'
ERROR: prepared transactions are disabled
HINT: Set max_prepared_transactions to a nonzero value.
И даже включённая, она не решает задачу сама:
Two-phase transactions are intended for use by external transaction management systems
(Двухфазные транзакции предназначены для использования внешними системами управления транзакциями).
Кто-то снаружи должен помнить, какие шарды пообещали, и довести дело до конца
после своего же падения. На практике поэтому операции через шарды проектируют
так, чтобы атомарность не требовалась: ключ выбирают так, чтобы транзакция
оставалась в одном шарде, а там, где это невозможно, запись во второй шард
делают отдельным шагом — через
исходящий журнал в той же транзакции
и повтор, безопасный для дублей.
Смена состава на живом PostgreSQL
Модель в начале статьи говорит, сколько ключей переезжает. Живой сервер добавляет вопрос, которого у модели нет: разрешает ли он нужное изменение вообще. Хеш-секционирование PostgreSQL — это деление по модулю, и первая попытка добавить девятый шард к восьми выглядит так:
1. A NINTH PARTITION WITH MODULUS 9 NEXT TO EIGHT WITH MODULUS 8
----------------------------------------------------------------
create foreign table orders_8 partition of orders
for values with (modulus 9, remainder 8) server s8 ...
ERROR: every hash partition modulus must be a factor of the next larger modulus
DETAIL: The new modulus 9 is not divisible by 8, the modulus of existing partition "orders_7".
Правило из документации: every modulus which occurs among the partitions of a hash-partitioned table is a factor of the next larger modulus
(каждый модуль, встречающийся среди секций таблицы с хеш-секционированием, должен быть делителем следующего по величине модуля).
Восемь и шестнадцать уживаются, восемь и девять — нет. Пока на шард приходится
одна секция, добавить равноправный девятый шард нельзя вовсе. Можно одно из двух: разрезать один шард
надвое или разложить все строки заново.
Первый путь документация описывает сама — ровно эту процедуру стенд и выполнил:
You can detach one of the modulus-8 partitions, create two new modulus-16 partitions covering the same portion of the key space... and repopulate them with data.
Можно отсоединить одну из секций с модулем 8, создать две новые секции с модулем 16, покрывающие ту же часть пространства ключей... и заново наполнить их данными.
2. SPLITTING SHARD 0 IN TWO: MODULUS 8 BECOMES 16 FOR THAT SHARD ONLY
---------------------------------------------------------------------
rows in shard 0 before the split 24770
rows copied to the new shard 8 12670
share of all rows that travelled 6.3%
count(*) over the table before 200000
count(*) after copying, before deleting 212670
rows deleted from shard 0 12670
count(*) after deleting 200000
rows per shard after, shards 0..8: [12100, 25430, 25700, 24150, 24950, 25050, 24810, 25140, 12670]
largest / smallest shard 2.12x
Три наблюдения: два совпали с моделью, третьего в модели нет вовсе.
Перевезено 6,3 % строк — столько же, сколько модель дала для разрезанного диапазона, и с тем же итогом: два шарда держат вдвое меньше остальных, перекос 2,12. Хеш-секционирование при добавлении по одному шарду ведёт себя как диапазоны, а не как слоты.
Пока строки едут, таблица считает их дважды. Копирование и удаление — два
шага, и между ними count(*) показывает 212 670 строк вместо 200 000: 12 670
перевезённых видны и в старом шарде, и в новом. Новая секция с модулем 16 и
остатком 0 указывает на всю таблицу шарда 0, а что лежит во внешней секции,
PostgreSQL не проверяет: пока перевезённые строки не удалены из шарда 0,
координатор видит их и там, и в новом шарде. Документация предупреждает об
этом заранее:
it is then the user's responsibility that the contents of the foreign table satisfy the partitioning rule
(тогда пользователь сам отвечает за то, чтобы содержимое внешней таблицы соответствовало правилу секционирования).
Это наблюдение шире PostgreSQL. Перешардирование — это перенос данных на ходу, и у переходного состояния должен быть свой протокол: кто владеет строкой, пока она в пути, и что видит читатель. Обычные ответы — номер версии раскладки, который знают и маршрутизатор, и шарды; чтение и запись в оба места на время переноса; переключение владельца одной операцией после того, как копия догнала оригинал. Простое «скопировать, потом удалить» без такого протокола оставляет окно, в котором данные видны дважды, — стенд показал его одним числом.
Второй путь — ровная раскладка на девять — стоит того, что предсказывает модель для деления по модулю. Посчитано собственной хеш-функцией PostgreSQL, строка за строкой:
3. WHAT AN EVEN NINE-SHARD LAYOUT WOULD MOVE, BY POSTGRESQL'S OWN HASH
----------------------------------------------------------------------
rows 200000
rows whose shard changes, modulus 8 -> 9 177100
share 88.5%
the new shard fair share, 1 / 9 11.1%
88,5 % строк против 88,9 % ключей в модели: хеш другой, ключи другие, а цена та же, потому что это цена самого деления по модулю.
Что из этого выбирать
Сначала запросы, потом ключ. Выпишите запросы и транзакции, которые должны оставаться быстрыми и атомарными, и выберите ключ так, чтобы они задевали один шард. Всё, что останется без ключа, будет платить каждому шарду — на стенде это пять кругов на шард, и настройками их не свести к одному.
Число шардов — один раз и с запасом. При одной секции на шард хеш-секционирование PostgreSQL не позволяет добавить шард ровно: либо перекос вдвое, либо перевозка почти всего. Схема «много логических шардов на немного узлов» — как слоты Valkey — отделяет решение «ключ → шард» от решения «шард → узел» и переносит при смене состава шарды целиком, около минимума. С хеш-секционированием это значит секций больше, чем шардов: 72 секции по девять на каждом из восьми шардов переходят на девять шардов по восемь, и переезжают 8 секций из 72 — те же 11,1 %. Такой переход стенд не запускал: это арифметика, а не замер.
Диапазоны — ради выборок по диапазону, и только если следить за возрастающими ключами. Они единственные держат упорядоченную выборку в одном-двух шардах, но возрастающий ключ складывает все новые записи в последний шард, а добавление узла даёт выбор между половиной данных в дороге и перекосом вдвое.
Перешардирование — с протоколом, а не скриптом. Переход между раскладками — это перенос данных на ходу. Заранее решите, кто владеет строкой в пути и что видит читатель: номер версии раскладки, запись в оба места на время переноса, переключение одной операцией. Без этого «скопировать, потом удалить» показывает данные дважды — 212 670 строк вместо 200 000 на стенде.
Горячий ключ — отдельная задача. Никакая раскладка его не разнесёт: у ключа один владелец. Разбиение на копии помогает, пока копии легли в разные шарды, и перекладывает цену на чтение.
Операции через шарды — без атомарности по умолчанию. Уникальность вне ключа — отдельной таблицей, шардированной по этому столбцу. Запись в два шарда — двумя шагами с журналом и повтором, безопасным для дублей, а не одной транзакцией: postgres_fdw может зафиксировать её частично и при этом сообщить об ошибке.
Чем измерено
Первый скрипт считает без сервера, три остальных работают против живого
PostgreSQL; адрес задаётся переменной DE_BENCH_DSN. Каждый открывается прямо
отсюда, вместе с записью прогона.
bench/sharding/placement.py— раскладка при переходе с восьми узлов на девять для пяти стратегий и ещё двух без таблицы слотов (rendezvous, jump hash), неравномерный доступ, горячий ключ, возрастающий ключ при диапазонах. Ключи и хеш — изbench/hashring/ring.py.bench/sharding/routing.py— план запроса, разбор протокола между координатором и шардом, время запроса и число последовательных кругов для 4 и 8 шардов.bench/sharding/crossshard.py— уникальность вне ключа шардирования, транзакция в два шарда с отказом на фиксации, двухфазная фиксация по умолчанию.bench/sharding/resharding.py— девятый шард к восьми: запрет модуля 9, разрезание одного шарда с подсчётом строк в пути, цена ровной раскладки по хешу самого PostgreSQL.
Общая часть трёх последних — bench/sharding/stand.py: координатор, шарды и
ретранслятор с задержкой между ними.
Версия — PostgreSQL 16.13. Шарды на стенде — базы одного сервера, поэтому абсолютные миллисекунды здесь говорят о задержке канала, а не о железе; выводы сформулированы в кругах по сети и переносятся на любую задержку умножением.
Это не пересказ и не отдельный текст: всё ниже взято из самой статьи — её выжимка, заголовки разборов, колонка «на самом деле» и таблица версий. Поэтому разойтись со статьёй эти тезисы не могут.
Суть
- Шардирование — это выбор ключа, а не нарезка таблицы. Ключ решает, в какой шард идёт строка, и тем самым делит все операции на два сорта: те, у которых ключ есть, идут в один шард, остальные — во все, если у маршрутизатора нет отдельных сведений о том, где искать. Всё остальное в статье — следствия этого деления.
- Добавить шард дорого по-разному. Восемь шардов становятся девятью: деление хеша по модулю перевозит 88,9 % ключей при минимуме 11,1 %, перерисованные диапазоны — 50,0 %. Кольцо, фиксированные слоты, rendezvous и jump hash укладываются около минимума, но кольцо оставляет перекос в 1,41 раза.
- Ровная раскладка ключей — не ровная нагрузка. При доступе по закону Ципфа с показателем 1,2 самый популярный ключ — 19,6 % всех запросов, и его шард несёт в 2,14 раза больше средней нагрузки. Один этот ключ — в 1,57 раза больше средней доли шарда, и ниже этого его шард не опустит никакая стратегия: ключ целиком живёт в одном шарде.
- Запрос без ключа платит каждому шарду. На стенде PostgreSQL с postgres_fdw координатор обменивается с каждым задетым шардом пять раз. При круге по сети 2 мс запрос с ключом стоит около 13 мс при любом числе шардов — 13,2 при четырёх и 13,4 при восьми, — без ключа 52,1 и 104,0.
- Настройки, которые обещают параллельность, совмещают не всё.
async_capableиparallel_commitвыключены по умолчанию; вместе они совмещают выборку и фиксацию, а три обмена из пяти по-прежнему идут шард за шардом. Запрос без ключа всё равно дорожает с каждым шардом: 67,9 мс на восьми. - Два шарда в одной транзакции через postgres_fdw — не атомарно. postgres_fdw не использует двухфазную фиксацию. Когда один шард отказал на
COMMIT, клиент получил ошибку во всех четырёх опытах, а строка в другом шарде осталась в трёх из них. Сparallel_commit— при отказе любого из двух. Другая распределённая база может давать другой договор, но атомарность через шарды всегда требует отдельного механизма и оплачивается задержкой и сложностью. - Уникальность вне ключа шардирования сам PostgreSQL не обеспечивает. На таблице, чьи секции — шарды, PostgreSQL не создаёт уникального индекса вовсе.
- Перешардирование — это перенос данных на ходу. Пока строки копировались в новый шард и ещё не были удалены из старого,
count(*)показал 212 670 строк вместо 200 000. Простое «скопировать, потом удалить» оставляет окно, в котором данные видны дважды.
На самом деле
- Он раскладывает ровно ключи, а нагрузку создают запросы, и к ключам они обращаются неравномерно. В модели с законом Ципфа и показателем 1,2 на самый популярный ключ приходится 19,6 % всех запросов, и его шард несёт в 2,14 раза больше средней нагрузки — при ровной раскладке ключей. Один этот ключ — в 1,57 раза больше средней доли шарда, и этого не снимет никакая стратегия: правило раскладки отображает ключ в одного владельца, и поток к ключу целиком идёт туда (
bench/sharding/placement.py, блок 2). - Перевозит мало: при переходе с восьми узлов на девять — 10,7 % ключей при минимуме 11,1 %. Но дуги кольца случайны, и даже со 128 точками на узел самый загруженный держит в 1,41 раза больше самого свободного. Ровно и около минимума одновременно раскладывают фиксированные слоты — 11,0 % перевезено, перекос 1,03 — и два правила без таблицы: rendezvous (11,3 % и 1,03) и jump consistent hash (10,9 % и 1,04) (
bench/sharding/placement.py, блоки 1 и 5). - Из пяти обменов с шардом эта настройка совмещает только выборку. Открытие удалённой транзакции, объявление курсора, его закрытие и — без
parallel_commit— фиксация по-прежнему идут шард за шардом. На восьми шардах с обеими настройками запрос без ключа стоит 67,9 мс против 13,4 мс с ключом и в 1,86 раза дороже, чем на четырёх (bench/sharding/routing.py, блоки 5 и 6). - Для транзакции через два шарда это неверно. Когда один шард отказал на
COMMIT, клиент получил ошибку во всех четырёх опытах, а строка в другом шарде осталась в трёх. postgres_fdw фиксирует шарды без двухфазной фиксации, и отменить уже зафиксированное нечем (bench/sharding/crossshard.py, блок 2). - Он ещё и меняет исход отказа. Без него частичная фиксация зависит от того, какой шард фиксируется первым; с ним все
COMMITуходят разом, и шард, который не отказывал, фиксируется всегда — какой бы из двух ни отказал (bench/sharding/crossshard.py, блок 2). - Если на шард приходится одна секция, к восьми хеш-секциям с модулем 8 сервер не даст добавить девятую с модулем 9: каждый модуль обязан делить следующий по величине. Остаётся разрезать один шард надвое — 6,3 % строк в пути и перекос 2,12 — или разложить всё заново: по хешу самого PostgreSQL сменили бы шард 88,5 % строк (
bench/sharding/resharding.py). - Они решают разные задачи. Реплика хранит все данные и принимает все записи: она масштабирует чтение и переживает отказ узла, но не прибавляет ни места, ни пропускной способности записи. Шард хранит свою часть строк — и именно поэтому операции, задевшие два шарда, перестают получать уникальность и атомарность даром: их приходится возвращать отдельным механизмом. Обычно эти два средства сочетают: у каждого шарда свои реплики.
Что разобрано
- База: шард, ключ и тот, кто знает, куда идти
- Как ключ становится адресом: пять стратегий и одна цена
- Ровные ключи — не ровная нагрузка
- Куда идёт запрос
- Что перестаёт работать само, когда операция задевает два шарда
- Смена состава на живом PostgreSQL
- Что из этого выбирать
- Чем измерено
Унести в работу
- Выбирая ключ шардирования, сначала выпишите запросы и транзакции, которые обязаны оставаться быстрыми и атомарными, и берите ключ, при котором они задевают один шард. Всё, что останется без ключа, платит каждому шарду: на postgres_fdw это пять обменов с каждым.
- Запрос без ключа шардирования в горячем пути считайте дефектом схемы, а не местом для настройки.
async_capableиparallel_commitстоит включить — на восьми шардах они снимают треть цены, — но три обмена на шард остаются, и цена продолжает расти с числом шардов. - Не доверяйте ошибке транзакции, задевшей два шарда: postgres_fdw фиксирует шарды без двухфазной фиксации, и один из них может остаться зафиксированным. Повтор такой операции обязан быть безопасен для дублей — ключом идемпотентности или проверкой перед записью.
- Число шардов закладывайте один раз и с запасом — много логических шардов или слотов на немного узлов, — а не добавляйте по одному. При одной секции на шард хеш-секционирование PostgreSQL не даёт добавить шард ровно: либо перекос вдвое, либо перевозка почти всех строк. А сам переход проектируйте как перенос данных на ходу, с протоколом переходного состояния: без него «скопировать, потом удалить» показывает данные дважды — на стенде
count(*)увидел 212 670 строк вместо 200 000. - Уникальность по столбцу вне ключа шардирования держите отдельной таблицей, шардированной по этому столбцу: на таблице, чьи секции — шарды, PostgreSQL уникального индекса не создаст вовсе. Регистрация при этом станет записью в два шарда — проектируйте её как два шага, а смену email, удаление и повтор после ошибки — так, чтобы две таблицы не разошлись.
Частые заблуждения
Хеш от ключа разложит нагрузку по шардам ровно
Он раскладывает ровно ключи, а нагрузку создают запросы, и к ключам они обращаются неравномерно. В модели с законом Ципфа и показателем 1,2 на самый популярный ключ приходится 19,6 % всех запросов, и его шард несёт в 2,14 раза больше средней нагрузки — при ровной раскладке ключей. Один этот ключ — в 1,57 раза больше средней доли шарда, и этого не снимет никакая стратегия: правило раскладки отображает ключ в одного владельца, и поток к ключу целиком идёт туда (bench/sharding/placement.py, блок 2).
Консистентное хеширование и перевозит мало, и раскладывает ровно
Перевозит мало: при переходе с восьми узлов на девять — 10,7 % ключей при минимуме 11,1 %. Но дуги кольца случайны, и даже со 128 точками на узел самый загруженный держит в 1,41 раза больше самого свободного. Ровно и около минимума одновременно раскладывают фиксированные слоты — 11,0 % перевезено, перекос 1,03 — и два правила без таблицы: rendezvous (11,3 % и 1,03) и jump consistent hash (10,9 % и 1,04) (bench/sharding/placement.py, блоки 1 и 5).
С async_capable запрос во все шарды стоит как запрос в один
Из пяти обменов с шардом эта настройка совмещает только выборку. Открытие удалённой транзакции, объявление курсора, его закрытие и — без parallel_commit — фиксация по-прежнему идут шард за шардом. На восьми шардах с обеими настройками запрос без ключа стоит 67,9 мс против 13,4 мс с ключом и в 1,86 раза дороже, чем на четырёх (bench/sharding/routing.py, блоки 5 и 6).
Если транзакция вернула ошибку, ничего не записано
Для транзакции через два шарда это неверно. Когда один шард отказал на COMMIT, клиент получил ошибку во всех четырёх опытах, а строка в другом шарде осталась в трёх. postgres_fdw фиксирует шарды без двухфазной фиксации, и отменить уже зафиксированное нечем (bench/sharding/crossshard.py, блок 2).
parallel_commit только ускоряет фиксацию
Он ещё и меняет исход отказа. Без него частичная фиксация зависит от того, какой шард фиксируется первым; с ним все COMMIT уходят разом, и шард, который не отказывал, фиксируется всегда — какой бы из двух ни отказал (bench/sharding/crossshard.py, блок 2).
Чтобы добавить шард в PostgreSQL, достаточно добавить секцию
Если на шард приходится одна секция, к восьми хеш-секциям с модулем 8 сервер не даст добавить девятую с модулем 9: каждый модуль обязан делить следующий по величине. Остаётся разрезать один шард надвое — 6,3 % строк в пути и перекос 2,12 — или разложить всё заново: по хешу самого PostgreSQL сменили бы шард 88,5 % строк (bench/sharding/resharding.py).
Шардирование — то же, что реплики, только сложнее
Они решают разные задачи. Реплика хранит все данные и принимает все записи: она масштабирует чтение и переживает отказ узла, но не прибавляет ни места, ни пропускной способности записи. Шард хранит свою часть строк — и именно поэтому операции, задевшие два шарда, перестают получать уникальность и атомарность даром: их приходится возвращать отдельным механизмом. Обычно эти два средства сочетают: у каждого шарда свои реплики.
Проверка знаний
Таблица заказов шардирована по хешу user_id на восемь шардов. Какой запрос координатор отправит в один шард?
Источники и что читать дальше
7 ИСТОЧНИКОВ
- PostgreSQL 16 — F.38. postgres_fdwОфициальная документация. Как устроена транзакция, задевшая удалённый сервер: «The remote transaction is committed or aborted when the local transaction commits or aborts» (Удалённая транзакция фиксируется или отменяется, когда фиксируется или отменяется локальная). И граница, из которой следует блок про атомарность: «Note that it is currently not supported by postgres_fdw to prepare the remote transaction for two-phase commit» (Заметим, что подготовка удалённой транзакции к двухфазной фиксации в postgres_fdw сейчас не поддерживается). Оттуда же обе настройки, которые меряются в статье, и их значения по умолчанию: async_capable — «allows foreign tables to be scanned concurrently for asynchronous execution» (разрешает сканировать внешние таблицы одновременно, в асинхронном исполнении), parallel_commit — «commits, in parallel, remote transactions opened on a foreign server» (фиксирует параллельно удалённые транзакции, открытые на внешнем сервере); у обеих «The default is false» (По умолчанию — false).https://www.postgresql.org/docs/16/postgres-fdw.html
- PostgreSQL 16 — 5.11. Table PartitioningОфициальная документация. Правило, которое запрещает уникальный email на таблице, шардированной по user_id, и объяснение его причины: «the constraint's columns must include all of the partition key columns. This limitation exists because the individual indexes making up the constraint can only directly enforce uniqueness within their own partitions» (столбцы ограничения должны включать все столбцы ключа секционирования. Это ограничение существует потому, что отдельные индексы, из которых состоит ограничение, могут непосредственно обеспечивать уникальность только внутри своих секций). И оговорка про внешние секции, то есть про шарды: «it is then the user's responsibility that the contents of the foreign table satisfy the partitioning rule» (тогда пользователь сам отвечает за то, чтобы содержимое внешней таблицы соответствовало правилу секционирования).https://www.postgresql.org/docs/16/ddl-partitioning.html
- PostgreSQL 16 — CREATE TABLE, hash partitionsОфициальная документация. Почему девятый шард нельзя добавить к восьми с тем же хешем: «every modulus which occurs among the partitions of a hash-partitioned table is a factor of the next larger modulus» (каждый модуль, встречающийся среди секций таблицы с хеш-секционированием, должен быть делителем следующего по величине модуля). И ровно та процедура, которую статья выполняет на стенде: «You can detach one of the modulus-8 partitions, create two new modulus-16 partitions covering the same portion of the key space... and repopulate them with data» (Можно отсоединить одну из секций с модулем 8, создать две новые секции с модулем 16, покрывающие ту же часть пространства ключей... и заново наполнить их данными).https://www.postgresql.org/docs/16/sql-createtable.html
- PostgreSQL 16 — 74.4. Two-Phase TransactionsОфициальная документация. Для кого двухфазная фиксация в PostgreSQL: «Two-phase transactions are intended for use by external transaction management systems» (Двухфазные транзакции предназначены для использования внешними системами управления транзакциями). То есть сервер умеет обещать фиксацию, но решать, когда её выполнить, должен кто-то снаружи — сам postgres_fdw этого не делает.https://www.postgresql.org/docs/16/two-phase.html
- Valkey — Cluster specificationОфициальная документация. Фиксированные слоты в промышленной системе: «HASH_SLOT = CRC16(key) mod 16384» (ХЕШ_СЛОТ = CRC16(ключ) по модулю 16384). Ключ отображается в слот, а не в узел, и при смене состава узлы передают друг другу слоты целиком — поэтому модель статьи считает вариант с 16 384 слотами отдельной строкой.https://valkey.io/topics/cluster-spec/
- A Fast, Minimal Memory, Consistent Hash AlgorithmИсточник. Джон Лампинг, Эрик Вич (Google), 2014. Jump consistent hash: номер узла вычисляется короткой арифметикой от хеша ключа, без таблицы и без памяти, ценой того, что узлы пронумерованы подряд. Реализация в
bench/sharding/placement.py, блок 5, — по статье; там же rendezvous (HRW) на тех же ключах.https://arxiv.org/abs/1406.2294 - Замеры этой статьи: модель раскладки и стенд PostgreSQLИсточник. Четыре скрипта: точное вычисление раскладки при смене состава и неравномерном доступе, и три прогона против живого PostgreSQL 16.13 — маршрут запроса с разбором протокола, операции через два шарда и перешардирование. Шарды — базы на одном сервере за postgres_fdw, у каждого свой ретранслятор с задержкой; почему этого достаточно для выводов статьи, разобрано в описании замеров./ru/bench/sharding/routing.py