Партиционирование и шардирование: стратегии, ребалансировка, горячие ключи
Репликация отвечает на вопрос «сколько копий у одних и тех же данных». Партиционирование отвечает на принципиально другой: какие данные вообще живут на этом узле. Реплики хранят одно и то же, секции — разное. Это ортогональные оси, и путают их чаще, чем любую другую пару понятий в предмете. В настоящей системе есть обе: данные режутся на секции, каждая секция реплицируется, и падение одной машины означает потерю по одной реплике у нескольких разных секций.
Партиционирование заводят по четырём причинам, и они, в отличие от причин для репликации, между собой не конфликтуют: объём (40 ТБ не влезут на машину с 4 ТБ SSD), пропускная способность записи (один лидер упирается в одну очередь fsync, десять секций — в десять), размер рабочего набора (индекс на 200 ГБ не влезет в кеш на 64 ГБ, а десять по 20 ГБ — влезут по частям) и радиус поражения (сбой затрагивает 1/n пользователей). Звучит как бесплатный обед — и поэтому партиционирование остаётся самым недооценённым решением в проекте. Цена платится не при внедрении, а через полгода, когда выясняется, что 60 % трафика даёт один арендатор, что «покажи последние 50 заказов» превратилось в опрос всех 128 секций, а обещанный exactly-once сломался ровно потому, что таблица дедупликации оказалась в другой секции, чем эффект.
1. Что мы режем и как это называется
Партиционирование — разбиение данных на непересекающиеся подмножества так, что каждый элемент принадлежит ровно одной секции, а функция partition(key) → partition_id детерминирована и известна всем участникам. Шардирование — то же самое, когда секции разнесены по разным машинам. Дальше «секция» — логическая единица разбиения, «шард» — секция вместе со своей группой репликации.
решение партиционирования:
КАКАЯ секция"] P --> G["Секция P7 = группа репликации"] G --> L["Лидер: узел B"] G --> R1["Реплика: узел D"] G --> R2["Реплика: узел F"] L --> D["Решение репликации:
СКОЛЬКО копий и когда отвечать OK"] R1 --> D R2 --> D P -.->|"меняется при ребалансировке"| G G -.->|"меняется при отказе узла"| L
Две пунктирные стрелки — два разных класса инцидентов, которые постоянно смешивают в постмортемах. «Секция переехала на другой узел» и «у секции сменился лидер» — разные события, с разными последствиями и разными логами.
Терминологическая ловушка. В PostgreSQL и ClickHouse «partitioning» означает разбиение таблицы внутри одного узла: локальные куски ради быстрого DROP PARTITION и pruning по предикату. Это не распределённая техника — PARTITION BY toYYYYMM(ts) в ClickHouse не имеет отношения к sharding_key в Distributed-таблице. В Kafka, Cassandra и DynamoDB «partition» — именно распределённая единица.
| Система | Единица | Что задаёт разбиение | Число секций меняется |
|---|---|---|---|
| Apache Kafka | partition топика | murmur2(key) % numPartitions |
только вверх, вручную |
| Cassandra / ScyllaDB | token range (vnode) | Murmur3(partition key) в [-2^63, 2^63) |
автоматически при вводе узла |
| Amazon DynamoDB | partition | внутренний хеш partition key | автоматически и скрыто |
| MongoDB | chunk / range | shard key: ranged или hashed | балансировщиком, автоматически |
| Elasticsearch | primary shard | murmur3(_routing) % number_of_shards |
нет, только _split / _shrink |
| Redis Cluster | hash slot (16384) | CRC16(key) mod 16384 |
слоты переносятся, число фиксировано |
| Google Spanner | split | диапазон первичного ключа | автосплит по размеру и нагрузке |
| TiKV / TiDB | region (~96 МиБ) | диапазон ключа, планировщик PD | автосплит и слияние |
Последний столбец делит мир пополам: системы с фиксированным числом секций (Kafka, Elasticsearch) и системы с динамическим делением и слиянием (HBase, Spanner, TiKV, Vitess). Это два разных инженерных подхода с разными отказами.
2. Ключ партиционирования — решение о том, что можно сделать одним запросом
Ключ задаёт сразу четыре вещи, и все четыре потом невозможно поменять дёшево:
- Локальность. Данные с одинаковым ключом лежат вместе; всё, что делается одним обращением к одному узлу, определяется этим.
- Границу транзакции. Однопартиционная транзакция — обычная локальная транзакция БД. Межпартиционная — двухфазный коммит или сага. Разница в стоимости — порядок величины.
- Границу порядка. В Kafka порядок гарантирован только внутри партиции. Если
order.createdиorder.cancelledпопали в разные партиции, потребитель законно увидит отмену раньше создания. - Единицу перекоса. Один ключ — всегда один узел. Никакая стратегия не размажет нагрузку на один ключ по кластеру.
Отсюда правило, звучащее как заимствование из DDD и им и являющееся: ключ партиционирования должен совпадать с границей агрегата (см. тактические строительные блоки). Если инвариант звучит как «баланс счёта не может стать отрицательным», то account_id обязан быть ключом партиционирования — иначе проверка инварианта требует координации.
В Cassandra это выражено в синтаксисе: PRIMARY KEY ((tenant_id, day), event_ts, event_id) — первая скобка определяет узел, остальное задаёт порядок внутри секции. DynamoDB называет то же самое partition key + sort key: хеш выбирает машину, диапазон работает внутри неё.
3. Стратегии разбиения
3.1. По диапазонам
Ключевое пространство упорядочено и режется на непрерывные интервалы, границы хранятся в карте секций. Ради чего это выбирают: диапазонный запрос читается одной-двумя секциями. Для аналитики и таймсерий аргумент решающий.
Сценарий отказа: горячий хвост. Ключ — метка времени или монотонный идентификатор (автоинкремент, Snowflake ID, UUIDv7). Все новые записи попадают в последний диапазон. В HBase:
WARN regionserver.HRegion: Region events,\x39\x2E,1741660800123.a3f1c8...
has too many store files; delaying flush up to 90000ms
WARN regionserver.MemStoreFlusher: Memstore size is 3.9 G, above high water mark,
blocking updates for region events,\x39\x2E,...
ERROR ipc.CallRunner: RegionTooBusyException: Above memstore limit,
regionName=events,\x39\x2E,1741660800123.a3f1c8...
Три RegionServer’а из четырёх при этом стоят на 5 % CPU, а дашборд средней загрузки показывает уютные 30 %. Это фирменная черта перекоса: средние метрики его скрывают, смотреть надо на максимум по секциям. Лечение — вынести перед временем что-то с высокой кардинальностью: (tenant_id, ts) вместо (ts), либо префикс-соль salt = hash(row_key) % 16; второе ломает глобальные диапазонные сканы, теперь их 16.
3.2. По хешу и композитный ключ
partition = hash(key) % P даёт равномерность независимо от вида ключей, но уничтожает порядок: диапазонный запрос превращается в опрос всех секций. Компромисс — композитный ключ: хеш по первой части, диапазон по второй.
CREATE TABLE shop.events (
tenant_id uuid, day date, event_ts timestamp, event_id timeuuid, payload text,
PRIMARY KEY ((tenant_id, day), event_ts, event_id)
) WITH CLUSTERING ORDER BY (event_ts DESC);
-- Один узел, один последовательный скан, уже отсортировано: это дёшево.
SELECT * FROM shop.events
WHERE tenant_id = ? AND day = '2026-03-11' AND event_ts > '2026-03-11 08:00:00' LIMIT 100;
-- А это опрос ВСЕХ секций кластера, и Cassandra потребует ALLOW FILTERING.
SELECT * FROM shop.events WHERE event_ts > '2026-03-11 08:00:00';
day в partition key стоит там не для красоты: без него секция арендатора растёт неограниченно и через год превращается в партицию на 600 МБ, которую нельзя ни быстро прочитать, ни перенести. Приём называется bucketing и обязателен для любых таймсерий в Cassandra.
3.3. Почему hash(key) mod N по числу узлов — почти всегда ошибка
Ключ остаётся на месте при переходе от N к N+1 узлам, только если hash(k) mod N == hash(k) mod (N+1) — а это верно примерно для доли 1/(N+1) ключей. Значит, переезжает N/(N+1) данных: при 4 узлах 80 %, при 10 узлах 91 %. Панель ② на схеме выше показывает это на двенадцати ключах.
Сценарий отказа: шторм промахов кеша. Слой memcached из 8 узлов, mod 8. Один узел падает, клиентская библиотека переключается на mod 7, и 87 % ключей мгновенно становятся промахами:
02:41:12 cache.hit_ratio=0.968 db.qps=1240 db.p99=8ms
02:41:19 cache.hit_ratio=0.121 db.qps=41800 db.p99=2140ms
02:41:24 postgres FATAL: sorry, too many clients already
02:41:26 upstream connect error or disconnect/reset before headers. reset reason: overflow
Упал один узел — лёг весь кеш и база под ним: это каскад, а не деградация. Именно на этом сценарии построена мотивация статьи Karger et al. 1997 года.
3.4. Фиксированное число секций
Приём, снимающий проблему, — отвязать число секций от числа узлов. Секций делают заведомо больше, чем машин (Kafka — сотни партиций, Redis Cluster — ровно 16384 слота, Elasticsearch — number_of_shards), а на узлы раскладывают целые секции. Ключ всегда отображается в один и тот же номер секции; при вводе узла переезжают секции целиком, а не отдельные ключи (панель ③).
Сценарий отказа: увеличили число партиций в Kafka. Потребитель не справляется, партиций 12, поднимаем до 24. Kafka разрешает операцию, и она мгновенно меняет murmur2(key) % P: ключ order-8842 раньше шёл в партицию 4, теперь — в 16. Старые события заказа остались в четвёртой, новые пойдут в шестнадцатую, и два потребителя читают их параллельно. В логах это выглядит как невозможное:
02:52:04Z INFO applying event order.cancelled order_id=8842 version=7
02:52:04Z ERROR aggregate not found: order 8842 (expected version 6, store empty)
02:52:11Z INFO applying event order.created order_id=8842 version=1
«Отмена пришла раньше создания» — сигнатура именно этой ошибки. Уменьшить число партиций Kafka не даёт вовсе: это сломало бы и смещения потребителей, и оставшиеся данные. Отсюда правило: партиций закладывают с запасом ×3–×5, потому что цена лишней партиции (дескрипторы, метаданные, время перевыборов) существенно ниже цены изменения их числа. В Elasticsearch изменение number_of_shards вообще запрещено, а _split требует read-only индекса и работает только с кратными множителями при заранее заданном number_of_routing_shards.
3.5. Справочник (directory / lookup-based)
Карта «диапазон ключей → шард» хранится явно: mongos + config servers в MongoDB, vindex в Vitess, Placement Driver в TiKV, Slicer в Google. Отображение произвольное, поэтому можно переносить что угодно куда угодно и балансировать по фактической нагрузке, а не по числу ключей. Цена — компонент, который обязан быть одновременно согласованным и высокодоступным, то есть консенсус со всеми его ограничениями, плюс кеш карты у клиентов и связанное с ним устаревание.
Квадрант 2 пуст не случайно (композитный ключ лёг бы примерно в его нижнюю границу): дешёвая ребалансировка и глобальный порядок ключей — противоречивые требования. Композитный ключ — единственный честный компромисс, порядок в нём есть только внутри секции.
4. Консистентное хеширование и виртуальные узлы
Идея Karger et al. (STOC 1997, PDF): отобразить и ключи, и узлы в одно кольцо хеш-значений; ключ принадлежит первому узлу по часовой стрелке. Когда узел выбывает, его дуга достаётся соседу, а все прочие дуги не двигаются вообще.
В наивном виде схема плоха двумя вещами: n случайных точек делят кольцо неравномерно (отношение максимальной дуги к средней растёт как O(log n)), а при отказе весь объём узла валится на одного соседа, который падает следом. Оба дефекта лечит одна конструкция — виртуальные узлы (токены): каждый физический узел ставит на кольцо v точек. Относительный разброс нагрузки убывает как 1/√v (при v = 256 — около 6 %), дуги выбывшего узла расходятся между всеми выжившими, а восстановление идёт параллельно со всех уцелевших машин. Так устроены Dynamo (SOSP 2007) и Cassandra (num_tokens — по умолчанию 16 начиная с 4.0, раньше 256; уменьшили потому, что много дуг повышает шанс, что случайная тройка отказавших узлов накроет полный набор реплик какой-нибудь из них).
route(key):
h ← hash(key)
i ← первый индекс в отсортированном массиве токенов, где token[i] > h
если i вышел за границу — i ← 0 # замыкание кольца
вернуть owner[token[i]]
import bisect
import hashlib
class HashRing:
"""Кольцо консистентного хеширования с виртуальными узлами."""
def __init__(self, vnodes: int = 256) -> None:
self.vnodes = vnodes
self._tokens: list[int] = [] # отсортированные позиции на кольце
self._owner: dict[int, str] = {} # позиция -> физический узел
@staticmethod
def _hash(data: str) -> int:
# blake2b быстрее sha1 и даёт равномерное распределение; берём 64 бита
return int.from_bytes(hashlib.blake2b(data.encode(), digest_size=8).digest(), "big")
def add(self, node: str) -> None:
for i in range(self.vnodes):
token = self._hash(f"{node}#{i}")
if token in self._owner: # маловероятная коллизия — пропускаем
continue
bisect.insort(self._tokens, token)
self._owner[token] = node
def remove(self, node: str) -> None:
self._tokens = [t for t in self._tokens if self._owner[t] != node]
self._owner = {t: n for t, n in self._owner.items() if n != node}
def route(self, key: str) -> str:
i = bisect.bisect_right(self._tokens, self._hash(key))
return self._owner[self._tokens[i % len(self._tokens)]] # % — замыкание кольца
Сложность. route — O(log(n·v)) по времени, O(1) дополнительной памяти. add — O(v·log(n·v)) сравнений плюс O(n·v) на сдвиги массива (в проде вместо insort берут пересборку отсортированного массива или B-дерево). Память кольца — O(n·v): 100 узлов по 256 токенов это 25 600 записей, десятки килобайт, помещается в кеш клиента. Ключевое свойство: доля переезжающих ключей при вводе узла — 1/(n+1), а не n/(n+1), как у mod N. Проверять это стоит эмпирически — на самописных «кольцах», оказывавшихся mod N, обжигались многие: возьмите 200 000 ключей, запомните маршруты, добавьте пятый узел к четырём и посчитайте долю изменившихся. Должно получиться ≈ 20 % переезда и 5–8 % разброса размеров; 40 % переезда означает, что у вас не кольцо, а разброс в 30 % — что токенов слишком мало.
Продолжив обход кольца дальше по часовой стрелке с пропуском повторов физических узлов, из того же массива токенов получают preference list — список реплик ключа в Dynamo и Cassandra (в проде — с учётом стоек и датацентров). Это и есть стык партиционирования с репликацией: одна структура данных отвечает на оба вопроса.
4.1. Rendezvous hashing и jump hash
Альтернатива Thaler и Ravishankar (1996, техотчёт): никакого кольца, для каждой пары «ключ, узел» считается вес, побеждает максимум.
import hashlib
import math
def rendezvous(key: str, nodes: list[str], weights: dict[str, float] | None = None) -> str:
"""HRW: узел с максимальным взвешенным хешем. Разделяемого состояния — ноль."""
weights, best, best_score = weights or {}, "", float("-inf")
for node in nodes:
digest = hashlib.blake2b(f"{node}|{key}".encode(), digest_size=8).digest()
x = (int.from_bytes(digest, "big") + 1) / 2 ** 64 # равномерно в (0, 1]
score = -weights.get(node, 1.0) / math.log(x) # взвешивание HRW
if score > best_score:
best, best_score = node, score
return best
O(n) на ключ против O(log(n·v)) у кольца: при 20 узлах несущественно, при 2000 уже больно (есть иерархический вариант за O(log n)). Взамен — идеальное распределение без токенов, честные веса и главное свойство: упорядоченный список кандидатов получается бесплатно, достаточно отсортировать по score. Для выбора реплик и фолбэков это удобнее кольца; HRW применяют в Ceph (внутри CRUSH) и в клиентских балансировщиках.
Третий вариант — jump consistent hash (Lamping, Veach, Google, 2014, arXiv:1406.2294): пять строк цикла вида while j < buckets, где key прокручивается 64-битным LCG, O(ln n) времени, O(1) памяти и идеальный баланс без единой таблицы. Ограничение жёсткое: корзины нумеруются 0..buckets-1, удалить произвольную из середины нельзя — только последнюю. Поэтому jump hash хорош для разбиения на логические корзины фиксированного набора, которые затем отображаются на физические узлы справочником, и непригоден как замена кольцу при произвольных отказах. Для балансировщиков нагрузки есть ещё Maglev (Google, NSDI 2016): таблица поиска фиксированного размера даёт O(1) маршрутизацию и почти идеальный баланс ценой чуть большего числа переездов.
5. Ребалансировка: как переносить секцию, не потеряв записи
Ребалансировка — перемещение секции с узла на узел. Звучит как копирование файла; на деле это самая опасная операция в жизни кластера, потому что во время неё два узла одновременно считают себя владельцами одних данных.
Критична одна-единственная стрелка — ПередачаВладения, всё остальное оптимизация. Смысл заморозки не в том, чтобы избежать окна недоступности записи (избежать нельзя), а в том, чтобы сделать его измеримо малым — десятки миллисекунд.
Сценарий отказа: переезд без фенсинга. Уберите epoch — получите классику. Источник теряет связь с координатором на 30 секунд, тот считает его мёртвым и отдаёт P7 приёмнику. Источник возвращается, у него в памяти по-прежнему «я владею P7», и клиенты с устаревшей картой продолжают писать в него. Обе копии принимают записи:
node-a 02:14:31 INFO shard P7: serving writes (owner since 02:03:11)
node-d 02:14:31 INFO shard P7: serving writes (owner since 02:14:29)
node-a 02:19:02 INFO shard P7: relinquishing ownership, 41822 writes since handover
Последняя строка — приговор: 41 822 записи принял узел, который уже не был владельцем, и они не попали никуда. В биллинге это «прошедшие» и исчезнувшие списания. Спасает ровно одно — монотонно растущий номер владения (fencing token, он же epoch), который хранится в системе с консенсусом и проверяется на каждой записи; механика подробно разобрана в статье про координацию.
5.1. Что и куда переносить
- Фиксированное число секций. Их много, они целые, планировщик просто раздаёт их узлам. Просто и предсказуемо, но не спасает от секции, выросшей в 50 раз больше соседей.
- Динамические split / merge. Секция делится при превышении порога: HBase — по размеру store file, MongoDB — по объёму данных, TiKV — по
region-split-size(≈ 96 МиБ), Spanner — по размеру и по нагрузке. Адаптивно, но старт с пустой БД означает одну секцию и один узел под всем трафиком первого дня — отсюда практика pre-splitting. - Пропорционально числу узлов. Число секций фиксировано на узел (
num_tokensв Cassandra): ввод машины автоматически добавляет токены и забирает по кусочку у всех остальных.
Сценарий отказа: автоматическая ребалансировка встречает детектор отказов. Самый коварный сюжет темы, прямо связанный с моделями отказов. Узел тормозит из-за долгого GC → детектор объявляет его мёртвым → планировщик переносит его секции → копирование съедает сетевую полосу → heartbeat’ы соседей начинают опаздывать → мёртвыми объявлены ещё двое → переносятся их секции. Кластер уходит в самоподдерживающийся шторм миграций и перестаёт обслуживать трафик, не потеряв при этом ни одной машины:
02:03:11 gossip: /10.0.4.12 is now DOWN (no heartbeat for 10.4s)
02:03:12 streaming 214 ranges to /10.0.4.19, /10.0.4.21
02:05:40 gossip: /10.0.4.19 is now DOWN (no heartbeat for 11.1s)
02:05:41 streaming 331 ranges to /10.0.4.21, /10.0.4.24
02:07:03 MutationStage pending tasks: 48211, dropped MUTATION in last 5000ms: 20144
Отсюда три правила. Первое: ребалансировка обязана быть с тормозом — stream_throughput_outbound_megabits_per_sec в Cassandra, store-limit в TiKV, окно балансировщика в MongoDB, --throttle у kafka-reassign-partitions.sh; дефолты часто слишком щедрые. Второе: автоматический перенос при отказе опасен, при плановом вводе узла — нормален; Kleppmann формулирует то же самое, а практический компромисс — система готовит план, человек нажимает кнопку. Третье: данные не переносят при кратковременной недоступности — Cassandra ждёт max_hint_window (по умолчанию 3 часа), прежде чем считать узел выбывшим окончательно, потому что реплика, отсутствовавшая 40 секунд, догонится по логу за секунды, а перенос терабайтов не догонится никогда.
6. Маршрутизация: кто знает, где лежит ключ
Карта секций существует; вопрос — у кого она есть. Три архитектуры: толстый клиент (драйвер держит карту и идёт сразу на нужный узел — Cassandra TokenAwarePolicy, Kafka, Redis Cluster: минимум прыжков, максимум версий клиентов в проде), прокси (mongos, VTGate, ProxySQL: клиенты простые, но добавляется прыжок и компонент в критическом пути) и маршрутизация любым узлом (клиент идёт куда попало, узел перенаправляет или проксирует — координатор запроса в Cassandra, Elasticsearch).
Во всех трёх случаях клиентский кеш карты неизбежно устаревает — это не баг, а режим работы: карта меняется, а синхронно рассылать её всем значит нарушить доступность, поэтому протоколы включают явный механизм коррекции:
127.0.0.1:6379> GET user:42:cart
(error) MOVED 8342 10.0.3.17:6379 # слот уехал насовсем — обнови карту
(error) ASK 8342 10.0.3.17:6379 # слот мигрирует — сходи туда, карту не трогай
127.0.0.1:6379> MGET user:42:cart user:43:cart
(error) CROSSSLOT Keys in request don't hash to the same slot
# Kafka: лидер партиции сменился, продюсер обязан перечитать метаданные
WARN o.a.k.c.NetworkClient - [Producer clientId=orders-api] Received invalid metadata error
in produce request on partition orders-7 due to NotLeaderOrFollowerException
Разница между MOVED и ASK тонкая и важная. Клиент, обрабатывающий ASK как MOVED, во время миграции слота начинает метаться между узлами. Сценарий отказа: бесконечный редирект. Клиент кеширует карту, но не обновляет её после MOVED — каждый запрос идёт узел A → MOVED → узел B → MOVED → узел A. Латентность утраивается, RPS падает, а метрика ошибок остаётся нулевой, потому что формально ответы приходят. Ищется по отношению числа redirect-ответов к числу запросов, и эту метрику стоит завести заранее — см. наблюдаемость распределённых систем. Хеш-теги (user:{42}:cart и user:{42}:orders попадают в один слот, потому что хешируется только содержимое фигурных скобок) — штатный способ заставить связанные ключи жить вместе.
7. Горячие ключи и перекос
Все стратегии выше распределяют ключи. Ни одна не распределяет нагрузку на один ключ: если 40 % запросов идут в product:iphone-17, то хеширование, кольцо, токены и справочник одинаково честно отправят все 40 % на один узел.
Как это выглядит. Главное свойство перекоса — невидимость в средних. Смотреть надо на распределение по секциям:
$ kafka-consumer-groups.sh --describe --group billing
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
billing payments 0 812004431 812004433 2
billing payments 3 774119002 812440118 38321116
billing payments 4 811998210 811998210 0
Одна партиция отстаёт на 38 миллионов сообщений, остальные — на единицы; средний лаг по группе выглядит в дашборде терпимо. Причина почти всегда одна: ключом взят tenant_id, а один арендатор даёт больше трафика, чем все остальные вместе. Cassandra предупреждает о больших секциях явно, DynamoDB — отказом, MongoDB — неперемещаемым чанком:
WARN BigTableWriter.java:246 - Writing large partition shop/events:tenant-000042 (612.804MiB)
WARN ReadCommand.java:569 - Read 5000 live rows and 148213 tombstone cells for query
SELECT * FROM shop.events WHERE tenant_id = 'tenant-000042' LIMIT 5000
ProvisionedThroughputExceededException: The level of configured provisioned throughput for one
or more global secondary indexes of the table was exceeded.
{"s":"W","c":"SHARDING","msg":"Chunk migration failed","attr":{"namespace":"shop.events",
"error":"ChunkTooBig: chunk too big to move; consider resharding on a higher-cardinality key"}}
У DynamoDB при этом есть жёсткие потолки на секцию — 3000 RCU и 1000 WCU, — которые не обходятся увеличением ёмкости таблицы: adaptive capacity умеет перераспределять ёмкость и даже изолировать горячий ключ в отдельную секцию, но не мгновенно, а единичный ключ на 5000 WCU не обслужит никто.
Обнаружение. Нужны метрики по секциям (максимум, p99, энтропия распределения), а не агрегаты, и они должны существовать до инцидента. Точный подсчёт частот всех ключей в память сервиса не помещается, зато Count-Min Sketch и Space-Saving дают top-k с ограниченной ошибкой за O(1) на событие — так работают redis-cli --hotkeys и детекторы горячих ключей во всех крупных CDN (подробнее — в статье про потоковые алгоритмы). И третье: ключ секции обязан быть в структурированных логах, иначе связать «медленно» и «горячий арендатор» нечем.
Лечение: соль (salting). Разбить один логический ключ на S физических: писать в key#random(S), читать веером из всех S.
import random
SALT_BUCKETS = 16 # ×16 к стоимости чтения — осознанная плата
def write_key(base_key: str, hot: frozenset[str]) -> str:
"""Горячие ключи размазываются по S секциям, остальные не трогаем."""
return f"{base_key}#{random.randrange(SALT_BUCKETS)}" if base_key in hot else base_key
def read_keys(base_key: str, hot: frozenset[str]) -> list[str]:
return [f"{base_key}#{i}" for i in range(SALT_BUCKETS)] if base_key in hot else [base_key]
Стоимость честная: запись остаётся O(1) на одной секции, чтение становится O(S) обращений плюс слияние, поэтому соль применяют избирательно, к списку горячих ключей, обновляя список по данным heavy-hitters. Для счётчиков она почти бесплатна (результат — сумма); для «последних N записей» придётся брать N из каждой корзины и мержить, то есть амплификация S·N. Прочие приёмы, в порядке предпочтения: кеш перед секцией с коалесцированием запросов (single-flight превращает 10 000 одновременных чтений одного ключа в одно обращение к БД — дешевле любой ресегментации, см. кеширование и масштабирование); реплики только для чтения (горячее чтение лечится репликацией, горячая запись — нет); выделенная секция для арендатора-гиганта (требует справочной схемы, зато решает радикально — Slicer в Google, OSDI 2016, и Shard Manager в Meta, SOSP 2021, построены вокруг балансировки по фактической нагрузке); и наконец смена ключа — последний вариант по стоимости и первый по эффективности (reshardCollection в MongoDB 5.0+, новый топик и двойная запись в Kafka, новая таблица в Cassandra).
8. Вторичные индексы: два способа, оба неприятные
Партиционирование задано по первичному ключу, а запрос приходит по другому полю. Вариантов ровно два, и оба имеют неустранимый недостаток.
Локальный индекс (document-partitioned) — каждая секция индексирует только свои документы. Запись дешёвая, чтение — scatter/gather по всем секциям (Elasticsearch, MongoDB, обычные вторичные индексы Cassandra). Опасность не видна, пока не посчитать хвост: пусть вероятность того, что отдельная секция ответит дольше 100 мс, равна 1 %; тогда запрос по k секциям превысит 100 мс с вероятностью 1 − 0.99^k — при k = 16 это 15 %, при k = 100 уже 63 %. p99 одной секции становится медианой всего запроса — ровно то, что Dean и Barroso описали в The Tail at Scale. Отсюда: ограничивайте веер, используйте hedged requests и не заводите 500 шардов «на вырост» — каждый лишний шард удорожает каждый запрос сегодня.
Глобальный индекс (term-partitioned) — индекс партиционирован по значению индексируемого поля. Чтение дешёвое, зато запись в одну строку затрагивает две секции, то есть это распределённая транзакция. Все реализуют её асинхронно: GetItem по GSI в DynamoDB сразу после записи может вернуть старое, а недопровизионированный GSI троттлит записи в базовую таблицу — индекс становится узким местом для того, что вы даже не индексируете. Материализованные представления Cassandra годами имели статус экспериментальных ровно по той же причине. Практический вывод: глобальный индекс — это скрытая межпартиционная запись, относитесь к нему как к асинхронной репликации со всеми её свойствами.
9. Границы транзакций совпадают с границами секций
Транзакция внутри секции — обычная локальная транзакция: один WAL, одна блокировка, один fsync. Транзакция через две секции требует 2PC, а с ним координатора, окна неопределённости и заблокированных строк на время его отказа — см. распределённые транзакции.
Разница настолько велика, что вокруг неё построены продукты. Spanner позволяет объявить INTERLEAVE IN PARENT: строки дочерней таблицы физически лежат рядом с родительской в том же split, и транзакция «заказ + его позиции» становится однопартиционной — в статье OSDI 2012 это называется directory, набор строк с общим префиксом ключа и минимальная единица размещения. Citus требует colocate_with при создании распределённых таблиц, чтобы JOIN и транзакции оставались локальными. Vitess различает односегментные и кросс-шардовые запросы и по умолчанию ограничивает вторые. Правило: если две сущности должны меняться атомарно — у них обязан быть общий ключ партиционирования. Если это невозможно (перевод между счетами разных пользователей), атомарности не будет — будет сага с компенсациями, и это архитектурное решение, а не деталь реализации.
10. Почему exactly-once обычно миф — и при чём тут партиционирование
«Exactly-once доставка» невозможна как свойство канала: отправитель не отличает потерянный запрос от потерянного ответа и потому обязан либо повторить (риск дубля), либо не повторять (риск потери). Это разобрано в статье про гарантии доставки и идемпотентность; здесь важно другое — партиционирование добавляет к этому собственный слой проблем, который обычно упускают.
Границы Kafka EOS. Kafka действительно даёт транзакции: идемпотентный продюсер (дедупликация брокером по тройке PID, epoch, sequence — заметьте, посекционно), транзакционный координатор с control-маркерами и isolation.level=read_committed у потребителя. Это exactly-once обработка в сценарии «читаем из Kafka — пишем в Kafka». Это не exactly-once, если побочный эффект — вызов платёжного шлюза или запись в чужую БД: атомарности между Kafka и внешним миром не существует.
Ребалансировка потребителей ломает иллюзию:
02:14:07,551 WARN [Consumer clientId=billing-3, groupId=billing] consumer poll timeout has
expired. This means the time between subsequent calls to poll() was longer than the
configured max.poll.interval.ms
02:14:09,118 INFO [Consumer clientId=billing-1, groupId=billing] Revoke previously assigned
partitions payments-4, payments-11
02:14:12,004 ERROR CommitFailedException: Offset commit cannot be completed since the consumer
is not part of an active group for auto partition assignment; it is likely that the consumer
was kicked out of the group.
02:14:19,663 INFO [Consumer clientId=billing-7, groupId=billing] Setting offset for partition
payments-4 to the committed offset 774119002
Читается так: billing-3 обработал сообщения, но не успел закоммитить смещение — его выкинули из группы; партиция ушла к billing-7, тот начал со старого закоммиченного смещения и обработал то же самое ещё раз. Ошибки в системе нет, так работает at-least-once. Смягчается max.poll.records, кооперативной инкрементальной ребалансировкой (KIP-429, CooperativeStickyAssignor) и статическим членством (KIP-345, group.instance.id), но не устраняется.
Главное следствие, специфичное для партиционирования. Работающая замена exactly-once — at-least-once плюс дедупликация; но дедупликация есть состояние, а у состояния есть собственный ключ партиционирования. Отсюда правило, которое почти никогда не формулируют явно:
Запись факта «команда
Xуже применена» и сам эффект команды должны попадать в одну секцию, иначе для их атомарности нужен двухфазный коммит — то есть ровно то, чего мы пытались избежать.
-- Обе таблицы шардированы по account_id — значит, транзакция однопартиционная.
CREATE TABLE processed_commands (
account_id bigint NOT NULL,
command_id uuid NOT NULL,
applied_at timestamptz NOT NULL DEFAULT now(),
PRIMARY KEY (account_id, command_id)
);
BEGIN;
INSERT INTO processed_commands (account_id, command_id) VALUES (42, '5c1f9a3e-...')
ON CONFLICT DO NOTHING;
-- Приложение смотрит на число вставленных строк:
-- 0 строк -> команда уже применялась: ROLLBACK и подтверждаем сообщение;
-- 1 строка -> применяем эффект в этой же транзакции.
UPDATE accounts SET balance = balance - 100 WHERE id = 42;
COMMIT;
В Cassandra та же идея выражается через lightweight transaction, и ограничение видно ещё жёстче — LWT работает только внутри одной секции, а BEGIN BATCH с разными partition key не атомарен:
INSERT INTO ledger.processed (account_id, command_id) VALUES (42, 5c1f9a3e-...) IF NOT EXISTS;
Неправильный и очень распространённый вариант: дедупликация в Redis (SET cmd:5c1f NX EX 86400), эффект — в PostgreSQL. Между SET и COMMIT есть окно: падение процесса внутри него даёт либо потерю (ключ поставлен, эффект не применён — повтор отбросится как дубль), либо дубль (эффект применён, ключ не поставлен). Плюс Redis может вытеснить ключ по TTL или потерять его при failover, потому что репликация там асинхронная. Схема снижает частоту дублей, но не устраняет их; называть это exactly-once — самообман.
Наконец, ребалансировка ломает окна дедупликации. Идемпотентность продюсера Kafka привязана к паре (PID, партиция) — изменение числа партиций меняет назначение ключей, и «то же самое сообщение» физически едет в другую партицию, где о нём никто не знает. Таблица дедупликации, шардированная отдельно от данных, при переезде секции остаётся на месте, и какое-то время дедуп-состояние и эффект живут на разных узлах. А TTL дедупликации обязан быть больше максимального времени повтора: retention в Kafka 7 дней при TTL дедупа в час означает, что реплей за прошлые сутки применит всё заново.
Чеклист идемпотентности при партиционировании: (1) ключ идемпотентности выводится из ключа партиционирования или содержит его — f(account_id, command_id); (2) отметка о применении и эффект — в одной транзакции одной секции; (3) TTL дедуп-записи ≥ retention источника, а не «сутки, чтобы не росло»; (4) операции по возможности делаются естественно идемпотентными — SET status = 'paid' вместо INCREMENT attempts; (5) если естественной идемпотентности нет — outbox и компенсации, с пониманием, что дубли всё равно возможны и разбирать их будет бизнес-логика, а не инфраструктура.
11. Как выбирать ключ
account_id / tenant_id / user_id"] C -- нет --> E["Ключ = id сущности,
хеш-распределение"] B -- нет --> F{"Основной паттерн — диапазон по времени?"} F -- да --> G{"Есть второе поле высокой кардинальности?"} G -- да --> H["Композит: hash по нему + range по времени
плюс bucketing по дню"] G -- нет --> I["Диапазоны + автосплит,
следить за горячим хвостом"] F -- нет --> J["Scatter-gather неизбежен:
считать хвост, ограничивать веер"] D --> K{"Есть объект в 100 раз крупнее остальных?"} E --> K H --> K K -- да --> L["Соль на запись + веер на чтение
либо выделенная секция"] K -- нет --> M["Готово. Мониторить p99 ПО СЕКЦИЯМ,
а не среднее по кластеру"]
Четыре проверки, которые стоит прогнать до внедрения, а не после. Кардинальность: число различимых значений ключа должно на порядки превышать число секций — country при 128 секциях гарантирует перекос. Неизменяемость: ключ не должен меняться у существующей строки, потому что изменение — это удаление из одной секции плюс вставка в другую, то есть распределённая транзакция. Наблюдаемость: ключ секции обязан попадать в логи и в атрибуты спанов — см. наблюдаемость. Проверка хаосом: перекос и переезд секций прекрасно воспроизводятся — убить узел под нагрузкой, запустить ребалансировку в момент разделения сети, подать зипфовский трафик — см. тестирование распределённых систем.
Типичные ошибки
hash(key) % число_узловв шардировании состояния. Работает до первого масштабирования, после которого переезжает 80 % данных.- «Число партиций Kafka увеличим потом». Увеличение рвёт порядок для существующих ключей, уменьшить нельзя вовсе.
- Мониторинг средних. Средняя загрузка 30 % при одном узле на 100 % — самый частый способ не заметить перекос.
- Ключ низкой кардинальности (
status,country,event_type): секции распределятся, нагрузка нет. - Монотонный ключ при диапазонном разбиении — весь трафик в последнюю секцию.
- Секция без ограничителя роста: партиция арендатора без bucketing по времени растёт годами, а потом её нельзя ни прочитать, ни перенести.
- Передача владения без fencing-токена — два владельца, тихая потеря записей, «в логах чисто».
- Автоматическая ребалансировка без ограничителя полосы — каскад из миграций, забитой сети и ложных срабатываний детектора отказов.
ASKобрабатывается какMOVED— клиент мечется между узлами всё время миграции слота.- Дедупликация в другом хранилище, чем эффект — то, что превращает «мы сделали exactly-once» в «мы перестали замечать дубли».
- Слишком много мелких шардов (каждый scatter-gather платит за хвост каждого) и, наоборот, отсутствие pre-split при старте — одна секция и один узел под всей нагрузкой первого дня.
Мини-итог
- Партиционирование отвечает на вопрос «какие данные где», репликация — «сколько копий». Смешивать их в постмортеме дороже всего.
- Ключ партиционирования одновременно задаёт локальность, границу транзакции, границу порядка и единицу перекоса; поменять его потом — миграция, а не настройка.
hash mod Nпо числу узлов переноситN/(N+1)данных при масштабировании; консистентное хеширование с виртуальными узлами —1/(n+1), с разбросом нагрузки~1/√v.- Фиксированное число секций отвязывает карту ключей от числа узлов ценой того, что это число выбирается один раз и навсегда.
- Переезд секции безопасен только с монотонным номером владения: без фенсинга два узла принимают записи одновременно, и в логах каждого из них всё выглядит нормально.
- Автоматическая ребалансировка в связке с детектором отказов — источник каскадов; тормоз обязателен, полная автоматика при отказах — нет.
- Горячий ключ не лечится ни одной стратегией разбиения: это всегда один узел. Лечится солью, кешем с коалесцированием, выделенной секцией или сменой ключа.
- Локальный вторичный индекс превращает p99 секции в медиану запроса; глобальный — это скрытая межпартиционная запись.
- Exactly-once как свойство канала не существует. Замена — at-least-once плюс дедупликация, причём дедуп-запись обязана жить в той же секции, что и эффект, иначе вы вернулись к двухфазному коммиту.
Источники
- Karger et al., Consistent Hashing and Random Trees (STOC 1997) — оригинал консистентного хеширования.
- Thaler, Ravishankar, A Name-Based Mapping Scheme for Rendezvous (1996) — HRW-хеширование.
- Lamping, Veach, A Fast, Minimal Memory, Consistent Hash Algorithm (2014) — jump consistent hash.
- Eisenbud et al., Maglev (NSDI 2016) — хеширование в балансировщиках нагрузки.
- DeCandia et al., Dynamo (SOSP 2007) — кольцо, виртуальные узлы, preference list.
- Corbett et al., Spanner (OSDI 2012) — splits, directories, интерливинг ради локальности.
- Adya et al., Slicer: Auto-Sharding for Datacenter Applications (OSDI 2016) и Shard Manager (SOSP 2021) — ребалансировка по нагрузке.
- Dean, Barroso, The Tail at Scale (CACM 2013) — почему scatter-gather дорог.
- Martin Kleppmann, «Designing Data-Intensive Applications», глава 6; Redis Cluster Specification — слоты,
MOVED/ASK, хеш-теги; KIP-429 — кооперативная ребалансировка потребителей Kafka. - AWS: Designing partition keys to distribute your workload evenly — соль, adaptive capacity, лимиты на секцию.
- Родственное на портале: Репликация и шардирование в БД, Cassandra и wide-column, NewSQL и распределённые БД.
Что дальше
Мы дважды упёрлись в одно и то же: карта секций должна быть согласованной, а номер владения — монотонным. И то, и другое невозможно без механизма, который заставляет группу узлов договориться о единственном значении, переживая отказы и разделения сети. Дальше — про этот механизм: Консенсус: Paxos, Raft, выбор лидера, репликация лога.