Когда данных или нагрузки становится больше, чем тянет один сервер, их разрезают на части — секции (partitions, они же шарды). Правило простое: каждая запись живёт ровно в одной секции, а секции раскладывают по разным узлам. Практику шардинга мы разбираем на MongoDB — replica set, shard key, chunks, balancer. Эта статья — на уровень выше и без привязки к базе: как выбрать способ разбиения, что происходит со вторичными индексами, как двигать данные при добавлении узлов и как клиент вообще понимает, в какой секции лежит нужная запись. Вопросы эти одинаковы для MongoDB, Cassandra, Elasticsearch и HBase — меняются только названия.
Цель — равномерность, враг — горячая точка
Смысл секционирования один: раскидать данные и нагрузку по узлам поровну, чтобы десять узлов тянули в десять раз больше одного. Если разбиение вышло неравномерным — на часть секций пришлось непропорционально много данных или запросов, — его называют перекошенным (skewed), а перегруженную секцию — горячей точкой (hot spot). В худшем случае вся нагрузка ложится на один узел, а остальные девять простаивают — и весь смысл шардинга теряется.
Есть два базовых способа решить, в какую секцию отправить запись — по диапазону ключа и по хешу ключа. По сути это два разных ответа на один вопрос: как не получить горячую точку.
Разбиение по диапазону против разбиения по хешу
По диапазону значений ключа (range). Каждой секции достаётся непрерывный отрезок ключей — как тома бумажной энциклопедии: A–C, D–F и так далее. Границы подбирают под данные, потому что распределение неравномерно (слов на «А» и на «Ъ» — разное количество). Плюс такого разбиения — запросы по диапазону: внутри секции ключи лежат по порядку, и «все события за март» читаются одним быстрым сканом. Минус — тот же порядок легко рождает горячую точку. Если ключ начинается с метки времени, то вся сегодняшняя запись валится в одну секцию (сегодняшний диапазон), а вчерашние узлы стоят без дела. Лечится составным ключом, где время не на первом месте: сначала, скажем, имя датчика, и только потом время — тогда запись за один момент растекается по разным датчикам.
По хешу ключа (hash). Хорошая хеш-функция превращает даже похожие ключи в равномерно разбросанные числа, и секции достаётся диапазон не самих ключей, а их хешей. Равномерность получается почти даром — так делают Cassandra и MongoDB. Плата — теряется сортировка: соседние по смыслу ключи разлетаются по всем секциям, и запрос по диапазону теперь вынужден опросить все узлы. Cassandra идёт на компромисс — составной ключ: хешируется только первая часть (она и определяет секцию), а остальные части работают как отсортированный индекс уже внутри секции. Тогда «все сообщения одного пользователя за период» читаются эффективно, а сами пользователи разбросаны по кластеру равномерно.
Есть отдельная беда, которую не чинит ни хеш, ни диапазон, — горячий ключ. Например, знаменитость с миллионами подписчиков: весь трафик бьёт в одну её запись, а хеш от одного и того же ключа всегда один и тот же, так что «размазать» его автоматически нельзя. Тут помогает только приложение: приклеить к ключу случайный суффикс из сотни вариантов — тогда запись размажется по ста секциям. Но за это платит чтение: теперь, чтобы собрать всё по этому ключу, придётся опросить все сто секций и склеить. Поэтому приём применяют к горстке заведомо горячих ключей, а не ко всем подряд.
Вторичные индексы: локальные против глобальных
Пока обращаемся к данным по первичному ключу — всё просто: по нему и определяем секцию. Но приложению нужны и вторичные индексы: «найти все машины красного цвета», «все статьи со словом hogwash». Вторичный индекс не совпадает с разбиением на секции, и разрезать его можно двумя способами — с диаметрально противоположной ценой.
Локальный индекс (индекс по документу). Каждая секция ведёт вторичный индекс только по своим документам. Положили красную машину в секцию 3 — запись «цвет: красный» появилась в индексе секции 3, и больше нигде. Запись дешёвая: всё локально, задета одна секция. Зато чтение — это scatter/gather («разослать и собрать»): красные машины разбросаны по всем секциям, поэтому запрос «найди красные» приходится слать во все узлы и склеивать ответы. Так работают MongoDB, Cassandra, Elasticsearch. Разбросанное чтение удлиняет «хвост» задержки (ждём самый медленный узел из всех), но запись остаётся локальной — поэтому это вариант по умолчанию.
Глобальный индекс (индекс по терму). Тут индекс один на весь кластер, но и он сам тоже разрезан на секции — по искомому значению. Вся запись «цвет: красный» целиком лежит в одной секции индекса. Чтение быстрое — идём в одну секцию вместо всех. Зато плата переезжает на запись: добавление одного документа с несколькими полями задевает сразу несколько секций индекса (значения «цвет» и «марка» живут в разных местах), а держать это согласованным на лету — это распределённая транзакция. Поэтому глобальные индексы обновляют асинхронно: после записи изменение появляется в индексе с задержкой (так устроен DynamoDB — «обычно доли секунды, при сбоях дольше»).
Выбор ровно как везде: локальный индекс — дешёвая запись, дорогое чтение; глобальный — наоборот. У большинства систем по умолчанию — локальный.
Ребалансировка: как двигать секции между узлами
Со временем узлы добавляют (нагрузка выросла) и убирают (сбой или вывод из эксплуатации). Перенос секций с одного узла на другой называют ребалансировкой, и от неё ждут трёх вещей: после — нагрузка снова ровная; во время — база продолжает читать и писать; переносится не больше данных, чем реально нужно.
- Как делать нельзя —
hash mod N. Соблазнительно назначать секцию по формулеhash(ключ) % N, где N — число узлов. Но стоит изменить N, и почти все ключи меняют секцию: переход с 10 на 11 узлов перетасует практически все данные. Ребалансировка становится неподъёмной. - Фиксированное число секций. Заводят секций сильно больше, чем узлов (например, 1000 секций на 10 узлов), и раздают по многу на узел. Новый узел просто «забирает» несколько секций у существующих. Двигаются при этом целые секции; их количество и привязка ключей не меняются — меняется только то, на каком узле какая секция лежит. Так делают Elasticsearch, Riak, Couchbase. Грабля — число секций надо угадать заранее: оно задаёт потолок роста, а если взять слишком большим, появятся лишние накладные расходы.
- Динамическое число секций. Секция переросла порог — делится надвое; усохла — сливается с соседкой (совсем как узел в B-дереве). Число секций само подстраивается под объём данных. Так работают HBase, MongoDB (chunks), RethinkDB. Грабля пустой базы: на старте секция одна, и вся запись летит в один узел, пока не случится первое разбиение (лечится предварительным разбиением, если распределение известно заранее).
И отдельный вопрос — автоматика против ручного управления. Полностью автоматическая ребалансировка удобна, но в паре с автоматическим обнаружением сбоев опасна: перегруженный узел отвечает медленно → система решает, что он «отказал» → затевает ребалансировку → добавляет нагрузки уже перегруженному узлу → каскадный сбой. Человек в цикле («система предлагает — админ подтверждает», как в Couchbase и Riak) медленнее, зато спасает от таких сюрпризов.
Маршрутизация запросов: где лежит «foo»?
Секции разъезжаются по узлам и переезжают при ребалансировке. Откуда клиент знает, к какому узлу идти за ключом foo? Это частный случай общей задачи «найти, кто за что отвечает» (service discovery), и решений три:
- Клиент стучится в любой узел. Если ключ там — узел отвечает сам; если нет — сам перенаправляет запрос нужному узлу и возвращает ответ. Так делают Cassandra и Riak через gossip-протокол: узлы сами между собой обмениваются картой кластера, и внешний координатор не нужен.
- Отдельное маршрутизирующее звено. Все запросы идут через него, и оно знает, какой узел за что отвечает (в MongoDB это mongos). Сами запросы оно не обрабатывает — это балансировщик, который просто в курсе, где какая секция.
- Клиент сам знает раскладку и подключается к нужному узлу напрямую, без посредников.
Во всех трёх случаях кто-то должен знать актуальную карту «секция → узел» и вовремя узнавать об изменениях. Классическое решение — отдельный сервис-координатор ZooKeeper: узлы регистрируются в нём, маршрутизаторы подписываются на изменения, и когда секция переезжает, ZooKeeper всех оповещает. Так делают HBase, SolrCloud, Kafka. Альтернатива — тот самый gossip (Cassandra, Riak): узлы устроены сложнее, зато нет зависимости от внешнего координатора.
Где это применяется
Развилка «как шардировать» встаёт ровно тогда, когда один узел перестал справляться. И почти каждое решение здесь необратимо: ключ секционирования потом меняется очень тяжело, схему индексов после запуска не переиграешь. Практическая рамка: разбиение по диапазону — когда нужны сканы по времени или порядку и вы готовы бороться с горячей точкой составным ключом; по хешу — когда важнее равномерность; локальные вторичные индексы — вариант по умолчанию; глобальные — только если чтение по вторичному ключу критично, а асинхронное обновление индекса приемлемо.
Где спотыкаются начинающие:
- Ключ секционирования начинается с времени или автоинкремента — и вся свежая запись валится в одну секцию. Классическая горячая точка. На первое место ключа нужно что-то разнообразное.
- Ждут от вторичного индекса скорости первичного — а при локальных индексах любой запрос не по ключу секционирования превращается в scatter/gather по всем узлам.
- Назначают секции через
hash mod N— и первое же изменение числа узлов перетасует почти все данные. Нужно фиксированное или динамическое число секций. - Совмещают полностью автоматическую ребалансировку с автообнаружением сбоев — это рецепт каскадного отказа. Держите человека в цикле ребалансировки.
- Забывают про горячий ключ — одна знаменитость или один популярный товар кладут секцию, хотя формально распределение выглядит ровным.
Что почитать дальше: репликация и шардинг MongoDB — shard key, chunks, balancer и mongos на практике; модели репликации — секционирование почти всегда идёт в паре с репликацией; строительные блоки — где шардинг живёт в общей картине.