Коротко:
- Партиционирование делит таблицу внутри одного узла, шардирование - распределяет данные по разным физическим серверам.
- Пора думать о горизонтальном масштабировании, когда вертикальный апгрейд уже не даёт нужного прироста или стоит неоправданно дорого.
- Ключ шардирования определяет, как данные распределятся по узлам: плохой выбор приводит к 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().
Чеклист перед внедрением
- Исчерпаны ли возможности вертикального масштабирования и оптимизации запросов?
- Определён ли основной паттерн доступа к данным (какие запросы самые частые)?
- Выбранный ключ имеет высокую кардинальность и равномерное распределение значений?
- Большинство критичных запросов попадает в один шард (не требует scatter-gather)?
- Проверена ли монотонность ключа - нет ли риска hotspot-а по времени или автоинкременту?
- Спроектировано ли co-location для таблиц, которые часто JOIN-ятся?
- Есть ли план для кросс-шардовых агрегаций (отдельный аналитический слой, ETL)?
- Понятно ли, как будет работать балансировщик при добавлении новых узлов?
- Протестирована ли схема на реальном распределении данных (не синтетическом)?
- Есть ли мониторинг per-shard метрик (CPU, IOPS, количество записей на шард)?