Про шардирование легко рассказать за пять минут: выбрали ключ, разложили данные по узлам, готово. Беда в том, что раскладывание — единственная простая часть. Всё интересное начинается после неё, когда выясняется, что запрос без ключа шардирования идёт во все шарды сразу, join двух таблиц с разными ключами база отказывается выполнять вовсе, а глубокая страница выдачи обходится в разы дороже первой.
Эта статья — про счёт, который приходит после раскладывания. Не про выбор ключа: он разобран в статье про шардирование MongoDB, где и выбор, и хотспоты, и ре-шардинг показаны на живом кластере. Здесь — про то, чем платят за уже случившуюся распределённость, и все числа сняты прогонами на кластере Citus, а не пересказаны по документации.
В статье
- Что здесь считается ценой
- Запрос без ключа идёт во все шарды
- Join: колокация или отказ
- Референсные таблицы — другой механизм локальности
- Пагинация: восемьсот тысяч строк ради двадцати
- Новый узел сам по себе не берёт нагрузку
- Что смотреть в проде
- Где проходит граница
- Демо и версии
- Документация
Что здесь считается ценой
Распределённость стоит трёх вещей, и все они измеримы.
Число задач по шардам, которые порождает один запрос. Если координатор может по значению фильтра понять, где лежат данные, он планирует одну задачу на один шард. Если не может — по задаче на каждый шард. Это не то же самое, что число узлов: задачи раскладываются по воркерам, и на стенде тридцать две задачи ложатся на двух.
Перемещение данных между узлами. Join, сортировка и агрегация иногда требуют, чтобы данные с одного узла оказались рядом с данными другого. Это одна из самых дорогих операций в распределённой схеме.
Объём, который координатор поднимает к себе. Он не всегда равен объёму ответа, и разница бывает разительной.
Дальше — каждая из трёх на замерах. Стенд: Citus 14.1 поверх PostgreSQL 17, координатор и два воркера, таблицы по 32 шарда.
Сразу оговорка, которая касается всех чисел ниже. Все узлы стенда живут на одном хосте, сетевой задержки между ними практически нет. Поэтому переносить отсюда можно характер зависимости и структурные величины — число задействованных задач, число шардов и размещений, объём поднятых строк, — но не абсолютные миллисекунды. На настоящем кластере, где узлы разнесены, добавятся сетевой объём и удалённые взаимодействия между узлами, и они не бесплатны. Насколько это сдвинет времена — здесь не измерено: строки идут между узлами потоком, а не отдельным round-trip на строку, и предсказывать по однохостовому прогону даже направление сдвига я не берусь.
И отдельно: сравнивать Citus с обычным PostgreSQL по скорости на таком стенде нельзя. На данных, влезающих в один узел, распределённая схема добавляет слой координации там, где обычная база просто читает таблицу, — и на таком стенде это чистые накладные расходы. Оговорюсь, что «распределённое всегда проигрывает на малых данных» из этого не следует: параллельное исполнение по шардам иногда выигрывает даже на одном хосте, если запрос упирается в процессор. Просто этот стенд такой вопрос не ставит и ответить на него не может. Ниже сравниваются распределённые схемы между собой.
Запрос без ключа идёт во все шарды
Два запроса к одной таблице, отличающиеся ровно одним: по какой колонке фильтр.
-- А: фильтр по ключу шардирования
SELECT count(*) FROM orders WHERE customer_id = 42;
-- Б: фильтр по обычной колонке
SELECT count(*) FROM orders WHERE total > 100;Главная величина здесь не время, а число задач в плане — сколько шардов координатор счёл нужным опросить:
| Ветка | Задач в плане | Время: min / медиана / p95 / max |
|---|---|---|
| По ключу шардирования | 1 | 11.1 / 11.7 / 13.1 / 13.8 мс |
| По обычной колонке | 32 | 54.4 / 56.0 / 63.9 / 72.7 мс |
Двадцать повторов на ветку, число задач совпало на всех двадцати. Это структурный признак: он не зависит от загрузки машины и воспроизводится точно, в отличие от времени.
Механика простая. Значение customer_id = 42 по хеш-функции однозначно указывает на один шард, и координатор исключает остальные 31, не обращаясь к ним. Фильтр по total ничего не говорит о размещении, поэтому исключить нельзя ни один шард — опрашиваются все.
Здесь стоит развести три величины, которые легко слипаются в одну. Task Count: 32 — это тридцать две задачи по шардам, а не тридцать два узла. Воркеров на стенде два, и тридцать две задачи ложатся на них, по шестнадцать на каждый. Цена складывается из трёх разных слагаемых: растёт число задач и суммарная работа; веер идёт по фактическим узлам, которых здесь всего два; и запрос завершается по самому медленному завершившемуся — задаче или узлу. Смешивать их нельзя: на стенде с двумя воркерами и на кластере из тридцати двух узлов при одинаковом Task Count: 32 вторая и третья составляющие будут совершенно разными.
Про форму хвоста я сначала написал вывод, а потом его снял. По первому прогону выходило, что у ветки без ключа рассеяние асимметричнее, и я объяснил это ожиданием самого медленного из веера. На повторном прогоне картина оказалась обратной — хвост тяжелее у ветки с одной задачей. На третьем стороны сравнялись. На четвёртом, который сейчас записан как канонический, снова тяжелее у ветки без ключа.
Четыре прогона, три разные картины. Совпадение канонического прогона с исходной догадкой ничего не возвращает: если бы я оставил тот вывод, он сейчас выглядел бы подтверждённым — и был бы ровно так же не доказан. Вдобавок эксперимент для такого вывода не годится в принципе: ветки отличаются не только числом задач, но и предикатом, и объёмом сканирования, так что влияние веера в нём не изолировано. Проверять форму хвоста нужно отдельным опытом, где меняется только число задач; здесь его нет, поэтому выводов о хвосте не будет.
Что до самих времён: двадцати наблюдений хватает, чтобы увидеть наблюдаемый разброс, а не шум одного числа, но мало для оценки «настоящего» p95 — доверительный интервал у перцентиля по двадцати точкам широк. Это порядок величины, как и остальное время в статье. Воспроизводится точно только Task Count.
Join: колокация или отказ
Два join одинаковой формы, отличающиеся тем, разложены ли соединяемые таблицы по одному ключу.
orders и customers шардированы обе по customer_id и лежат в одной группе колокации — парные шарды физически на одном узле. shipments шардирована по order_id, в отдельной группе.
-- А: таблицы со-расположены (обе по customer_id)
SELECT c.city, count(*) FROM orders o
JOIN customers c USING (customer_id) GROUP BY c.city;
-- Б: ключи разные (orders по customer_id, shipments по order_id)
SELECT s.carrier, count(*) FROM orders o
JOIN shipments s USING (order_id) GROUP BY s.carrier;Первый результат оказался неожиданнее, чем я планировал. Второй запрос база не выполняет вовсе:
ERROR: the query contains a join that requires repartitioning
HINT: Set citus.enable_repartition_joins to on to enable repartitioningТо есть речь не о «медленно», а об отказе. Citus по умолчанию не строит план, требующий переброски данных между узлами, — и заставляет вас включить это осознанно.
Включаем и сравниваем планы:
| Ветка | Задач в плане | Стадии переброски | Время |
|---|---|---|---|
| Со-расположены | 32 | нет | ~58–70 мс |
| Ключи разные | 8 | две (Map 32 → Merge 8) | ~866–874 мс |
У первой ветки join выполняется целиком внутри каждого шарда — соединяемые строки уже рядом. У второй в плане появляются две стадии MapMergeJob, по одной на каждую перераспределяемую таблицу: данные читаются с 32 шардов, перебрасываются и сливаются в 8 задач.
Ловушка в этой таблице. Число задач у репартиционного join меньше — 8 против 32. Соблазн прочитать это как «эффективнее» велик, и он ошибочен: величины структурно про разное. У первой ветки это число шардов, у второй — число задач финальной стадии слияния; сами шарды читаются все 32, это видно как Map Task Count внутри стадии. Откуда берётся именно восемь, я не выяснял и придумывать не буду — записано как наблюдение.
Референсные таблицы — другой механизм локальности
Маленький справочник можно не распределять, а объявить референсным. Сравним два варианта одного и того же справочника перевозчиков: как референсная таблица и как обычная распределённая копия.
По планам запросов они различаются так же, как ветки предыдущего раздела: у референсной переброски нет, у распределённой копии — две стадии. Но интереснее не это, а почему локальность возникает:
| Шардов | Размещений | |
|---|---|---|
| Референсная таблица | 1 | 2 (по числу воркеров) |
| Распределённая копия | 32 | 32 |
Референсная таблица — это один шард, скопированный на каждый воркер. Именно на воркер, а не на каждый узел кластера: координатор — тоже узел, но размещений здесь два, по числу воркеров. Join с ней локален потому, что копия справочника уже физически рядом с любыми данными, независимо от их ключа.
Колокация из предыдущего раздела добивается локальности иначе: там оба участника разложены по 32 шарда с одним размещением каждый, и локальность возникает из согласованного разбиения по общему ключу — без единой лишней копии.
По тексту плана запроса эти два механизма неразличимы. В обоих случаях просто нет стадии переброски. Различает их только таблица размещений — и это ровно тот случай, когда план запроса не отвечает на вопрос «почему», а системный каталог отвечает.
Практический вывод следует из размещений, а не из замера: копия на каждом воркере означает, что запись в справочник идёт на все воркеры сразу. Референсные таблицы хороши ровно там, где они и задуманы — маленькие, читаемые часто, меняющиеся редко.
Пагинация: восемьсот тысяч строк ради двадцати
Здесь пришлось строить отдельную таблицу на миллион строк: на исходных четырёх тысячах эффект насыщается — при 125 строках на шард смещение в тысячу уже превышает всё, что шард может отдать, и замер показывал бы остаток эффекта вместо самого эффекта.
Запрос обычный: двадцать строк, отсортированных не по ключу шардирования, с растущим смещением.
SELECT order_id, total FROM orders_big
ORDER BY total DESC LIMIT 20 OFFSET :n;Структурная величина — сколько строк координатор поднял с шардов ради этих двадцати:
| Смещение | Строк поднято с шардов | Объём |
|---|---|---|
| 0 | 640 | 12 кБ |
| 1 000 | 32 640 | 637 кБ |
| 10 000 | 320 640 | 6 256 кБ |
| 25 000 | 800 640 | 15 МБ |
Закономерность точная: 32 × (смещение + 20). Координатор не знает, на каком шарде лежат нужные строки, поэтому каждый шард обязан отдать ему свои смещение + лимит кандидатов — и лишь потом всё это сортируется и отбрасывается.
Восемьсот тысяч строк ради двадцати. Само число строк детерминировано и совпадает с формулой до единицы на каждом прогоне — а вот время ведёт себя куда хуже. Медиана растёт, но кратность гуляет — от неполных четырёх раз до почти семнадцати при переходе от нулевого смещения к 25 000. Границы я тут называть не буду, и вот почему. Сначала я написал «в 6.3–8.0 раза, воспроизводимо»; это оказалось диапазоном пары удачных прогонов. Потом вписал в стенд наблюдавшуюся огибающую — и первый же прогон после этого вышел за её верхний край. Величина, для которой очередной запуск пробивает только что записанную границу, не имеет опубликованного значения: у неё есть только порядок и разброс. Во всех имеющихся прогонах медиана выросла — но это наблюдение, а не структурный инвариант, в отличие от числа строк. Во сколько именно раз — вопрос без ответа на таком стенде.
А если сортировать по ключу шардирования? Соблазнительная мысль, и ответ нетривиален. Строк с шардов поднимается ровно столько же — те же 800 640, до единицы. Но планы различаются:
| Сортировка не по ключу | Сортировка по ключу | |
|---|---|---|
| Чтение шарда | Seq Scan + сортировка |
Index Scan, сортировки нет |
| Объём с шардов | 15 МБ | 21 МБ |
| Сортировка на координаторе | top-N heapsort в памяти |
external merge на диске |
Сортировка по ключу убирает работу внутри шарда, но Citus добавляет служебную колонку сортировки — и объём растёт, а сортировка на координаторе перестаёт помещаться в память. Итог разнонаправленный: на первой странице такой запрос обычно быстрее, на глубокой — обычно медленнее. Именно «обычно»: направление тоже держится не всегда. В одном из прогонов оно менялось трижды по разным точкам смещения — то есть и это наблюдение, а не закон. Структурная часть от прогона к прогону при этом не дрогнула ни разу: Index Scan против Seq Scan, top-N heapsort против сортировки с уходом на диск, и точное число поднятых строк.
Две оговорки, без которых этот вывод обманет. Первая: Index Scan здесь возможен потому, что первичный ключ таблицы начинается с колонки шардирования — это свойство конкретной схемы, а не общее свойство сортировки по ключу. Вторая: уход сортировки на диск зависит от work_mem, и при большем значении проигрыш на глубокой странице может уменьшиться, исчезнуть или сменить знак. Запас тут тоньше, чем кажется: при work_mem в 4096 кБ сортировка серии А дошла до 3118 кБ — то есть чуть больше данных или чуть меньше памяти, и на диск уедет уже она, а весь контраст между сериями исчезнет. Это ограничение конфигурации, а не свойство Citus.
Настоящее лечение глубокой пагинации — не в выборе колонки сортировки, а в отказе от смещения в пользу курсора («покажи следующие двадцать после вот этой строки»). Это работает и на одном узле, а на шардах цена ошибки заметно выше — потому что лишние строки поднимаются с каждого шарда сразу. Курсор не универсален (он плохо дружит с произвольными переходами «на страницу N» и с сортировкой по неуникальному ключу), но там, где выбор есть, на шардах он предпочтителен.
Новый узел сам по себе не берёт нагрузку
Кластер из двух воркеров, 96 размещений распределённых таблиц (три таблицы по 32 шарда) плюс две копии референсного справочника — 98 размещений всего. Добавляем третий узел.
| Момент | w1 | w2 | w3 | Копий справочника |
|---|---|---|---|---|
| До добавления | 49 | 49 | — | 2 |
| После регистрации узла | 49 | 49 | 0 | 2 |
| После ребаланса | 34 | 34 | 31 | 3 |
Регистрация узла не двигает ни одного шарда. Это и есть главное наблюдение раздела: новый узел стоит пустым, пока ребаланс не запущен явно. Если бы данные переезжали сами, ребаланс был бы не нужен — он существует именно потому, что не переезжают.
Второе: копию справочника на новый узел кладёт ребаланс, а не регистрация. Причём делает это первой задачей, блокируя перенос обычных шардов до её завершения.
Здесь стоит быть точным с числами, потому что перепутать их легко. Всего размещений становится 99 вместо 98, но 96 размещений распределённых таблиц остаются 96 — они лишь перераскладываются по трём узлам. Прибавка в единицу — это третья копия справочника, а не потерянная или задвоенная строка данных. Общее число размещений и число распределённых — два разных инварианта, и путать их не стоит: при сборке стенда самопроверка на этом и споткнулась, объявив появление копии потерей данных.
Доступен ли переезжающий шард во время ребаланса. Раз в секунду скрипт бьёт в него двумя пробами: читающей и пишущей. Пишущая не просто вставляет строку, а проверяет, что вставленное сразу видно — иначе «INSERT не отказал» ничего не значило бы.
Восемь прогонов, 154 пробы каждого вида: отказов ноль ни у чтения, ни у записи, потерь видимости нет ни разу. Максимумы записи — 183, 153, 1301, 316, 959, 190, 538 и 2668 мс; чтения — 241, 146, 894, 314, 145, 163, 507 и 603 мс. Медиан и квантилей не привожу: скрипт печатает не каждую итерацию, так что распределения у меня попросту нет — только максимумы и факт отсутствия отказов.
И сразу главное ограничение, потому что без него весь абзац читается сильнее, чем есть. Пишущая проба не определяет наличие краткой блокировки. Заблокированный INSERT обычно ждёт и потом завершается успешно — проба его засчитает. statement_timeout в стенде не задан, поэтому долгая блокировка не дала бы ошибку, а подвесила бы вызов. И отказ сам по себе блокировку не доказывал бы: для этого нужен разбор текста ошибки или pg_locks, чего стенд не делает. Проверяется ровно две вещи: что INSERT завершился и что записанное видно.
Дальше — почему проб именно две и почему бьют они именно в этот шард. Обе детали появились не сразу, и каждая закрывала собственную дыру в доказательстве.
Первый вопрос: куда именно бил пробный запрос. Здесь у меня была подмена, которую я не заметил. В первой версии ключ пробы был константой, а его шард в план перемещений не входил. То есть успешные пробы доказывали доступность неподвижного шарда во время фоновой задачи — не чтение данных в процессе их переноса. Утверждение звучало сильно, а измеряло не то.
Сейчас ключ выбирается из фактического плана: скрипт читает get_rebalance_table_shards_plan(), находит customer_id, чей шард orders в этом плане есть, и проверяет это гейтом. В записанном прогоне гейт выбрал customer_id = 3, шард orders=102055, переезжающий с первого воркера на третий.
Одного этого гейта, впрочем, мало, и на это мне указали отдельно. Он сверяет выбранный шард с тем же предварительным планом, из которого шард и взят, — то есть проверяет источник на согласие с самим собой. План рассчитывается до старта задачи и теоретически может с ней разойтись. Поэтому после завершения ребаланса стоит второй гейт: активное размещение шарда пробы обязано оказаться именно на том узле, который план назвал целевым, иначе прогон объявляется провалившимся. План обещал citus-w3, факт по pg_dist_placement — citus-w3.
А пишущая проба появилась потому, что читающая доказывала не то. Сначала я писал, что логическая репликация даёт перенос «без блокировки чтения», и подпирал это успешными SELECT. Связка неверная: режимы Citus различаются по записям. block_writes копирует данные через COPY с блокировкой записи, а читать в нём можно точно так же. Значит читающая проба вернула бы успех в обоих режимах — она их не различает вообще и ничего про выбранный режим не доказывает. Наблюдение было верным, вывод к нему не относился.
Границы назову прямо, их четыре. Ребаланс здесь занимает десятки секунд на четырёх тысячах заказов и двадцати перемещениях: в двух прогонах, где время засечено часами, окно опроса составило 31.5 и 39.9 секунды. До этого я писал «21–26 секунд», выводя длительность из числа проб, — прямой замер показал, что оценка занижала её на треть и больше. Ещё один случай, когда правдоподобное число оказалось невычисленным — на реальных объёмах он идёт часами, и окно для неприятностей несоизмеримо шире. Дальше: пробы идут раз в секунду, а окно переключения могло оказаться короче этого интервала — тогда кратковременную блокировку они просто не застали бы. Длительность самого окна я не измерял, так что «заведомо короче» сказать нельзя; но отсутствие отказов и без того не доказывает, что записи не блокируются никогда. Третье: контрольного прогона в режиме block_writes я не делал, поэтому стенд показывает, что записи выживают при auto, но не показывает, что при block_writes они бы не выжили. И четвёртое, самое узкое: скрипт следит за состоянием задачи целиком, а не за фазой конкретного шарда, — то есть доказано, что пробы шли всё время работы задачи, но не что хотя бы одна пришлась ровно на момент переключения именно нашего шарда.
Разовые всплески длительности есть, и объяснить их я не могу. Сначала они попадались только на чтении, и я написал, что всплеска, характерного для ожидания блокировки, не видно. Потом прогон дал обратное — максимум записи выше максимума чтения. То есть соотношение максимумов не воспроизвелось и признаком служить не может, и объяснение пришлось снять вместе с эвристикой, которая по этому соотношению пыталась ставить диагноз прямо в скрипте. В нынешнем наборе четыре прогона держались в пределах 145–316 мс. Дальше четыре разных отклонения: в одном обе пробы поднялись умеренно и вместе (507 и 538 мс), в другом обе выросли крупно (894 и 1301), в третьем выросла только запись — 959 против 145, — а в четвёртом, снятом уже не мной, запись дошла до 2668 мс при чтении 603. Ни в одном из них ни одна проба не отказала и ни одна запись не потерялась.
Одна деталь прогона с умеренным ростом обеих проб (в записях стенда это job_id=30) стоит упоминания, хотя объяснения у меня и для неё нет: максимум там пришёлся на первую итерацию, когда задача ещё стояла в scheduled и ни один шард не двигался. Похожее, слабее выраженное, есть ещё в двух прогонах. Напрашивается прогрев — первое обращение в свежем соединении, — но замера холодного и тёплого вызова стенд не делает, так что это догадка, а не вывод.
Прогон, где выросла только запись (job_id=27, 959 мс против 145 у чтения), стоит прочесть аккуратно — напрашивается ровно неверный вывод. Он не показывает, что блокировки не было: выше я сам написал, что заблокированный INSERT подождёт и завершится успешно, а значит 959 мс совместимы и с блокировкой, и с посторонней задержкой. Он не показывает и обратного. Единственное, что из него следует: без pg_locks и без привязки к фазе переноса конкретного шарда измерение неоднозначно — а снятая эвристика на этих данных уверенно объявила бы блокировку. То же относится и к прогону job_id=31 с максимумом записи 2668 мс: цифра крупная, вывода из неё никакого.
Сырой результат такой: отказов и потерь видимости не было ни в одном прогоне, максимумы гуляют в обе стороны, природа задержек не установлена. В длительность входят docker exec и подключение, а не только сам INSERT, и ни один отсчёт не соотнесён с фазой переноса конкретного шарда. Списать это на шум хоста было бы удобно, но такого замера я не делал — значит это тоже была бы догадка.
Мимоходом там же нашлась ловушка в счёте: план перемещений вернул тридцать строк, а Citus запланировал двадцать перемещений. Расхождения нет — customers и orders лежат в одной группе колокации, и парные шарды переезжают вместе, одним перемещением на пару. Строка плана и перемещение — разные величины, как задача и узел двумя разделами выше.
Отдельно замечу про соблазн набрать наблюдений побольше. В одном из ранних прогонов этого стенда набралось 1200 успешных проб — но лишь потому, что ребаланс тогда завис, и почти всё это время шарды не двигались. Большая выборка измеряла простой. Скажу точнее, чем хотелось бы: два десятка проб покрывают время жизни задачи от постановки до первого наблюдения завершения — а это не только перенос шардов. Туда входит и ожидание планировщика, и идущая первой задача копирования справочника, и возможные паузы между перемещениями; за прогрессом конкретных перемещений скрипт не следит. Доказано лишь, что многоминутного зависания в этом окне не было. И всё же два десятка проб по работающей задаче честнее тысячи, снятых с зависшей.
Что смотреть в проде
Четыре вещи, каждая из которых обнаружилась при сборке стенда и стоила времени.
Ключ шардирования обязан входить в каждое ограничение уникальности. Citus не даст распределить таблицу, если её первичный ключ, UNIQUE или EXCLUDE не включает колонку шардирования:
ERROR: cannot create constraint on "shipments"
DETAIL: Distributed relations cannot have UNIQUE, EXCLUDE, or PRIMARY
KEY constraints that do not include the partition columnДекларативно обойти нельзя — распределённого UNIQUE без колонки шардирования не будет. Прикладные варианты остаются: проверка уникальности в коде или генерация заведомо глобально уникальных идентификаторов (UUID, снежинки), — но за них база уже не отвечает. На практике это означает, что при шардировании существующей таблицы уникальность придётся пересмотреть — и, возможно, не только в базе, но и в приложении, если исходный уникальный идентификатор не совпадает с ключом шардирования.
Неблокирующий ребаланс потребовал wal_level = logical — и без него не упал, а завис. Оговорка про «неблокирующий» существенная: Citus умеет переносить шарды двумя способами. По умолчанию (shard_transfer_mode = 'auto') идёт логическая репликация — она позволяет обойтись без блокировки записи, и именно она требует logical. Есть и второй режим, block_writes: там перенос делается через COPY с блокировкой записи, и требования к wal_level нет. Различие между режимами — про записи; читать можно в обоих. То есть это цена конкретного режима, а не свойство ребаланса вообще. Стенд режим не менял, так что проверено только про auto.
Со стандартным wal_level = replica на всех сервисах фоновая задача уходила в затяжной повтор с ошибкой про логическое декодирование и блокировала собой все запланированные перемещения.
Внешне это выглядит как «ребаланс идёт», а на деле не двигается ничего. Сколько задача повторяется, я не выяснял — ждать до конца не стал; наблюдал повторы без единого сдвинутого шарда, а «бесконечно» тут было бы догадкой.
Про «на каких именно узлах» скажу отдельно, потому что сначала написал шире, чем измерил. В стенде logical включён на всех сервисах разом, и сравнивались ровно два состояния: replica везде против logical везде. Такой опыт не разделяет роли — ни координатор от воркеров, ни отправителя от получателя. А разделять есть что: PostgreSQL требует wal_level = logical на стороне публикации, к подписчику такого требования нет, и какой из воркеров становится публикующим в конкретной операции переноса, стенд не проверяет. Поэтому честная формулировка узкая: для стенда logical включён на всех сервисах консервативно, а минимально необходимый набор ролей этим экспериментом не установлен.
Слив узла асинхронен. citus_drain_node ставит фоновую задачу и возвращает управление сразу — снятие узла следующей же командой отказывает, шарды ещё на месте. Хуже: слив помечает узел как «не должен принимать шарды», и если в этот момент выполняется ребаланс, планировавший перенос на этот узел, его задачи начинают падать с сообщением Moving shards to a node that shouldn't have a shard is not supported. Наблюдал на Citus 14.1 одиннадцать повторов за семь минут без единого сдвинутого шарда — дальше прервал, так что это нижняя граница наблюдения, а не доказанная бесконечность. Лечится отменой зависшей задачи через citus_job_cancel.
Уборку после слива стоит звать, пока узел ещё жив. citus_drain_node оставляет в pg_dist_cleanup отложенную запись на удаление старой копии перенесённого шарда — сам перенос её не убирает. Дальше — конкретный проверенный сценарий на Citus 14.1, из которого я не берусь делать общего правила. Я снял регистрацию узла и погасил контейнер до того, как запись обработали; после этого запись осталась в каталоге, и явный CALL citus_cleanup_orphaned_resources();, вызванный уже после удаления узла, её не разобрал — node_group_id в ней указывает на несуществующий узел. Вызов уборки после citus_drain_node, но до citus_remove_node в том же сценарии сработал. Перебирать варианты отказа я не стал, так что читать это стоит как «в этом порядке точно работает, в обратном — застряло», а не как описание всех возможных путей.
Сюда же — деталь, которая стоила стенду пяти минут на каждом прогоне: citus_drain_node не трогает копии референсных таблиц. Слив переносит только размещения распределённых таблиц, а копия справочника остаётся на узле до citus_remove_node. Цикл ожидания, который ждал опустошения узла от всех размещений, ждал структурно недостижимого нуля.
Отсюда набор для мониторинга: состояние фоновых задач ребаланса и число повторов у них, распределение размещений по узлам, значение pg_dist_cleanup после каждого слива, и — если пользуетесь референсными таблицами — совпадает ли число их копий с числом активных воркеров (не узлов кластера: координатор в этот счёт не входит).
Где проходит граница
Всё описанное выше — цена, которую платят после того, как решение шардировать принято. Само решение — отдельный разговор, и до него стоит исчерпать более дешёвые шаги: вертикальное масштабирование, реплики для чтения, партиционирование внутри одного узлаготовится, с 2 октября.
Выбор ключа шардирования, хотспоты и ре-шардинг разобраны в статье про шардирование MongoDB — там всё это показано на живом кластере, и повторять не буду. Механика распределённых транзакций между шардами — в статье про сильную и итоговую согласованность.
Отдельная развилка — не шардировать руками вовсе, а взять базу, которая делает это сама. Это Distributed SQLСкоро, и выбор между ним и middleware стоит того, чтобы считать его отдельно. В классе middleware Citus не одинок: для MySQL ту же нишу занимает Vitess — здесь я его не разбираю, но при выборе он в списке.
Механику самого разбиения — как ключ превращается в номер шарда и почему добавление узла не должно перетасовывать всё — разбирает статья про consistent hashingСкоро.
Что до самого Citus: расширение версии 14.1 работает с PostgreSQL 16, 17 и 18 — это заявлено в release notes 14.1, и готовый образ под восемнадцатую ветку тоже есть. Тут легко ошибиться, и я ошибся: суффиксированные теги 14.1.0-pg16 и 14.1.0-pg17 перечисляют только старшие ветки, а PostgreSQL 18 лежит в основном теге без суффикса — у citusdata/citus:14.1.0 внутри PG_MAJOR=18. Отсутствие тега -pg18 означает ровно обратное тому, на что похоже. Стенд пинован на 14.1.0-pg17 сознательно, отсюда семнадцатая ветка в замерах.
Отставание у Citus есть, но не в поддержке движка, а в документации: docs.citusdata.com и под stable, и под latest отдаёт 13.0.1, а версионированного пути к 14.x попросту нет. Сверяться придётся с описанием предыдущей мажорной версии — вот это и стоит закладывать, планируя обновление.
Демо и версии
Стенд: digital-cookbook/databases/citus — координатор и два воркера в compose, схема с со-расположенными, референсной и по-разному распределёнными таблицами, пять демонстраций из этой статьи. Полный прогон — около трёх минут.
Каждая демонстрация печатает падающий вариант: что наблюдалось бы, если бы измеряемого эффекта не было, — и умеет сама объявить провал вместо выдачи правдоподобных чисел. Это не украшение. При сборке стенда так нашлись четыре дефекта, каждый из которых давал верное на вид, но бессмысленное измерение: демонстрация референсных таблиц повторяла результат демонстрации колокации, не различая два разных механизма; пагинация мерила эффект на объёме, где он физически не может проявиться; проверка ребаланса объявляла потерей данных появление той самой копии справочника, которую сама же и демонстрировала; а проба доступности била в неподвижный шард, доказывая совсем не то, что было заявлено.
Последний дефект нашёл не я, а внешнее ревью — и он поучительнее остальных. Проба возвращала честные успешные ответы, гейты были зелёными, числа воспроизводились. Сломано было не измерение, а его связь с утверждением: скрипт нигде не проверял, что читаемые данные действительно переезжают. Такое не ловится перепрогоном, потому что перепрогон подтверждает ровно ту же подмену. Теперь эту связь удерживают два гейта: один проверяет, что читаемый шард стоит в плане перемещений, второй — что после завершения задачи он действительно оказался на обещанном узле.
Версии на момент прогона: Citus 14.1, PostgreSQL 17.6, образ citusdata/citus:14.1.0-pg17.
Все числа сняты на одной машине, где узлы кластера живут контейнерами на общем хосте. Переносить стоит конкретные инварианты, проверенные при фиксированных условиях: число задач у запроса по ключу и без него, наличие или отсутствие стадий переброски, точное равенство поднятых строк формуле шарды × (OFFSET + LIMIT), сохранение числа размещений распределённых таблиц при ребалансе. «Структурные величины вообще» переносить нельзя — часть из них прямо зависит от топологии: число копий референсной таблицы равно числу воркеров и меняется вместе с ним, а направления по времени зависят ещё и от сети, параллелизма и work_mem. Абсолютные миллисекунды на настоящем кластере будут другими: добавится сетевой объём и удалённые взаимодействия, которых здесь почти нет. Каким станет разрыв между «один шард» и «все шарды» — этот стенд не измеряет, и гадать я не буду.
Документация
- Citus 14.1: release notes — основной официальный источник по версии, на которой снят стенд
Остальное — документация ветки 13.0: под stable и под latest docs.citusdata.com отдаёт именно её, описания 14.x на момент написания нет. Расхождения с поведением 14.1 возможны, поэтому сверяйтесь с release notes выше.
- Citus 13.0: распределённые таблицы и колокация
- Citus 13.0: референсные таблицы
- Citus 13.0: ребаланс шардов и
citus_rebalance_start - Citus 13.0:
citus_drain_node - Citus 13.0: ограничения на уникальность в распределённых таблицах
- PostgreSQL:
wal_levelи логическая репликация - PostgreSQL:
work_memи методы сортировки
Комментарии