Distributed Key-Value Store (Dynamo/Cassandra-style) with consistent hash ring (RF=3, W=2, R=2 quorum), vector-clock conflict resolution, hinted handoff, anti-entropy via Merkle trees, and read repair. Includes 5 scenarios: PUT quorum, GET with read repair, partition + hinted handoff, Merkle anti-entropy, concurrent-write conflict resolution.
Распределенное key-value хранилище лежит в основе сессий, feature flags, user profiles, shopping carts, metadata stores, counters, time-to-live данных и многих high-throughput сервисов. Dynamo-like дизайн учит выбирать между availability и consistency, строить consistent hashing ring, реплицировать данные без single point of failure, восстанавливаться после partition и объяснять, почему "простая hash map на кластере" быстро превращается в сложную distributed database.
Типичный scale assumption: 100 nodes, hundreds of virtual nodes per physical node, replication factor 3, миллионы операций в секунду, средний value несколько килобайт, storage до petabyte-scale. Ключевые требования: get, put, delete, TTL, горизонтальное добавление nodes без downtime, tunable consistency, no SPOF и понятное поведение при conflicts.
Этот кейс важен еще и потому, что Dynamo-style подход встречается в реальных системах: Amazon DynamoDB, Cassandra, Riak, ScyllaDB и внутренние KV stores. Он показывает, как система сознательно покупает availability ценой eventual consistency и conflict resolution.
Связанные темы: ::concept{slug="partitioning-strategies"}, ::concept{slug="replication"}, ::concept{slug="quorum-reads-writes"}, ::concept{slug="gossip-protocol"}, ::concept{slug="vector-clocks"}, ::concept{slug="consistency-models"}.
Ментальная модель: весь keyspace раскладывается по consistent hash ring. Физические nodes получают много virtual nodes, чтобы нагрузка и данные распределялись ровнее. Для ключа система считает hash, находит token range и preference list: например replicas A, B, C при RF=3. Любая node может быть coordinator: она принимает запрос клиента, определяет replicas и собирает нужные acknowledgements.
Quorum читается как пересечение множеств. При N=3, W=2, R=2 любая успешная запись и успешное чтение пересекаются хотя бы в одной replica. Это не магия linearizability во всех случаях, но в обычном случае помогает увидеть последнюю версию. При network partition и sloppy quorum система может принять запись на fallback node, а потом через hinted handoff доставить ее исходной replica.
На каждом node storage engine обычно LSM-based: write-ahead log для durability, memtable для быстрых writes, immutable SSTables на диске, bloom filters для пропуска отсутствующих keys и compaction для слияния файлов. Поэтому writes быстрые, а reads могут требовать проверки нескольких SSTables, если compaction отстает.
Диаграмма показывает smart client, cluster ring, coordinator, replicas, per-node engine и background operations. Smart SDK может знать ring topology и отправлять запрос сразу подходящей node, но даже без этого любая node может принять запрос и стать coordinator.
Группа cluster ring показывает nodes A, B, C, D, E. Node A часто выступает coordinator, B и C replicas, D fallback для sloppy quorum, E joining node для online rebalance. Это не значит, что роли жестко закреплены: в реальном кластере каждая node владеет множеством token ranges и для разных keys может быть leader/coordinator/replica.
Per-node engine показывает WAL, memtable, SSTable и bloom filter. Это объясняет write path и read path: запись сначала durable append в WAL, затем memtable, позже flush в SSTable. Read сначала проверяет memtable и bloom filters, затем ищет ключ в relevant SSTables.
Background ops показывают gossip, anti-entropy, hinted handoff, compaction и failure detector. Именно эти фоновые процессы превращают набор nodes в самовосстанавливающуюся систему.
PUT W=2 учит quorum write. Coordinator пишет локально, реплицирует на B и C, но отвечает клиенту после двух подтверждений. Третья replica может догнаться позже. Это снижает latency и сохраняет доступность при одном сбое, но оставляет окно divergence.
GET R=2 показывает read repair. Coordinator читает две replicas, видит, что одна отстала, выбирает доминирующую версию по vector clock и асинхронно чинит stale replica. Это полезно, но работает только для читаемых keys; холодные keys требуют anti-entropy.
Concurrent writes показывают siblings. Если два клиента одновременно пишут разные значения в разные replicas, vector clocks могут показать, что ни одна версия не доминирует. Система не должна молча выбрасывать одну запись. Она возвращает обе версии, а client/application делает merge, например объединяет shopping cart.
Sloppy quorum + hinted handoff учит availability under failure. Если replica C недоступна, coordinator может записать временную копию на D с hint для C, ответить клиенту, а потом replay hint. Это повышает write availability, но усложняет guarantees: до replay данные лежат не на своем canonical owner.
Anti-entropy через Merkle trees показывает ремонт без пользовательского read traffic. Replicas сравнивают hashes token ranges, спускаются к отличающимся leaves и синхронизируют только различающиеся keys.
Add node scenario показывает consistent hashing. При добавлении E не нужно перемещать весь dataset; переезжает часть ranges. Но online migration требует throttling, dual writes или hints на время cut-over.
Availability против consistency. Настройки ONE дают минимальную latency и максимум доступности, но stale reads и lost updates становятся вероятнее. QUORUM дороже, но лучше для пользовательских данных. ALL сильнее, но теряет availability при одном падении replica.
Vector clocks против last-write-wins. LWW просто реализовать и удобно для кэшей, но clock skew и concurrent writes могут silently lose data. Vector clocks сохраняют обе версии, но переносят merge logic в application и добавляют metadata overhead.
Consistent hashing с vnodes упрощает rebalance, но усложняет observability: одна физическая node владеет сотнями ranges, и hot key может ломать баланс даже при ровном распределении tokens.
LSM engine дает высокий write throughput, но compaction может съедать диск и IO. Если compaction lag растет, read amplification увеличивается, bloom filters занимают память, а p99 latency ухудшается.
Tunable consistency гибкая, но опасна для продуктовой модели. Если разные сервисы читают один bucket с разными consistency levels, пользователи могут видеть противоречивые состояния.
Amazon Dynamo описал идеи consistent hashing, vector clocks, sloppy quorum и hinted handoff. DynamoDB как managed service скрывает большую часть operational details, но concepts похожи: partitioning, replication, provisioned/on-demand throughput, hot partitions.
Apache Cassandra использует consistent hashing, replication strategies, tunable consistency, LSM storage, gossip и repair. ScyllaDB совместима с Cassandra API и оптимизирована под high throughput на shard-per-core архитектуре. Riak исторически был близок к Dynamo и использовал siblings/vector clocks.
Redis Cluster тоже key-value и sharded, но его модель другая: in-memory, hash slots, primary-replica, обычно не Dynamo-style quorum. Etcd/Consul дают strong consistency через Raft, но хуже подходят для massive high-throughput arbitrary values.
Первый anti-pattern: использовать KV store для сложных ad-hoc queries. Если нужно фильтровать по множеству полей и делать join, нужна другая модель или secondary index с отдельными trade-offs.
Второй anti-pattern: выбрать partition key с низкой cardinality. Ключи вроде tenant_id для огромного tenant или country создают hot partitions. Нужны composite keys, bucketing или write sharding.
Третий anti-pattern: включить LWW для данных, где потеря concurrent update недопустима. Shopping cart, collaborative editing и counters требуют merge-aware модели или CRDT.
Четвертый anti-pattern: не запускать repair. Eventual consistency не означает "само когда-нибудь исправится" для холодных ключей. Нужны anti-entropy jobs, read repair и monitoring divergence.
Пятый anti-pattern: считать quorum абсолютной linearizability. Без consensus per key, clock discipline и строгой coordination остаются edge cases при partitions, retries и stale coordinators.
Не используйте Dynamo-like KV, если нужен SQL, joins, complex transactions и глобальные constraints. Реляционная база или NewSQL может быть правильнее.
Не используйте eventual KV для денег, inventory с жесткими ограничениями или уникальных username без дополнительного coordination. Там нужны conditional writes, transactions или consensus-backed service.
Не используйте distributed KV ради маленького объема данных. Один Postgres с индексом или Redis primary-replica будет проще, дешевле и надежнее в эксплуатации.
Не используйте KV как search engine. Point lookup by key быстрый, но поиск по тексту и facets требуют отдельного индекса.
Читайте Dynamo paper, Cassandra architecture, LSM tree, SSTables, bloom filters, compaction strategies, hinted handoff, read repair и Merkle tree repair. Затем разберите PACELC и CAP для понимания latency/consistency trade-offs: ::concept{slug="pacelc-theorem"}, ::concept{slug="consistency-models"}, ::concept{slug="quorum-reads-writes"}. Для практики полезно сравнить DynamoDB, Cassandra, ScyllaDB, Redis Cluster и etcd.