Consistent hashing concept page: hash ring with 4 nodes (A/B/C/D) with 256 vnodes each, demonstrating naive modulo disaster vs consistent hashing add/remove, vnode balance properties, and node failure rebalance via clockwise next-on-ring.
KEKey · @kuzminykh_igor_b3550a9b
0 stars
0 views
92d ago · last update
consistent-hashing.js·4 scenarios
Loading canvas…
Зачем
Шардирование «в лоб» — shard = hash(key) % N — работает ровно до первого изменения топологии. Стоит добавить или убрать одну ноду, как формула меняется для всех ключей сразу: модуль перестаёт «попадать» в те же бакеты. Для кластера 4 → 5 нод это означает: ~80% ключей переедут на другие физические серверы. На production-кластере в десятки TB это часы downtime, cache-miss storm на бэкенды и каскадные отказы.
Consistent hashing (Karger et al., 1997, оригинально для Akamai CDN) решает эту проблему элегантно: при добавлении или удалении одной ноды переезжает только ~1/N доли данных — остальные ключи продолжают мапиться туда, где и были.
«Положи ключи и ноды на одно кольцо хэшей. Ключ "идёт по часовой" до первой встречной ноды — это и есть его owner. Добавление ноды двигает только её ближайшего соседа против часовой. Чтобы баланс был справедливым — каждый физический сервер ставится на кольцо в 128–256 мест через virtual nodes.»
Три ключевых факта, без которых дальше нет смысла:
Кольцо хэшей — пространство [0, 2^32) или [0, 2^64), замкнутое в круг. И ключи, и ноды на одном пространстве.
Принцип "по часовой" — для ключа k ищем первую ноду с позицией ≥ hash(k). Если упёрлись в конец — заворачиваем в начало.
Vnodes решают балансировку — с одной точкой на ноду распределение получится перекошенным в 3–5 раз. С 256 точками — выравнивается до ±10%.
Что показывает диаграмма
В кольце четыре ноды-кеша (A, B, C, D), каждая владеет примерно четвертью пространства ключей. Клиент слева — это любой вызов cache.get(key) из приложения. Стрелки от клиента ко всем нодам показывают, что роутинг детерминирован hash-функцией от ключа, а не выбирается случайно или round-robin'ом.
В реальном production здесь жило бы не 4 точки, а 4 × 256 = 1024 vnode-позиций на кольце — но визуально это сливается в шум, поэтому в диаграмме оставлены только physical-ноды с подписью «~25%», обозначающей их статистическую долю.
Сценарии анимации не показывают сам ring (это статичная структура), а симулируют события поверх него: добавление ноды, отказ, перебалансировку.
Сценарии
Baseline: 4 ноды, hash(key) % 4. Стабильно, каждая владеет ~25%. Добавляем 5-ю ноду — формула становится hash(key) % 5, и остаток от деления меняется почти для всех ключей. Получаем remap ~80% данных: cache miss storm на бэкенды, миграция терабайтов по сети одновременно, downtime. Это причина, по которой naive modulo не применяют в production — только в студенческих туториалах.
То же добавление 5-й ноды, но через ring. Новая нода получает свои 256 vnode-позиций, разбросанных по кольцу. Каждая существующая нода теряет небольшой slice к новому соседу против часовой. Итог: переезжает ~1/5 = ~20% ключей, 80% остаются на месте — cache hit rate сохранён, миграция идёт в фоне, downtime ноль. Параллельность: 4 source-ноды шлют данные на 1 target-ноду одновременно маленькими порциями, а не одна гигантская труба.
Что будет, если у каждой ноды только одна точка на кольце (V=1)? Random hash placement даёт load skew до 5x: одна нода может случайно получить 50% ключей, другая — 10%. Hot node перегружена, остальные простаивают. Vnodes (V=256) фиксят это статистически — закон больших чисел: 1024 равномерно распределённые точки дают баланс ±10%. Cassandra использует 256 vnodes по умолчанию с версии 1.2. Trade-off: больше vnodes = больше metadata в gossip-трафике, потолок практический ~1024.
Нода падает — её слайс кольца уходит к next clockwise соседу. Если у node-C было 25%, то node-D временно владеет ~50% (свои + чужие). Это hot-spot risk: если node-D тоже не выдержит — каскад. Решение: vnodes делают этот удар равномерным (упавшая нода передаёт по маленькому slice каждому соседу, а не один большой блок одному), плюс заранее настроенный replication factor с rack/AZ awareness гарантирует, что реплики живут.
Trade-offs (ADR)
Контекст. Sharded cache/storage cluster, который должен переживать +1/-1 ноду без даунтайма и без перетряхивания всех данных. Альтернативы: naive modulo (отметаем сразу), статичные partitions (Redis Cluster slots), proxy-based routing (mcrouter), consistent hashing с vnodes.
Решение. Consistent hashing с vnodes — индустриальный стандарт для cluster size 10–1000 нод.
Последствия — что получаем.
При +1 ноде переезжает ~1/N данных вместо ~(N-1)/N. Для cluster из 100 нод это 1% данных вместо 99%.
Миграция параллельная: все existing nodes шлют свои маленькие куски одновременно.
Lookup O(log(V × N)) через binary search в sorted vnode array — ~1 µs.
Cache hit rate выживает при изменении топологии (ключи в основном остаются на старых местах).
Последствия — за что платим.
Metadata overhead. Кольцо из тысяч nodes × 256 vnodes = сотни тысяч записей. Cassandra-команда советует ≤256 vnodes как раз из-за gossip-трафика.
Сложность в коде клиента. Не три строки % N, а sorted ring, hash function pinned to version, binary search.
Hash function — критический выбор. Bad hash (плохая uniformity) = bad balance даже с vnodes. MurmurHash3, xxHash, CityHash — да; MD5/SHA1 — слишком медленно для каждого lookup.
Hot keys не решаются. Consistent hashing про равномерность по nodes, но один популярный ключ всё равно идёт на одну ноду. Нужны отдельные механики: replicate hot keys, request coalescing, separate hot tier.
Rack/AZ awareness нужно делать отдельно. Naive CH может разместить все реплики в одной AZ. Cassandra решает через snitch (NetworkTopologyStrategy), DynamoDB — через partition placement constraints.
Альтернативы — когда лучше другое:
Jump consistent hash (Lamping & Veach, Google, 2014) — O(log N) без хранения ring, ~5 ns. Минус: только sequential bucket numbering (0..N-1), нельзя удалить произвольный bucket. Используется в gRPC LB, Google internal.
Rendezvous (HRW) hashing (1996) — O(N) per lookup, без хранения ring, идеально для N < 100. Можно weighted. Используется в Apache Druid (segment placement).
Maglev hashing (Google, 2016) — lookup table размера M (prime, e.g. 65537), O(1) lookup ~10 ns, table read. При changes — rebuild с гарантией что большинство entries не меняется (connection persistence). Используется в Google Maglev L4 LB, Envoy.
Static slots (Redis Cluster: 16384 hash slots) — фиксированное число partitions, явный mapping slot → node. Проще для миграции slot-by-slot, но не масштабируется до тысяч нод.
Consistent hashing with bounded loads (Mirrokni et al., 2018) — расширение с лимитом на max load. Vimeo использует для cache.
Реальные системы
Akamai — изобретатели (Karger et al. 1997 paper писали под их CDN).
Memcached — client-side via libketama (Last.fm) — оригинальный killer use case.
Cassandra / ScyllaDB — token ring с 256 vnodes по умолчанию.
DynamoDB — partition assignment, thousands of partitions per node.
Riak — ring with 64–1024 partitions.
Couchbase — vBucket map (1024 vbuckets), та же идея.
Ceph — CRUSH algorithm (вариант CH с deterministic placement и failure domains).
Anti-patterns
Naive hash(key) % N в production. Любое изменение топологии = массовая миграция. Допустимо только в тестах и обучающих примерах.
Один vnode на ноду (V=1). Random placement → load skew до 5x. Используйте минимум 64, типично 128–256.
Слишком много vnodes (>1024 на ноду). Раздувает gossip и metadata — обратный эффект на performance кластера.
MD5 / SHA1 для key hash. Cryptographic hashes слишком медленные для каждого lookup. MurmurHash3, xxHash, CityHash — норма.
Consistent hashing как решение hot-key problem. CH про равномерность по нодам, single-key hotspot всё равно идёт на одну ноду. Нужны hot key replication, request coalescing, кеш перед кешем.
Игнорировать rack/AZ awareness. Без явных constraints все N реплик могут оказаться в одной AZ — теряем durability при failure домена.
Jump hash для кластера с произвольным удалением nodes из середины. Алгоритм работает только в sequential bucket numbering — нельзя выкинуть bucket #5 из 10.
Замена hash function на живом кластере. Backwards-incompatible — потребует full data migration. Hash function pin'ится на старте кластера.
Хранение ring config в БД с lookup на каждый запрос. Должно быть в gossip/config plane, in-memory у каждого клиента.
Когда НЕ использовать
Кластер ≤ 5 нод и редкие изменения. Накладные расходы на ring + vnodes + gossip не окупаются. Берите Rendezvous hashing или даже статичную таблицу маппинга.
Сильная схема данных требует co-location связанных ключей. Например, все данные одного user_id должны жить на одной ноде, плюс соседние user_id тоже. CH разбросает их случайно — используйте range-based sharding (как HBase, BigTable).
Запросы по range (SELECT * WHERE id BETWEEN 100 AND 200). CH хэширует key и теряет упорядоченность — придётся scatter-gather по всем нодам. Range partitioning эффективнее.
Hot keys доминируют workload. Если 5% ключей дают 80% запросов, балансировка по hash бесполезна. Сначала чините hot keys (replication, caching, sharding по дополнительному измерению), потом думайте про CH.
L4 load balancing где нужен O(1) lookup на packet. Берите Maglev hashing (variant CH с lookup table) — он ровно для этого. Классический ring с binary search даст лишний µs на packet, что заметно при миллионах PPS.
Single-node deployment. Очевидно, но: не нужно консистентного хэширования, когда нода одна.
Гарантированно фиксированное число шардов навсегда. Если бизнес-требование «ровно 16 шардов и точка» — статичная hash table проще, отлаживается за минуту.