Когда данных или нагрузки становится больше, чем тянет один сервер, их разрезают на части — секции (partitions), они же шарды. Правило простое: каждая запись живёт ровно в одной секции, а секции раскладывают по разным узлам.
Со словами тут путаница, и о неё легко споткнуться. В разговоре про распределённые системы «секция» и «шард» — одно и то же: кусок данных, уехавший на свой узел. А внутри одной базы то же слово значит другое: партиционирование в PostgreSQL режет таблицу на куски на том же сервере, и шардированием там называют только разъезд по разным серверам. В этой статье речь всегда про второе — про куски на разных узлах. Практику шардинга мы разбираем на 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 — «обычно доли секунды, при сбоях дольше»). Следствие важнее, чем кажется: глобальный индекс может вернуть документ, которого уже нет, и не вернуть только что записанный. Поэтому проверки вида «если такой записи нет — создаём» на нём строить нельзя, только на первичном ключе.
Один и тот же запрос «найти красные»: локальному индексу нужны все секции и склейка ответов, глобальному хватает одной секции индекса.
Секция — это не машина
Слово «узел» в разговоре о разбиении означает то одно, то другое, и путаница мешает. Разведём.
Секция — это логический кусок данных, определённый ключом (диапазон или хеш). Реплика — физическая копия секции. Узел — машина, на которой лежат несколько секций, обычно в разных ролях: для одних секций он ведущий, для других держит копию.
В живой системе эти две оси работают вместе: данные сначала режутся на секции по ключу, а каждая секция живёт в нескольких копиях на разных машинах. Отказ машины означает, что часть секций потеряла одну копию — их ведущие переезжают на другие машины, а данные остаются доступными. Добавление машины означает переезд части секций на неё.
Отсюда и практический вывод: число секций и число машин — разные числа, и первое обычно заметно больше второго. Именно это делает возможной балансировку: перекладывать секции по машинам проще, чем перерезать данные.
Сколько секций заводить
Ориентир идёт от двух величин. Снизу: секций должно быть заметно больше, чем машин (обычно в несколько раз, а в системах с фиксированным числом секций — на порядок), иначе балансировка становится грубой и добавление машины не даёт равномерности. Сверху: каждая секция стоит служебных расходов — соединения, метаданные, отдельные файлы, отдельные фоновые процессы; тысячи секций на скромном кластере означают, что заметная часть работы уходит на их обслуживание, а запросы без ключа обходят их все.
По объёму практический ориентир такой же, как у партиций в базе: держать секцию в пределах десятков гигабайт, чтобы её можно было перелить за разумное время. Слишком крупная секция — это долгая балансировка и невозможность разделить горячий кусок.
Многие системы решают это фиксированным числом секций, заданным при создании (тысячи), и дальше только раскладывают их по машинам. Плата — число выбрано навсегда.
Что ломается, когда данные разрезаны
Список стоит знать до того, как разрезать.
Соединения между секциями. Соединить две таблицы, разложенные по разным ключам, — значит переслать данные по сети. Поэтому связанные данные стараются класть в одну секцию (общий ключ секционирования), а справочники дублируют на все узлы.
Уникальность. Глобально уникальный индекс поверх секций требует проверки на всех узлах при каждой вставке. Обычно уникальность обеспечивают в пределах секции, а глобальную — либо отдельным сервисом, либо выбором ключа так, чтобы она сводилась к локальной.
Транзакции на несколько секций. Их либо нет вовсе, либо они дороги (двухфазная фиксация). Отсюда правило: границы транзакции стараются уместить в одну секцию, а остальное делают согласованием в приложении.
Агрегаты и сортировка по всему набору. Любой запрос без ключа секционирования идёт на все секции и собирает результат на инициаторе (разослать и собрать). Это работает, но масштабируется плохо: чем больше секций, тем больше ожидание самой медленной.
Пагинация по всему набору. Глубокая страница по всем секциям означает, что каждая отдала свои кандидаты, а инициатор отсортировал и выбросил лишнее. Дорого и потому в интерфейсах почти всегда заменяется листанием от последней показанной записи.
Горячий ключ: второй ходспросят на собеседовании
Суффикс к ключу размазывает нагрузку, но делает чтение дороже (собирать надо со всех суффиксов). Часто дешевле не размазывать, а убрать горячий ключ из хранилища: положить его значение в кэш перед базой (популярный товар читают из Redis, а не из секции), или вынести в отдельное хранилище, рассчитанное на такую нагрузку.
Практическое правило: горячее чтение лечится кэшем, горячая запись — суффиксом или отдельным счётчиком с периодическим сведением.
Сменить ключ = перелить данные
Утверждение «менять ключ очень тяжело» стоит договорить: менять его нельзя вовсе — можно только переразложить данные по новому ключу. Технически это переезд: новая коллекция или таблица с новым ключом, перелив порциями, сверка, переключение. Некоторые системы умеют делать это на живых данных своим механизмом (в MongoDB это перераскладка коллекции), но и он означает перезапись всего объёма с фоновым догоном изменений.
Отсюда же и главный совет по выбору: ключ выбирают под запросы, которые уже есть, и под равномерность, которую можно проверить на реальных данных — а не «на будущее».
Где это применяется
Развилка «как шардировать» встаёт ровно тогда, когда один узел перестал справляться. И почти каждое решение здесь необратимо: ключ секционирования потом меняется очень тяжело, схему индексов после запуска не переиграешь.
Где спотыкаются начинающие:
- Ключ секционирования начинается с времени или автоинкремента — и вся свежая запись валится в одну секцию. Классическая горячая точка. На первое место ключа нужно что-то разнообразное.
- Ждут от вторичного индекса скорости первичного — а при локальных индексах любой запрос не по ключу секционирования превращается в scatter/gather по всем узлам.
- Забывают про горячий ключ — одна знаменитость или один популярный товар кладут секцию, хотя формально распределение выглядит ровным.
Глубже: ребалансировка: как двигать секции между узламирасширенное
Со временем узлы добавляют (нагрузка выросла) и убирают (сбой или вывод из эксплуатации). Перенос секций с одного узла на другой называют ребалансировкой, и от неё ждут трёх вещей: после — нагрузка снова ровная; во время — база продолжает читать и писать; переносится не больше данных, чем реально нужно.
- Как делать нельзя —
hash mod N. Соблазнительно назначать секцию по формулеhash(ключ) % N, где N — число узлов. Но стоит изменить N, и почти все ключи меняют секцию: переход с 10 на 11 узлов перетасует практически все данные. Ребалансировка становится неподъёмной. - Фиксированное число секций. Заводят секций сильно больше, чем узлов (например, 1000 секций на 10 узлов), и раздают по многу на узел. Новый узел просто «забирает» несколько секций у существующих. Двигаются при этом целые секции; их количество и привязка ключей не меняются — меняется только то, на каком узле какая секция лежит. Так делают Elasticsearch, Riak, Couchbase. Грабля — число секций надо угадать заранее: оно задаёт потолок роста, а если взять слишком большим, появятся лишние накладные расходы. Ту же идею Cassandra раскладывает кольцом: пространство хешей замкнуто в круг, каждый узел отвечает за несколько его участков (их называют виртуальными узлами), и новый узел забирает по кусочку у всех сразу, а не у одного соседа. Этот приём и называют согласованным хешированием.
- Динамическое число секций. Секция переросла порог — делится надвое; усохла — сливается с соседкой (совсем как узел в B-дереве). Число секций само подстраивается под объём данных. Так работают HBase и MongoDB (chunks). Грабля пустой базы: на старте секция одна, и вся запись летит в один узел, пока не случится первое разбиение (лечится предварительным разбиением, если распределение известно заранее).
«Почти все» — это не фигура речи, а счёт. Возьмём сто тысяч ключей и посмотрим, сколько из них сменит узел при переходе с десяти узлов на одиннадцать: сначала по формуле hash mod N, потом при тысяче фиксированных секций, которые раздаются узлам целиком.
живой пример
public class Rebalance {
public static void main(String[] args) {
int keys = 100000;
int partitions = 1000;
int target = partitions / 11;
int movedByMod = 0;
for (int i = 0; i < keys; i++) {
if (hash(i) % 10 != hash(i) % 11) movedByMod++;
}
int[] load = new int[10];
for (int p = 0; p < partitions; p++) load[p % 10]++;
boolean[] moved = new boolean[partitions];
for (int p = 0, taken = 0; p < partitions && taken < target; p++) {
if (load[p % 10] > target) { load[p % 10]--; moved[p] = true; taken++; }
}
int movedByPartition = 0;
for (int i = 0; i < keys; i++) {
if (moved[hash(i) % partitions]) movedByPartition++;
}
System.out.println("hash mod N: " + Math.round(100f * movedByMod / keys) + "% ключей переезжает");
System.out.println("1000 секций: " + Math.round(100f * movedByPartition / keys) + "% ключей переезжает");
}
static int hash(int key) {
return ("order-" + key).hashCode() & 0x7fffffff;
}
}
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
Хеш здесь тот же, что у строк в Java (31·h + символ в 32-битном целом), чтобы проценты совпали с остальными ветками.
живой пример
package main
import (
"fmt"
"math"
)
func hash(key int) int {
var h int32
for _, c := range fmt.Sprintf("order-%d", key) {
h = 31*h + int32(c)
}
return int(h & 0x7fffffff)
}
func main() {
keys, partitions := 100000, 1000
target := partitions / 11
movedByMod := 0
for i := 0; i < keys; i++ {
if hash(i)%10 != hash(i)%11 {
movedByMod++
}
}
load := make([]int, 10)
for p := 0; p < partitions; p++ {
load[p%10]++
}
moved := make([]bool, partitions)
for p, taken := 0, 0; p < partitions && taken < target; p++ {
if load[p%10] > target {
load[p%10]--
moved[p] = true
taken++
}
}
movedByPartition := 0
for i := 0; i < keys; i++ {
if moved[hash(i)%partitions] {
movedByPartition++
}
}
fmt.Printf("hash mod N: %.0f%% ключей переезжает\n", math.Round(100*float64(movedByMod)/float64(keys)))
fmt.Printf("1000 секций: %.0f%% ключей переезжает\n", math.Round(100*float64(movedByPartition)/float64(keys)))
}
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
Хеш здесь тот же, что у строк в Java (31·h + символ через Math.imul), чтобы проценты совпали с остальными ветками.
живой пример
function hash(key) {
let h = 0;
for (const c of `order-${key}`) h = (Math.imul(31, h) + c.charCodeAt(0)) | 0;
return h & 0x7fffffff;
}
const keys = 100000, partitions = 1000, target = Math.floor(partitions / 11);
let movedByMod = 0;
for (let i = 0; i < keys; i++) if (hash(i) % 10 !== hash(i) % 11) movedByMod++;
const load = new Array(10).fill(0);
for (let p = 0; p < partitions; p++) load[p % 10]++;
const moved = new Array(partitions).fill(false);
for (let p = 0, taken = 0; p < partitions && taken < target; p++) {
if (load[p % 10] > target) { load[p % 10]--; moved[p] = true; taken++; }
}
let movedByPartition = 0;
for (let i = 0; i < keys; i++) if (moved[hash(i) % partitions]) movedByPartition++;
console.log(`hash mod N: ${Math.round(100 * movedByMod / keys)}% ключей переезжает`);
console.log(`1000 секций: ${Math.round(100 * movedByPartition / keys)}% ключей переезжает`);
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
Хеш здесь тот же, что у строк в Java (31·h + символ с обрезкой до 32 бит), чтобы проценты совпали с остальными ветками.
живой пример
def hash_key(key: int) -> int:
h = 0
for c in f"order-{key}":
h = (31 * h + ord(c)) & 0xFFFFFFFF
if h >= 1 << 31:
h -= 1 << 32
return h & 0x7FFFFFFF
keys, partitions = 100_000, 1000
target = partitions // 11
moved_by_mod = sum(1 for i in range(keys) if hash_key(i) % 10 != hash_key(i) % 11)
load = [0] * 10
for p in range(partitions):
load[p % 10] += 1
moved = [False] * partitions
taken = 0
for p in range(partitions):
if taken >= target:
break
if load[p % 10] > target:
load[p % 10] -= 1
moved[p] = True
taken += 1
moved_by_partition = sum(1 for i in range(keys) if moved[hash_key(i) % partitions])
print(f"hash mod N: {round(100 * moved_by_mod / keys)}% ключей переезжает")
print(f"1000 секций: {round(100 * moved_by_partition / keys)}% ключей переезжает")
Запустить
Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →
Девяносто один процент против девяти. Фиксированные секции переносят ровно ту долю данных, которая нужна новому узлу, — примерно одну одиннадцатую от всего.
Добавили одиннадцатый узел: при hash mod N переезжает 91% ключей, а при тысяче фиксированных секций двигаются целые секции и только 9% ключей.
Кто даёт команду на перенос
И отдельный вопрос — автоматика против ручного управления. Полностью автоматическая ребалансировка удобна, но в паре с автоматическим обнаружением сбоев опасна: перегруженный узел отвечает медленно → система решает, что он «отказал» → затевает ребалансировку → добавляет нагрузки уже перегруженному узлу → каскадный сбой. Спасает тут не отказ от автоматики, а развод двух решений: обнаружить, что узел молчит, машина может сама, а вот команду «переносим данные» лучше отдать человеку («система предлагает — админ подтверждает», как в Couchbase). Медленнее, зато без каскада.
Глубже: маршрутизация запросов: где лежит «foo»?расширенное
Секции разъезжаются по узлам и переезжают при ребалансировке. Откуда клиент знает, к какому узлу идти за ключом foo? Это частный случай общей задачи «найти, кто за что отвечает» (service discovery), и решений три:
- Клиент стучится в любой узел. Если ключ там — узел отвечает сам; если нет — сам перенаправляет запрос нужному узлу и возвращает ответ. Так делают Cassandra и Riak через gossip-протокол: узлы сами между собой обмениваются картой кластера, и внешний координатор не нужен.
- Отдельное маршрутизирующее звено. Все запросы идут через него, и оно знает, какой узел за что отвечает (в MongoDB это mongos). Сами запросы оно не обрабатывает — это балансировщик, который просто в курсе, где какая секция.
- Клиент сам знает раскладку и подключается к нужному узлу напрямую, без посредников.
Во всех трёх случаях кто-то должен знать актуальную карту «секция → узел» и вовремя узнавать об изменениях. Классическое решение — отдельный сервис-координатор ZooKeeper: узлы регистрируются в нём, маршрутизаторы подписываются на изменения, и когда секция переезжает, ZooKeeper всех оповещает. Так делают HBase и SolrCloud. Kafka тоже так работала, но с версии 3.3 научилась хранить карту сама, а в 4.0 зависимость от ZooKeeper убрали совсем — это общий тренд: внешний координатор добавляет ещё одну систему, которую надо эксплуатировать. Альтернатива — тот самый gossip (Cassandra, Riak): узлы устроены сложнее, зато нет зависимости от внешнего координатора.
Коротко
- Цель разбиения — раскидать по узлам и объём, и нагрузку. Перекос даёт горячую точку, и смысл шардинга теряется.
- По диапазону — быстрые сканы по порядку ключа, но время в начале ключа гонит весь свежий поток в одну секцию. Лечит составной ключ.
- По хешу — равномерность почти даром, платой уходит сортировка. Горячий ключ хеш не чинит: его размазывает приложение суффиксом, и платит за это чтение.
- Локальный вторичный индекс — дешёвая запись и scatter/gather на чтении, вариант по умолчанию. Глобальный — наоборот и обновляется асинхронно.
hash mod Nнепригоден: переход с 10 узлов на 11 двигает 91% ключей, тысяча фиксированных секций — 9%.- Карту «секция → узел» держит внешний координатор (ZooKeeper) или сами узлы через gossip. Авторебалансировку вместе с автообнаружением сбоев не включают.
- Секция, реплика и узел — три разные сущности: данные режут по ключу, каждую секцию держат в нескольких копиях, а машина несёт много секций в разных ролях.
- Секций заводят заметно больше, чем машин (иначе балансировка грубая), но каждая стоит служебных расходов; по объёму держат десятки гигабайт.
- Разрезанные данные ломают соединения между секциями, глобальную уникальность, транзакции на несколько секций и глубокую пагинацию — всё это становится «разослать и собрать».
- Горячее чтение лечат кэшем перед хранилищем, а не только суффиксом к ключу; сменить ключ секционирования нельзя — только перелить данные заново.
Что почитать дальше
- Репликация и шардинг MongoDB — shard key, chunks, balancer и mongos на практике.
- Модели репликации — секционирование почти всегда идёт в паре с репликацией.
- Партиционирование и шардирование в PostgreSQL — RANGE, LIST и HASH внутри одной базы.
- Строительные блоки — где шардинг живёт в общей картине.