Цели главы. Репликация (глава 5) копирует данные целиком; когда объём или поток записей не помещается на узел, данные шардируют — режут на части. Глава разбирает три вопроса, из которых состоит любая схема шардирования: по какому правилу резать (диапазон или хеш), как перекраивать при изменении числа узлов (согласованное хеширование против явных таблиц) и что делать с запросами, не знающими ключа (вторичные индексы). Плюс вечный спутник шардирования — перекос нагрузки. Лабораторная — собственная реализация согласованного хеширования с измерениями.
Ключ шардирования — атрибут, по которому запись приписывается шарду; выбор правила — первое архитектурное решение:
Наивное правило «шард = 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.
Альтернатива алгоритмической геометрии — прозаичная таблица соответствия: пространство ключей режется на много мелких единиц (партиций/регионов/таблетов — тысячи на кластер), а карта «партиция → узел» хранится в отказоустойчивом сервисе метаданных, обычно реплицированном консенсусом. Это может быть отдельный etcd/ZooKeeper или собственная подсистема СУБД; единственный центральный процесс для такого подхода не обязателен. Ребаланс — фоновое переселение партиций с обновлением карты; горячая или распухшая партиция делится надвое. Подход тяжелее в реализации, но выразительнее: точечное управление размещением (геопривязка, изоляция арендаторов, вытеснение горячих партиций на свободные узлы) — поэтому его выбирают полновесные СУБД (HBase, TiDB, CockroachDB, YDB). Маршрутизация в обоих мирах: либо через слой прокси, либо клиент кэширует карту и умеет обрабатывать ответ «не мой ключ, обратись туда» при переездах.
Шардирование по user_id мгновенно отвечает «данные пользователя X», но «все заказы со статусом = ожидает» не знает ключа. Два устройства вторичного индекса:
Инженерное правило: локальные индексы + сдержанность в scatter-gather для оперативных запросов; тяжёлая аналитика по неключевым атрибутам — в отдельную аналитическую систему, а не по шардам продакшена.
Равномерность хеша защищает от систематического перекоса, но не от горячего ключа: хеш безупречно направит миллион чтений страницы знаменитости в один шард. Приёмы по нарастанию сложности: кэш перед шардом (горячее — по определению кэшируемое); репликация чтений горячего ключа по нескольким узлам; для записей — подсаливание: ключ размножается на key#0…key#R по случайному суффиксу, записи размазываются, чтения собирают все R (осмысленно для агрегатов — счётчиков, лайков); в системах с таблицами партиций — автоматическое деление горячей партиции. Общая мораль: перекос — свойство нагрузки, а не схемы; его мониторят per-shard метриками и лечат адресно.
Ответы и указания. 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: этапы — копирование с журналом изменений, переключение карты, дренаж; старый владелец после передачи отвечает «не мой, эпоха такая-то». Эпоха владения ограждает запоздавшие записи клиентов со старой картой.
Цель — реализовать кольцо с виртуальными узлами (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: измерьте выравнивание нагрузки и цену чтения.