Стратегии партицирования: hash, range, composite, geographic + hot key проблемы и live resharding (Pinterest pattern). 4 сценария.
Один сервер не вмещает данные — либо по объёму, либо по нагрузке. Партицирование (sharding) разбивает датасет на куски по ключу, и каждый кусок живёт на отдельной ноде. Но выбор ключа — это не вкусовщина: ошибиться в стратегии = получить hot partition, scatter-gather на каждом запросе, или 6-месячный resharding-проект с downtime.
sharding как концепт объясняет «зачем». Этот урок — про детали: range vs hash vs composite vs directory vs geographic, как избегать hot tail и celebrity keys, как Pinterest и Discord перешивают кластер вживую без потери единого write.
«Стратегия партицирования = выбор того, что чаще всего трогает запрос. Если запросы по
user_id— партицируй поuser_id. По времени — по time bucket. По геолокации — по region. Партиционный ключ должен совпадать с самым горячим access pattern, иначе каждый запрос превращается в scatter-gather.»
Три закона:
@elonmusk_posts, tenant_id=enterprise_giant, region=us-east-1 — даже хеш не спасёт, если 80% нагрузки на один ключ.Канвас — упрощённый партиционированный кластер: Application → Router/Coordinator → 4 шарда. Координатор знает mapping key → partition и инкапсулирует стратегию. Реальные системы:
mongos процесс роутит запросы к нужному mongod.Диаграмма — образ. Mapping может быть встроен в клиент (Cassandra), вынесен в lookup-сервис (Pinterest директория), или жить в metadata cluster (Vitess + etcd). Важно — что где-то решение «какой шард» принимается перед первым диск-сиком.
Четыре сценария показывают ровно те ситуации, которые ломают наивные дизайны: идеальный hash, range hot tail, celebrity hot key и live resharding.
hash-even — идеальный hashWrites на разных user_id → hash(id) % 4 равномерно раскладывает по 4 шардам. ~25% нагрузки на каждый, никаких hot spots. Это baseline: hash partitioning по высококардинальному ключу с равномерным распределением запросов. DynamoDB, Cassandra, MongoDB hashed sharding — все начинают здесь.
Почему работает: хорошая hash-функция (MurmurHash3, xxHash) даёт uniform output даже на skewed input. 64-битный hash space + modulo по N шардов = approximately N/total нагрузки на каждый шард.
Когда ломается: см. следующие два сценария.
range-hot-tail — классический антипаттернПартиция по месяцу: Jan→p1, Feb→p2, ..., May→p4 (текущий). Все today-writes идут в один partition. Остальные три простаивают. p4 уходит в throttling, 503 на клиента.
Это самая частая ошибка в системах с time-series данными (events, logs, metrics). Range partitioning сам по себе не плох — он идеален для «дай всё за июль». Но партицировать только по timestamp без compound — гарантированный hot tail.
Фикс: compound key (bucket, timestamp), где bucket = hash(user_id) % 16. Теперь writes размазаны по 16 buckets × time, при этом range query по конкретному user_id всё ещё работает (один bucket → один шард, range по времени внутри шарда). Cassandra-паттерн: ((user_id), message_ts) — partition key размазывает, clustering key сортирует внутри.
hot-celebrity — hash не панацея@elonmusk_posts хешируется в один конкретный шард. Hash распределение «равномерное» в среднем — но один ключ, который генерирует 80% reads, утопит свой шард независимо от количества partitions. p3 на 100% CPU, p1/p2/p4 idle.
Mitigation hierarchy (от простого к сложному):
(elonmusk_posts, shard_id) для writes; reads scatter-gather (но scatter по фиксированному N — приемлемо).Ключевая мысль: hash защищает от skew в ключевом пространстве, но не от skew в access pattern.
live-resharding — Pinterest/Discord patternКогда все шарды в red zone (disk 95%+, CPU saturating), нужно расти. Из 4 шардов в 8 — без downtime, без потери write.
Канонический алгоритм:
key → {old_shard, new_shard} — слой абстракции в coordinator. Старая логика: hash(key) % 4, новая: hash(key) % 8.Discord так перешивал свои Cassandra/ScyllaDB кластера для триллионов сообщений. Pinterest — для MySQL шардов с фотографиями. Vitess умеет это автоматически (workflow MoveTables / Reshard).
ADR-001: Hash partitioning by default; composite (partition_key, clustering_key) для range-внутри-сущности; никогда не голый timestamp.
user_id, ~15% — «последние N событий юзера», ~5% — аналитика по time range.user_id (point lookups покрыты). Для «последние N событий» — composite key (user_id, event_ts): partition key даёт шард, clustering key сортирует. Аналитика по time range — отдельный OLAP store (ClickHouse / BigQuery), куда CDC льёт events; OLTP не партицируется по времени никогда.created_at — hot tail, отброшено сразу. (b) Directory-based (Pinterest стиль) — lookup сервис как SPOF, операционная сложность не оправдана для greenfield. (c) Geographic-only — нужен мульти-региональный продукт, у нас single-region MVP.ADR-002: Resharding только live (dual-write/backfill/cutover), никогда STOP-THE-WORLD.
| Система | Стратегия | Детали |
|---|---|---|
| DynamoDB | Hash + range | Hash partition key + optional sort key. Adaptive capacity для hot partitions, automatic resharding под капотом. |
| Cassandra / ScyllaDB | Token ring (hash) + clustering | ((partition_key), clustering_key). Virtual nodes (vnodes) для финер-grained распределения. Любая нода = coordinator. |
| Discord | Composite | ((channel_id), message_ts) на ScyllaDB. Триллионы сообщений, читают по channel + time range — partition key даёт шард, clustering сортирует. |
| Vitess (YouTube, Slack, Etsy) | Hash/range/lookup vindexes | MySQL шарды + vtgate router. Online schema changes и resharding через MoveTables / Reshard workflow. |
| PostgreSQL declarative | Range/list/hash | Native с PG 10+. PARTITION BY RANGE/LIST/HASH. Pgroonga, Citus для distributed sharding поверх. |
| MongoDB | Hashed / ranged | mongos router + config servers. Auto-balancing chunks между шардами. |
| Pinterest (historical) | Directory-based | Lookup table «key → shard» в собственном сервисе поверх MySQL. Гибко, но lookup = SPOF. |
| Geographic + composite | Geo-pinned user data + (user_id, tweet_ts) для timeline-like запросов. | |
| Bigtable / Spanner | Lexicographic ranges | Auto-splitting tablets когда range распухает. Spanner добавляет global TrueTime для cross-shard transactions. |
(bucket_hash, ts) — обязательно.country_code при N=4 шардах — большинство трафика «US» утопит один шард, остальные три почти пусты. Высококардинальный ключ обязателен (user_id, request_id, заведомо много значений).