Шардирование базы данных: когда партиционирования уже недостаточно и как выбрать ключ шардирования

Шардирование базы данных: когда партиционирования уже недостаточно и как выбрать ключ шардирования

Коротко:

  • Партиционирование делит таблицу внутри одного узла, шардирование - распределяет данные по разным физическим серверам.
  • Пора думать о горизонтальном масштабировании, когда вертикальный апгрейд уже не даёт нужного прироста или стоит неоправданно дорого.
  • Ключ шардирования определяет, как данные распределятся по узлам: плохой выбор приводит к hotspot-ам - одни шарды перегружены, другие простаивают.
  • Consistent hashing снижает объём перебалансировки при добавлении новых узлов.
  • Кросс-шардовые запросы и распределённые транзакции - самые болезненные последствия неудачной схемы разбиения.

Почему партиционирования в какой-то момент перестаёт хватать

Представьте таблицу заказов на PostgreSQL: 800 миллионов строк, партиции по месяцам, индексы настроены. Запросы по конкретному периоду работают быстро. Но база живёт на одном сервере, и как только нагрузка на запись вырастает до десятков тысяч операций в секунду, упираешься не в схему - а в железо. Один диск, одна оперативная память, один процессор.

Партиционирование - это организационный приём внутри одного инстанса. Оно помогает планировщику запросов отсекать ненужные блоки данных, упрощает архивирование и ускоряет DELETE старых партиций. Но физически все данные по-прежнему на одной машине. Когда этот сервер исчерпал возможности вертикального масштабирования, нужен другой подход.

Шардирование базы данных решает именно эту задачу: данные разбиваются на независимые части - шарды - и размещаются на разных серверах. Каждый шард хранит свой срез данных и обрабатывает только свои запросы. Совокупная пропускная способность системы растёт вместе с количеством узлов.

Это не бесплатно. Вместе с горизонтальным масштабированием приходят новые проблемы: кросс-шардовые JOIN-ы, распределённые транзакции, сложность операций. Поэтому шардировать стоит осознанно, а не потому что «у всех так».

Партиционирование vs шардирование: ключевые различия

Путаница между двумя понятиями встречается часто. Внешне они похожи: оба делят данные на части по какому-то критерию. Но уровень, на котором это происходит, принципиально разный.

ХарактеристикаПартиционированиеШардирование
Где живут данныеОдин сервер, разные файлы/табличные пространстваРазные физические узлы
Кто управляет разбиениемСУБД прозрачно для приложенияПриложение или слой маршрутизации
Масштабирование записиНет - всё идёт на один инстансДа - каждый шард принимает свой поток
JOIN между частямиПрозрачный, оптимизированный планировщикомТребует координации между узлами
Сложность операцийСтандартная SQL-семантикаОграниченные транзакции, сложные агрегации

В PostgreSQL партиционирование встроено нативно начиная с версии 10: declarative partitioning по диапазонам, спискам или хэшу. Шардирование же обычно реализуется через расширения вроде Citus, через промежуточные прокси (pgpool-II, Pgbouncer с маршрутизацией) или переносится на уровень приложения.

MongoDB шардирует данные нативно через mongos - роутер, который направляет запросы к нужному шарду на основе shard key. Конфигурация хранится в config server replica set.

Когда реально стоит переходить к распределённому хранению

Нет универсального порога. Но есть признаки, которые говорят, что пора серьёзно об этом думать:

  • Вертикальный апгрейд больше не окупается. Переход с 32 до 64 ядер или удвоение RAM дало 10-15% прироста вместо ожидаемых 50%. Закон убывающей отдачи работает.
  • Запись стала узким местом. Читать можно масштабировать репликами. Когда не справляется запись - реплики не помогут. Это первый сигнал, что нужна горизонталь.
  • Объём данных выходит за пределы одного диска или разумного SSD-бюджета. Хранить 10 ТБ на одном nvme можно, но ценник быстро становится неадекватным.
  • Время бэкапа и восстановления стало слишком большим. Если pg_dump одной базы занимает 8 часов, RTO при аварии выглядит пугающе.
  • Запросы деградируют даже с хорошими индексами. Планировщик честно читает 200 миллионов строк после прунинга партиций - и всё равно медленно.

Если из пяти пунктов совпадает два-три, архитектурный разговор стоит начать. Если только один - сначала попробуйте партиционирование, кэш и оптимизацию запросов.

Как выбрать ключ шардирования: критерии и антипаттерны

Выбор ключа (sharding key) - самое важное решение во всей схеме. Его почти невозможно поменять без полного перераспределения данных. Ошиблись на старте - будете жить с последствиями годами.

Хороший ключ: что это значит на практике

Хороший ключ равномерно распределяет данные по узлам и при этом позволяет большинству запросов обходиться одним шардом. Эти два требования часто противоречат друг другу, и здесь начинается настоящая инженерная задача.

Высокая кардинальность. Ключ должен иметь достаточно различных значений, чтобы данные не сгрудились в нескольких точках. Булевый флаг или поле из трёх значений - ужасный выбор: весь трафик пойдёт на два-три шарда.

Равномерность распределения. Ключ не должен создавать скосы. Классический hotspot - время создания записи в сервисе с пиковой нагрузкой. Все новые заказы попадают в «горячий» шард текущего месяца, остальные простаивают.

Локализация запросов. Большинство рабочих запросов должно попадать на один шард. Если типичный сценарий - получить все заказы конкретного пользователя, то user_id как ключ позволяет обслужить этот запрос с одного узла. Если же нужно агрегировать заказы по категориям - user_id плохой выбор, придётся опрашивать все шарды.

Стабильность. Значение ключа не должно меняться. Если ключ - email пользователя, а пользователь может его поменять, при обновлении запись придётся физически переместить на другой шард. Это дорого и сложно.

Типичные ошибки при выборе ключа

Монотонно возрастающий ключ (auto-increment ID, временная метка). Новые записи всегда попадают на последний шард. Он перегружен по записи, остальные - недогружены. Это классический hotspot, который убивает смысл горизонтального масштабирования.

Составной ключ с плохой первой частью. В MongoDB shard key составной? Проверьте, по какому полю чаще фильтруют. Если первое поле составного ключа имеет низкую кардинальность - распределение будет плохим.

Ключ, который не входит в большинство запросов. Если шард-ключ - регион пользователя, а большинство запросов фильтрует по product_id без региона, каждый такой запрос пойдёт ко всем шардам одновременно (scatter-gather). При 10 шардах нагрузка возрастает в 10 раз.

Слишком «горячие» значения. Представьте маркетплейс, где 5% продавцов генерируют 80% трафика. Если ключ - seller_id, эти продавцы создадут hotspot-ы независимо от общего числа шардов. Нужна дополнительная логика: либо пре-разбиение горячих продавцов на несколько виртуальных сегментов, либо другой ключ.

Стратегии разбиения: range, hash и directory

После выбора поля нужно решить, как именно значения этого поля маппятся на шарды.

Range sharding. Шард 1 хранит user_id от 1 до 1 000 000, шард 2 - от 1 000 001 до 2 000 000 и так далее. Удобно для диапазонных запросов: если запрашиваете данные конкретного пользователя, сразу ясно, на какой узел идти. Проблема: при монотонном росте ID все записи снова попадают на последний шард.

Hash sharding. Значение ключа хэшируется, остаток от деления на число шардов определяет узел: shard = hash(key) % n. Распределение равномерное, hotspot-ы исчезают. Но диапазонные запросы теперь требуют опроса всех шардов - планировщик не может предсказать, где лежат записи с ID от 100 до 200.

Directory sharding. Отдельная таблица маппинга: значение ключа явно указывает на конкретный шард. Максимально гибко: можно перемещать конкретные значения между узлами без изменения схемы. Цена - лишний lookup при каждом запросе и единая точка отказа в виде таблицы маппинга.

Consistent hashing: зачем нужен и как работает

Обычный хэш-шардинг с формулой hash(key) % n имеет неприятное свойство: при добавлении одного нового шарда меняется n, а значит переопределяется маппинг почти всех существующих ключей. Придётся перемещать огромный объём данных.

Consistent hashing решает это иначе. Представьте числовое кольцо от 0 до 2^32. Каждый шард занимает несколько точек на этом кольце. Ключ хэшируется в точку кольца и маппится на ближайший шард по часовой стрелке. При добавлении нового узла он забирает часть диапазона только у своего соседа - остальные шарды не затронуты. В среднем перемещается только 1/n часть данных.

Cassandra, DynamoDB и Redis Cluster используют именно этот принцип. В Cassandra каждый узел отвечает за диапазон токенов на кольце; виртуальные узлы (vnodes) дополнительно улучшают равномерность распределения.

Пример: В кластере из 4 узлов добавляете 5-й. При обычном хэше нужно переместить ~80% данных (так как меняется формула для почти всех ключей). С consistent hashing - около 20%, то есть только те записи, которые попадали на «соседа» нового узла.

Кросс-шардовые запросы и их цена

Это главная боль после внедрения горизонтального масштабирования. Запросы, которые не попадают в один шард, требуют опроса нескольких узлов, получения частичных результатов и их объединения на стороне приложения или координатора.

Scatter-gather запрос против 8 шардов создаёт в 8 раз больше сетевых вызовов. Если добавить сортировку и LIMIT, координатор должен получить N*limit строк от каждого шарда, отсортировать их и вернуть финальные limit строк. При больших OFFSET это становится катастрофой.

JOIN между таблицами, которые шардированы по разным ключам - ещё хуже. Либо придётся делать broadcast join (отправить одну таблицу целиком на каждый шард), либо переносить логику соединения в приложение.

Практическое правило: спроектируйте шарды так, чтобы данные, которые часто запрашиваются вместе, лежали вместе. Это называется co-location. Если заказы и детали заказов шардированы по одному и тому же ключу (например, order_id с одинаковым хэшом), JOIN между ними никогда не пойдёт на соседний узел.

Hotspot-ы: как диагностировать и что делать

Hotspot - ситуация, когда один шард принимает непропорционально большую нагрузку. Это буквально уничтожает смысл масштабирования: весь трафик всё равно идёт в одну точку.

Как обнаружить. Смотрите на метрики CPU, disk I/O и количество операций в секунду по каждому шарду отдельно. Если один узел на 90% по записи, а остальные на 10-20% - это hotspot. В MongoDB такую картину видно через db.adminCommand({serverStatus: 1}) на каждом mongod и через встроенный mongos balancer stats.

Причины:

  • Монотонно возрастающий ключ - все новые записи идут в один шард.
  • Несколько «звёздных» сущностей с огромным числом связанных записей (крупный продавец, популярный хэштег).
  • Время как часть ключа в сервисах с временной локальностью чтения - пользователи читают свежие данные, и шард с последними записями перегружен и по чтению тоже.

Решения:

  • Добавить искусственный суффикс к ключу: вместо seller_id использовать seller_id + random_bucket (например, от 0 до 9). Запись распределяется по 10 виртуальным сегментам. При чтении нужно опросить все 10 и объединить результаты - компромисс между записью и чтением.
  • Переосмыслить ключ: если hotspot системный, значит выбор поля был ошибочным изначально.
  • Использовать виртуальные шарды (vnodes) с consistent hashing - узлы получают несколько диапазонов на кольце, что сглаживает пики.

Шардирование в PostgreSQL: практика через Citus

Нативного шардирования в PostgreSQL нет. Расширение Citus (сейчас часть Microsoft Azure) превращает Postgres в распределённую БД: есть coordinator node, который принимает запросы, и worker nodes, которые хранят данные.

Таблица размечается как distributed:

SELECT create_distributed_table('orders', 'user_id');

Citus по умолчанию создаёт 32 шарда (параметр citus.shard_count) и распределяет их по workers через consistent hashing. При добавлении нового worker запускается rebalancer, который перемещает шарды.

Reference tables - небольшие таблицы (справочники, конфиги), которые реплицируются на все workers. Это решает проблему JOIN-а со справочными данными: каждый worker может соединить свой срез с локальной копией справочника без сетевых вызовов.

Важный нюанс в Citus: запросы с фильтром по ключу шардирования идут напрямую на нужный worker. Запросы без него - broadcast на все workers. Это нужно учитывать при проектировании индексов и типичных паттернов доступа.

Шардирование в MongoDB: shard key и зоны

MongoDB поддерживает шардирование нативно. Роутер mongos принимает запросы от приложения и маршрутизирует их к нужным шардам на основе shard key. Метаданные о распределении чанков хранятся в config server replica set.

Данные делятся на chunks (по умолчанию 128 МБ). Балансировщик автоматически перемещает чанки между шардами, чтобы выровнять их количество.

Выбор shard key в MongoDB необратим без resharding (появилось в версии 5.0, но дорого по ресурсам). До 5.0 изменение ключа требовало полной миграции данных.

Зоны (zone sharding) позволяют привязать диапазоны ключей к конкретным шардам. Это полезно для geo-locality: европейские пользователи хранятся на шардах в европейском датацентре. Настраивается через sh.addTagRange().

Чеклист перед внедрением

  1. Исчерпаны ли возможности вертикального масштабирования и оптимизации запросов?
  2. Определён ли основной паттерн доступа к данным (какие запросы самые частые)?
  3. Выбранный ключ имеет высокую кардинальность и равномерное распределение значений?
  4. Большинство критичных запросов попадает в один шард (не требует scatter-gather)?
  5. Проверена ли монотонность ключа - нет ли риска hotspot-а по времени или автоинкременту?
  6. Спроектировано ли co-location для таблиц, которые часто JOIN-ятся?
  7. Есть ли план для кросс-шардовых агрегаций (отдельный аналитический слой, ETL)?
  8. Понятно ли, как будет работать балансировщик при добавлении новых узлов?
  9. Протестирована ли схема на реальном распределении данных (не синтетическом)?
  10. Есть ли мониторинг per-shard метрик (CPU, IOPS, количество записей на шард)?

Часто задаваемые вопросы

Репликация создаёт копии одних и тех же данных на нескольких узлах - для отказоустойчивости и масштабирования чтения. При шардировании каждый узел хранит разные данные - это масштабирование записи и объёма. На практике оба подхода используются вместе: каждый шард имеет свои реплики.

Да, и многие так делают. Логика маршрутизации живёт в коде: при вставке или запросе вычисляется шард по ключу, соединение открывается к нужному инстансу. Это даёт полный контроль, но усложняет код и перекладывает всю ответственность за балансировку на команду. Специализированные инструменты (Citus, Vitess, mongos) берут эту сложность на себя.

Resharding - изменение ключа шардирования или перераспределение данных по новым шардам. Это дорогая операция: нужно прочитать, переместить и переиндексировать большой объём данных. MongoDB 5.0+ поддерживает online resharding без остановки, но это всё равно нагрузочная процедура на часы или дни. Именно поэтому ключ нужно выбирать правильно с самого начала.

Если транзакция затрагивает данные на одном шарде - всё стандартно. Если нужна атомарность между несколькими шардами, требуется распределённая транзакция (2PC или SAGA-подобный паттерн). Это медленнее и сложнее. MongoDB поддерживает multi-document transactions, включая cross-shard с версии 4.2, но с ограничениями по производительности. Лучший подход - проектировать схему так, чтобы транзакции не пересекали границы шардов.

Когда кластер часто меняет состав: добавляются узлы при росте нагрузки или выводятся при её снижении. Формула hash % n при каждом изменении n требует массового перемещения данных. Consistent hashing ограничивает объём перебалансировки долей 1/n, что критично при больших объёмах.

Небольшие таблицы-справочники лучше реплицировать на все шарды (reference tables в Citus, broadcast tables в других системах). Шардировать имеет смысл только таблицы с большим объёмом записей и высокой нагрузкой. Шардирование маленькой таблицы только добавит сложности без выгоды.

Итог

Горизонтальное разбиение хранилища решает задачи, которые вертикальный апгрейд и партиционирование закрыть не могут: масштабирование записи, хранение объёмов за пределами одного сервера, снижение нагрузки на отдельные узлы. Но это решение, которое нельзя отменить по щелчку - особенно в части выбора ключа.

Правильный ключ равномерно распределяет данные и при этом позволяет большинству запросов не выходить за пределы одного шарда. Hotspot-ы, кросс-шардовые агрегации и распределённые транзакции - это не абстрактные риски, а конкретные проблемы, которые придётся решать, если схема выбрана небрежно. Consistent hashing и co-location - инструменты, которые снижают эту боль, но не устраняют её полностью.

Прежде чем переходить к шардированию, убедитесь, что вы исчерпали более простые варианты: оптимизацию запросов, партиционирование, кэш и добавление реплик для чтения. Если всё это уже сделано и узкое место остаётся - тогда да, время думать о горизонтали.