2026 г.

Курс «Распределённые системы». Глава 8. Шардирование

Цели главы. Репликация (глава 5) копирует данные целиком; когда объём или поток записей не помещается на узел, данные шардируют — режут на части. Глава разбирает три вопроса, из которых состоит любая схема шардирования: по какому правилу резать (диапазон или хеш), как перекраивать при изменении числа узлов (согласованное хеширование против явных таблиц) и что делать с запросами, не знающими ключа (вторичные индексы). Плюс вечный спутник шардирования — перекос нагрузки. Лабораторная — собственная реализация согласованного хеширования с измерениями.

8.1. Два способа резать

Ключ шардирования — атрибут, по которому запись приписывается шарду; выбор правила — первое архитектурное решение:

  • По диапазону: шард владеет отрезком упорядоченных ключей [a, b). Достоинство — дешёвые диапазонные запросы: соседние ключи лежат рядом (сканы по времени, по префиксу). Так устроены HBase, TiKV, YDB. Опасность — горячие диапазоны: монотонный ключ (время, автоинкремент) направляет все вставки в последний шард — вся мощь кластера простаивает, пишет один узел.
  • По хешу: шард определяется hash(ключ). Нагрузка размазывается равномерно — и ровно поэтому диапазонный запрос рассыпается на все шарды. Компромиссная классика — составной ключ: хеш по первой части (равномерность между шардами), порядок по второй (сканы внутри шарда): (user_id, timestamp) даёт равномерность по пользователям и дешёвую историю одного пользователя.

8.2. Проблема mod N и согласованное хеширование

Наивное правило «шард = hash(k) mod N» рушится на первом же изменении N: при переходе N → N+1 новое значение hash(k) mod (N+1) совпадает со старым лишь у ~1/(N+1) ключей — почти все данные должны переехать. Для кластера в терабайты это означает, что добавление одного узла устраивает многочасовую внутреннюю миграцию с проседанием всего.

Согласованное хеширование (Karger et al., 1997) чинит это геометрией. Пространство хешей замыкается в кольцо; на кольцо хешируются и узлы, и ключи; ключ принадлежит первому узлу по часовой стрелке. При добавлении узла к N равноправным узлам ему должна достаться в среднем доля 1/(N+1) данных, то есть переезжает около K/(N+1) из K ключей — теоретический минимум. Удаление узла зеркально: его дуга достаётся соседу.

Сырое кольцо страдает неравномерностью: n случайных точек делят окружность на дуги, различающиеся в разы, и сосед погибшего узла принимает всё его наследство один. Лечение — виртуальные узлы: каждый физический узел присутствует на кольце сотней-другой точек (hash(узел#1), hash(узел#2), …). Дисперсия дуг падает (по закону больших чисел), наследство погибшего дробится по всему кластеру, а заодно появляется взвешивание: мощному узлу — больше виртуальных точек. Всё это вы измерите в лабораторной собственным кодом. Согласованное хеширование живёт в Dynamo-наследниках (Cassandra, Riak), клиентском шардировании memcached, балансировке CDN.

8.3. Явные таблицы шардов

Альтернатива алгоритмической геометрии — прозаичная таблица соответствия: пространство ключей режется на много мелких единиц (партиций/регионов/таблетов — тысячи на кластер), а карта «партиция → узел» хранится в отказоустойчивом сервисе метаданных, обычно реплицированном консенсусом. Это может быть отдельный etcd/ZooKeeper или собственная подсистема СУБД; единственный центральный процесс для такого подхода не обязателен. Ребаланс — фоновое переселение партиций с обновлением карты; горячая или распухшая партиция делится надвое. Подход тяжелее в реализации, но выразительнее: точечное управление размещением (геопривязка, изоляция арендаторов, вытеснение горячих партиций на свободные узлы) — поэтому его выбирают полновесные СУБД (HBase, TiDB, CockroachDB, YDB). Маршрутизация в обоих мирах: либо через слой прокси, либо клиент кэширует карту и умеет обрабатывать ответ «не мой ключ, обратись туда» при переездах.

8.4. Вторичные индексы

Шардирование по user_id мгновенно отвечает «данные пользователя X», но «все заказы со статусом = ожидает» не знает ключа. Два устройства вторичного индекса:

  • Локальный (document-partitioned): каждый шард индексирует своё. Запись дёшева (индекс обновляется там же, где данные, атомарно), поиск — scatter-gather: запрос ко всем шардам, слияние ответов; задержка — по худшему шарду, стоимость растёт с числом шардов.
  • Глобальный (term-partitioned): индекс — сам по себе шардированная структура, разрезанная по индексируемому значению. Поиск — точечный (один шард индекса), зато запись трогает чужие шарды — асинхронно (индекс отстаёт: снова итоговая согласованность, глава 4) либо транзакционно (глава 9, с её ценами).

Инженерное правило: локальные индексы + сдержанность в scatter-gather для оперативных запросов; тяжёлая аналитика по неключевым атрибутам — в отдельную аналитическую систему, а не по шардам продакшена.

8.5. Перекос и горячие ключи

Равномерность хеша защищает от систематического перекоса, но не от горячего ключа: хеш безупречно направит миллион чтений страницы знаменитости в один шард. Приёмы по нарастанию сложности: кэш перед шардом (горячее — по определению кэшируемое); репликация чтений горячего ключа по нескольким узлам; для записей — подсаливание: ключ размножается на key#0…key#R по случайному суффиксу, записи размазываются, чтения собирают все R (осмысленно для агрегатов — счётчиков, лайков); в системах с таблицами партиций — автоматическое деление горячей партиции. Общая мораль: перекос — свойство нагрузки, а не схемы; его мониторят per-shard метриками и лечат адресно.

Итоги главы

  • Диапазон дружит со сканами и опасен монотонными ключами; хеш равномерен и враждебен диапазонам; составной ключ — рабочий компромисс.
  • mod N перемещает почти всё при смене N; при добавлении узла согласованное хеширование перемещает в среднем K/(N+1) ключей, а виртуальные узлы дают равномерность, дробление наследства и веса.
  • Явные таблицы партиций в отказоустойчивом сервисе метаданных — выразительная альтернатива: точечный ребаланс и деление горячих партиций; это частый выбор полновесных СУБД.
  • Вторичные индексы: локальные (дешёвая запись, scatter-gather чтение) против глобальных (точечное чтение, асинхронная запись).
  • Горячие ключи хешем не лечатся: кэш, реплики чтения, подсаливание, деление партиций — и обязательный пошардовый мониторинг.

Упражнения

  1. События логируются с ключом «timestamp». Объясните, что произойдёт при диапазонном шардировании, и предложите два исправления с сохранением дешёвых запросов «за последний час по сервису X» (подсказка: составной ключ; какой атрибут — в хеш-часть?).
  2. Посчитайте для mod N: какая доля ключей сохраняет шард при N=10 → 11? А при удвоении N=10 → 20? Сравните с согласованным хешированием в обоих случаях.
  3. Кольцо без виртуальных узлов, три узла. Покажите (можно на конкретных значениях хешей), что один узел может владеть >70% кольца. Оцените, сколько виртуальных точек на узел нужно, чтобы дисперсия долей стала приемлемой (эксперимент — в лабораторной; здесь — качественное рассуждение через закон больших чисел).
  4. Магазин шардирован по customer_id; аналитик просит «топ-10 товаров по суммарным продажам за вчера». Разберите исполнение запроса при локальных индексах и покажите, почему слияние только локальных топ-10 в общем случае неверно. Постройте контрпример из двух шардов и предложите точный и приближённый способы расчёта.
  5. Ключ «счётчик просмотров матча» принимает 50 000 записей/с — один шард столько не держит. Спроектируйте подсаливание: выбор R, путь записи, путь чтения, во что превратилась атомарность инкремента (вспомните G-Counter из главы 5 — что общего?).
  6. При переезде партиции клиенты со старой картой шардов продолжают писать в прежний узел. Спроектируйте протокол переезда без потерь и двойных записей: этапы, состояние «партиция переезжает», ответы старого владельца. Какой механизм главы 5 здесь обязателен (подсказка: эпоха владения)?

Ответы и указания. 1: один растущий timestamp направляет свежие записи в последний диапазон. Вариант ключа (hash(service_id), timestamp) распределяет разные сервисы и сохраняет временной скан внутри одного сервиса. Если один сервис сам горячий, вводят фиксированное число полос: (hash(service_id, stripe), timestamp), где stripe выбирается по идентификатору события; запрос за час объединяет только S известных полос. 2: при N→N+1 ключ сохраняет шард, если hash mod 10 = hash mod 11 — примерно 1/11 ≈ 9%; при удвоении hash mod 10 = hash mod 20 ровно для половины значений — сохраняется 50%; согласованное хеширование перемещает около 1/11 и 1/2 соответственно. 3: точки 0,05; 0,10; 0,80 оставляют одному узлу дугу 0,70, а сколь угодно малое смещение даёт больше 70%. Виртуальные точки усредняют независимые дуги, поэтому относительный разброс убывает примерно как 1/√V; универсального V нет, его выбирают экспериментом под требуемый перекос и вес узлов. 4: пусть товар X имеет по 9 продаж на каждом из двух шардов и занимает на каждом 11-е место, а десять локальных лидеров каждого шарда имеют по 10 продаж и отсутствуют на другом. X с суммой 18 — глобальный лидер, но ни в один локальный топ-10 не попал. Точный вариант суммирует все частичные агрегаты по товару или заранее шардирует агрегат по product_id; приближённый использует алгоритм частых элементов со sketch и заданной погрешностью. 5: R выбирается так, чтобы 50000/R было ниже лимита шарда; чтение суммирует R ячеек, а единого атомарного значения в момент записи больше нет. 6: этапы — копирование с журналом изменений, переключение карты, дренаж; старый владелец после передачи отвечает «не мой, эпоха такая-то». Эпоха владения ограждает запоздавшие записи клиентов со старой картой.

Лабораторная работа 3. Согласованное хеширование своими руками

Цель — реализовать кольцо с виртуальными узлами (100–200 строк Python) и получить численно все утверждения раздела 8.2. Никакой инфраструктуры: стандартная библиотека (hashlib, bisect).

Задание 1. Кольцо. Реализуйте класс: add_node(name), remove_node(name), get_node(key). Точки кольца — sha1 от строк вида "имя_узла#номер_виртуальной_точки"; отсортированный список точек + bisect для поиска первого узла по часовой стрелке.

Задание 2. Равномерность. Для 5 узлов и 100 000 случайных ключей постройте распределение «узел → доля ключей» при 1, 10, 100, 500 виртуальных точках на узел. В отчёт: таблица (или график) максимального и минимального узла; при скольких точках разброс входит в ±10%?

Задание 3. Стоимость изменения кластера. Зафиксируйте отображение 100 000 ключей для 10 узлов; добавьте одиннадцатый; посчитайте долю ключей, сменивших узел. Сравните с теоретической 1/11 и — реализовав в три строки — с mod N. Повторите для удаления узла. В отчёт: доли переездов для трёх схем (кольцо, mod N, теория).

Задание 4. Наследство погибшего. Удалите узел и посчитайте, между сколькими выжившими распределились его ключи — при 1 виртуальной точке и при 200. В отчёт: гистограмма распределения наследства; вывод о роли виртуальных узлов при отказах.

Задание 5*. Веса и горячий ключ. (а) Дайте одному узлу вдвое больше виртуальных точек и подтвердите измерением его двойную долю. (б) Смоделируйте горячий ключ (30% запросов в один ключ), покажите бесполезность любых схем хеширования против него и реализуйте подсаливание с R=8: измерьте выравнивание нагрузки и цену чтения.

Литература к главе

  1. D. Karger et al., "Consistent Hashing and Random Trees," STOC, 1997.
  2. G. DeCandia et al., "Dynamo," SOSP, 2007 — разделы о партиционировании и виртуальных узлах.
  3. M. Kleppmann, "Designing Data-Intensive Applications," O'Reilly, 2017 — гл. 6 (партиционирование) — основной текст к главе.
  4. Документация Cassandra: Dynamo-архитектура и токены. cassandra.apache.org/doc

Предыдущая глава || Содержание курса || Следующая глава

404 Not Found

404 Not Found


nginx/1.24.0 (Ubuntu)

Связь с редакцией