← назад к разделу

Когда данных или нагрузки становится больше, чем тянет один сервер, их разрезают на части — секции (partitions), они же шарды. Правило простое: каждая запись живёт ровно в одной секции, а секции раскладывают по разным узлам.

Со словами тут путаница, и о неё легко споткнуться. В разговоре про распределённые системы «секция» и «шард» — одно и то же: кусок данных, уехавший на свой узел. А внутри одной базы то же слово значит другое: партиционирование в PostgreSQL режет таблицу на куски на том же сервере, и шардированием там называют только разъезд по разным серверам. В этой статье речь всегда про второе — про куски на разных узлах. Практику шардинга мы разбираем на MongoDB — replica set, shard key, chunks, balancer. Эта статья — на уровень выше и без привязки к базе: те же вопросы одинаково встают в MongoDB, Cassandra, Elasticsearch и HBase, меняются только названия.

четыре записи от датчиков A, B, C и D за одну дату по диапазону: секцию выбирает начало ключа — дата секция 1 секция 2 секция 3 янв–апр май–авг сен–дек A 09-16 B 09-16 C 09-16 D 09-16 горячая точка: свежая запись вся в одной секции, соседние простаивают по хешу составного ключа: секцию выбирает hash(датчик), дата — внутри секции секция 1 секция 2 секция 3 hash 0–5 hash 6–a hash b–f A 09-16 B 09-16 C 09-16 D 09-16 поток растёкся по секциям — узлы нагружены ровно

Разбиение решает одно: в какую секцию уйдёт очередная запись. Когда ключ начинается с даты, весь свежий поток валится в одну секцию — это и есть горячая точка. Тот же поток при составном ключе, где первым идёт имя датчика, расходится по всем секциям; сортировка по дате при этом сохраняется внутри каждой секции.

Обязательно

Цель — равномерность, враг — горячая точка

Смысл секционирования один: раскидать данные и нагрузку по узлам поровну, чтобы десять узлов тянули в десять раз больше одного. Если разбиение вышло неравномерным — на часть секций пришлось непропорционально много данных или запросов, — его называют перекошенным (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 ключ → узел узлов стало 11 переехал 91% 1000 секций ключ → секция секция → узел переехало 9%

Добавили одиннадцатый узел: при hash mod N переезжает 91% ключей, а при тысяче фиксированных секций двигаются целые секции и только 9% ключей.

Кто даёт команду на перенос

И отдельный вопрос — автоматика против ручного управления. Полностью автоматическая ребалансировка удобна, но в паре с автоматическим обнаружением сбоев опасна: перегруженный узел отвечает медленно → система решает, что он «отказал» → затевает ребалансировку → добавляет нагрузки уже перегруженному узлу → каскадный сбой. Спасает тут не отказ от автоматики, а развод двух решений: обнаружить, что узел молчит, машина может сама, а вот команду «переносим данные» лучше отдать человеку («система предлагает — админ подтверждает», как в Couchbase). Медленнее, зато без каскада.

Глубже: маршрутизация запросов: где лежит «foo»?расширенное

Секции разъезжаются по узлам и переезжают при ребалансировке. Откуда клиент знает, к какому узлу идти за ключом foo? Это частный случай общей задачи «найти, кто за что отвечает» (service discovery), и решений три:

  1. Клиент стучится в любой узел. Если ключ там — узел отвечает сам; если нет — сам перенаправляет запрос нужному узлу и возвращает ответ. Так делают Cassandra и Riak через gossip-протокол: узлы сами между собой обмениваются картой кластера, и внешний координатор не нужен.
  2. Отдельное маршрутизирующее звено. Все запросы идут через него, и оно знает, какой узел за что отвечает (в MongoDB это mongos). Сами запросы оно не обрабатывает — это балансировщик, который просто в курсе, где какая секция.
  3. Клиент сам знает раскладку и подключается к нужному узлу напрямую, без посредников.

Во всех трёх случаях кто-то должен знать актуальную карту «секция → узел» и вовремя узнавать об изменениях. Классическое решение — отдельный сервис-координатор ZooKeeper: узлы регистрируются в нём, маршрутизаторы подписываются на изменения, и когда секция переезжает, ZooKeeper всех оповещает. Так делают HBase и SolrCloud. Kafka тоже так работала, но с версии 3.3 научилась хранить карту сама, а в 4.0 зависимость от ZooKeeper убрали совсем — это общий тренд: внешний координатор добавляет ещё одну систему, которую надо эксплуатировать. Альтернатива — тот самый gossip (Cassandra, Riak): узлы устроены сложнее, зато нет зависимости от внешнего координатора.

Коротко

  • Цель разбиения — раскидать по узлам и объём, и нагрузку. Перекос даёт горячую точку, и смысл шардинга теряется.
  • По диапазону — быстрые сканы по порядку ключа, но время в начале ключа гонит весь свежий поток в одну секцию. Лечит составной ключ.
  • По хешу — равномерность почти даром, платой уходит сортировка. Горячий ключ хеш не чинит: его размазывает приложение суффиксом, и платит за это чтение.
  • Локальный вторичный индекс — дешёвая запись и scatter/gather на чтении, вариант по умолчанию. Глобальный — наоборот и обновляется асинхронно.
  • hash mod N непригоден: переход с 10 узлов на 11 двигает 91% ключей, тысяча фиксированных секций — 9%.
  • Карту «секция → узел» держит внешний координатор (ZooKeeper) или сами узлы через gossip. Авторебалансировку вместе с автообнаружением сбоев не включают.
  • Секция, реплика и узел — три разные сущности: данные режут по ключу, каждую секцию держат в нескольких копиях, а машина несёт много секций в разных ролях.
  • Секций заводят заметно больше, чем машин (иначе балансировка грубая), но каждая стоит служебных расходов; по объёму держат десятки гигабайт.
  • Разрезанные данные ломают соединения между секциями, глобальную уникальность, транзакции на несколько секций и глубокую пагинацию — всё это становится «разослать и собрать».
  • Горячее чтение лечат кэшем перед хранилищем, а не только суффиксом к ключу; сменить ключ секционирования нельзя — только перелить данные заново.

Что почитать дальше