Cassandra (Apache) — concept page. Wide-column store, masterless ring, gossip-based membership, consistent hashing on token ring, tunable consistency (LOCAL_QUORUM/QUORUM/ONE/ALL), LSM storage (CommitLog + Memtable + SSTables), compaction strategies (STCS/LCS/TWCS), Lightweight Transactions via Paxos. Anti-entropy via hinted handoff, read repair, nodetool repair. Three scenarios: write path with CL=QUORUM (RF=3, W=2 via coordinator + 3 replicas), read repair healing a stale replica via timestamp reconciliation, hinted handoff during partition with coordinator stashing and replaying mutations. Three ADRs covering masterless vs leader-based replication, Cassandra vs ScyllaDB vs DynamoDB tool selection, and LOCAL_QUORUM as production default for tunable consistency.
Cassandra — оригинальный leaderless wide-column store, спроектированный в Facebook (2008) на пересечении идей Dynamo (Amazon, 2007) и Bigtable (Google, 2006). Сегодня знание Cassandra — это базовый класс для понимания целого семейства систем: ScyllaDB, AstraDB, Bigtable, HBase, частично DynamoDB. Три причины разобраться.
RF=3, CL=QUORUM, вы понимаете 80% теории distributed databases.Дорогой урок: команды приходят из Postgres и пытаются использовать Cassandra «как SQL без транзакций» — нормализованные таблицы, JOIN client-side, ALLOW FILTERING в каждом query. Потом обнаруживают tombstone hell, scatter-gather по всему кластеру, p99 в секундах. Этот документ — про то, как сразу не наступить.
«Cassandra — это leaderless ring of LSM-нод с tunable quorum. Дизайн схемы начинается от queries (denormalize всё), partition key выбирает ноду через consistent hashing, clustering key сортирует rows внутри партиции, а background compaction постоянно пересобирает SSTables. Под капотом — Dynamo (распределение) + Bigtable (storage engine).»
Три измерения, в которых принимаются все решения:
R + W > N — единственный способ получить strong consistency для одного ключа. Всё остальное — eventual с разным lag.nodetool repair в течение gc_grace_seconds (10 дней) — получаете «zombie data» (удалённые rows возвращаются). Не следите за wide partitions (>100 MB) — compaction страдает, p99 деградирует. Не настроили compaction strategy под workload — read amplification растёт линейно с числом SSTables.Любое решение по Cassandra — точка в этом пространстве. LOCAL_QUORUM reads + LOCAL_QUORUM writes на RF=3 (multi-DC), STCS для write-heavy / LCS для read-heavy / TWCS для time-series, nodetool repair раз в неделю — это не «опции», это обязательный production baseline.
На канвасе — типовая Cassandra installation: один DC, кластер из пяти нод, RF=3.
node-1 ... node-5), у каждой свой token range на консистентном хеш-кольце [0, 2^64). Ключ → токен → ближайший по часовой стрелке owner + следующие RF-1 реплики. Edges между нодами — представление gossip-меша (каждая нода обменивается digest'ом со случайным peer каждую секунду; полный mesh не рисуем, чтобы не загромождать).CommitLog (WAL, fsync перед ack), Memtable (in-RAM sorted skiplist), Bloom filter (in-RAM, false positive ~1%, отсекает SSTables, в которых ключа точно нет), три SSTable (immutable on-disk, newest first). Compaction — фоновая операция, на схеме обозначена edges sst1 → sst2 → sst3.Edges:
client → n1 — CQL соединение, n1 в этом примере выступает координатором.n1 ↔ n2/n3/n4/n5 плюс другие пары — gossip-меш (peer-to-peer).n2 → commitlog (fsync), n2 → memtable (insert), memtable → sst1 (flush), n2 → bloom → sst1/sst2/sst3 (read path), sst1 → sst2 → sst3 (compaction).Каждая нода в реальности — это копия этой LSM-структуры. На диаграмме показан только один «срез», чтобы не повторять одно и то же пять раз.
write-quorum — Канонический write path при CL=QUORUMClient → coordinator (n1) → hash(partition_key) → token=147 → owners: n3 (primary), n5, n2. Coordinator отправляет мутацию параллельно всем трём; на каждой реплике сначала CommitLog (fsync ~1ms), потом Memtable (in-RAM insert). Как только 2 из 3 acks дошли — coordinator отвечает клиенту OK. Третий ack приходит async, не блокирует latency. Когда memtable наполняется (~256 MB) — фоновый flush в SSTable.
Чему учит: R + W > N — единственная формула, которую надо помнить. R=2, W=2, N=3 → любой последующий read с CL=QUORUM пересечётся хотя бы с одной репликой, видевшей write. Это и есть «strong consistency для одного ключа». CL=ONE даёт ~1ms latency, но stale reads возможны; CL=ALL даёт максимум consistency, но availability падает до AND всех реплик — одна нода вниз = writes fail.
read-repair — Async-коррекция стейла во время чтенияCoordinator (n1) для CL=QUORUM выбирает 2 ближайшие реплики: с n3 запрашивает full row, с n5 — только digest (md5 hash для bandwidth-efficient compare). Digest не сходится — n5 stale. Coordinator делает второй запрос (full row) к n5, видит timestamp T1 < T2 (n3 новее), резолвит last-write-wins, возвращает клиенту свежее значение. Параллельно отправляет n5 мутацию-починку через нормальный write path. Клиент про divergence не знает.
Чему учит: last-write-wins по timestamp — основной (и единственный из коробки) механизм разрешения конфликтов в Cassandra. Это работает только если у нод синхронизированные часы (NTP обязателен, clock skew = silent data corruption). read_repair_chance (default <1.0) означает, что не каждое чтение триггерит repair — поэтому nodetool repair раз в неделю обязателен, иначе entropy накапливается и tombstone'ы могут expire'нуть до того, как достигнут отстающей реплики (= deleted data возвращается, «zombie»).
hinted-handoff — Availability во время partitionn5 падает (или партиция изолирует n1 от n5). Клиент шлёт write с CL=QUORUM на тот же ключ. Coordinator (n1) пытается отправить мутацию на n2, n3, n5; первые две успешны, n5 — connection refused. Coordinator локально записывает «hint»: «доставить эту мутацию на n5, когда она снова жива». W=2 удовлетворено, клиент получает OK — availability сохранена несмотря на отказ узла. Когда n5 поднимается, gossip распространяет ALIVE state за O(log N) раундов; n1 читает свой hint store и проигрывает накопленные мутации. Если n5 отсутствует дольше max_hint_window_in_ms (default 3 часа) — hints выбрасываются, единственный механизм reconciliation — nodetool repair.
Чему учит: это и есть «AP» в CAP. Под partition Cassandra предпочитает принять write (с локальным hint) и догнать позже, а не отказать. Trade-off: hint живёт на coordinator'е; если coordinator падает до доставки hint'а — hint потерян. Финальная сетка безопасности — periodic nodetool repair строго до истечения gc_grace_seconds (10 дней по умолчанию).
Status: accepted.
Context. Два мира распределённой OLTP. Leader-based (Postgres streaming, MongoDB replica set, MySQL Group Replication, Spanner / CockroachDB на уровне range): один узел координирует writes per shard, простая ментальная модель, легко получить strong consistency и ACID, но failover занимает секунды/минуты, лидер — потолок write throughput. Leaderless / masterless (Dynamo, Cassandra, Riak, ScyllaDB): каждый узел принимает каждый write, реплицирует на N peers, никакого failover не нужно — некого терять. Цена: tunable-не-default consistency, last-write-wins резолюция конфликтов, обязательные read-repair и nodetool repair для борьбы с entropy.
Decision. Masterless выбираем когда: (1) write throughput выше чем способен сустейнить один primary (Discord держал ~120M сообщений/день до миграции на Scylla), (2) multi-DC active-active (active-active с leader'ом — это уже Spanner/CockroachDB класс), (3) толерантность к single-node failure без паузы на election. Cassandra явно выбирает AP в CAP: при partition каждая minority-партиция продолжает принимать writes (с hint'ами), reconciliation позже по timestamp. Это неправильно для inventory или balances; правильно для messages, time-series, sensor data, write-heavy логов и фидов.
Consequences. Никакой ACID-транзакции «из коробки» — есть только LWT (Paxos round, 4× latency) для compare-and-set на одной партиции. Никакого referential integrity. Schema design жёстко привязан к access pattern. Operational ответственность смещена с «failover automation» (как у leader-based) на «anti-entropy hygiene» — repair, compaction tuning, monitoring tombstones.
Status: accepted (per-workload).
Context. Три wide-column store с одинаковой data model и кардинально разным operational профилем. Cassandra — Java, JVM GC pauses, зрелый OSS, runs anywhere; default-выбор когда нужен контроль над железом. ScyllaDB — C++ переписка на Seastar framework, per-core sharding (shared-nothing), DPDK NIC, нет GC; drop-in CQL replacement, ~5-10× throughput per node, ~5× ниже p99. Discord в 2022 мигрировал message store с 177 Cassandra нод на 72 Scylla. DynamoDB — AWS managed, проприетарный, pay-per-request или provisioned; ноль operations, auto-sharding, multi-region Global Tables — но lock-in, дорогой выше ~50K rps, нет CQL/SQL (только GetItem/Query).
Decision. Дерево решений:
| Условие | Выбор |
|---|---|
| AWS-shop, <50K rps sustained, want zero ops | DynamoDB |
| On-prem / multi-cloud / >100K rps, команда комфортна с CQL | ScyllaDB |
| Нужен максимум ecosystem (Spark, Solr/Elastic через DSE), есть Cassandra ops expertise, или железо с очень малым числом cores где per-core модель Scylla проигрывает | Cassandra |
В 2026 году выбор Cassandra для greenfield large-scale deployment без серьёзной оценки ScyllaDB — anti-pattern: тот же API, кардинально лучше $/rps. DynamoDB anti-pattern: hot partition (один key >3K rps) silently throttles, нет способа вырастить per-key throughput без redesign схемы.
Consequences. Переход Cassandra ↔ Scylla — смена бинарника, не приложения (CQL совместимость). Переход на/с DynamoDB — переписывание data access layer (нет CQL). Команды должны оценивать полный TCO включая ops headcount, не только лицензию/cloud bill.
LOCAL_QUORUM как production defaultStatus: accepted.
Context. Cassandra выставляет CL per-query: ONE, TWO, THREE, QUORUM, ALL + DC-aware варианты LOCAL_ONE, LOCAL_QUORUM, EACH_QUORUM, ANY. Математика: если R + W > N (реплики) для одного ключа — strong consistency для этого ключа. CL=ONE — самый быстрый (один ack), stale reads возможны. CL=ALL — самый строгий, но любая нода вниз = ошибка. QUORUM = floor(N/2)+1, толерирует 1 ноду вниз на RF=3. EACH_QUORUM — QUORUM в каждом DC, durable но добавляет cross-DC RTT.
Decision. Production default на multi-DC кластере — LOCAL_QUORUM на чтение и запись. Причины: (1) толерирует одну ноду вниз per DC без потери writes, (2) избегает cross-DC WAN на каждой операции (это +50-150ms p99), (3) даёт R+W>N внутри локального DC — read-your-writes within DC. Trade-off: write, успешный в DC-A, может ещё не быть реплицирован в DC-B; subsequent LOCAL_QUORUM read в DC-B может вернуть stale до прибытия async replication (обычно <1s).
Для глобально-строгих требований (например, уникальность email при регистрации) — LWT (Paxos, 4× latency) или внешний координатор. CL=ONE — только для analytics / observability writes где потеря приемлема; CL=ALL — только для одноразовых admin операций (availability проседает до AND всех реплик).
Consequences. App-код должен явно документировать места, где CL ≠ LOCAL_QUORUM. LWT — silver bullet, но 4× медленнее обычного write; не использовать на high-throughput путях. На multi-DC active-active дизайн должен учитывать: один и тот же ключ может одновременно меняться в двух DC, last-write-wins по timestamp решит конфликт (= тихая потеря).
Общий паттерн: write-heavy, schema стабильна, eventual consistency приемлема, нужен многонодовый scale без single-leader bottleneck. Большинство этих компаний сегодня либо мигрировали часть на ScyllaDB (Discord), либо рассматривают, либо построили внутренние Cassandra-like решения (Uber Schemaless, FB Tao).
PRIMARY KEY ((user_id, YYYY-MM), message_id).ALLOW FILTERING в production. Делает scatter-gather по всему кластеру — каждая нода сканирует свои SSTables. Латентность в секундах, нагрузка на весь кластер. Если query не помещается в (partition_key, clustering_key) — нужна другая таблица с другим PK, не FILTERING.gc_grace_seconds (10 дней). Range read с >100K tombstones → query fails. Для time-series — TWCS + TTL, чтобы старые window выбрасывались целиком.nodetool repair. В течение gc_grace_seconds (10 дней) каждая реплика должна получить tombstone. Иначе tombstone expire'ит локально, потом приходит «оживший» row с старой реплики → zombie data. Repair раз в 5-7 дней — обязательная фоновая работа.LOCAL_QUORUM. Каждое чтение через WAN = +50-150ms p99. Используйте LOCAL_QUORUM по умолчанию, EACH_QUORUM только когда globally durable обязательно.ConsistencyLevel=ONE для критичных writes. Один ack ≠ durable — если получившая нода падает до replication, write потерян. ONE приемлемо для observability/analytics, не для бизнес-логики.IF NOT EXISTS и IF column=X — 4× latency из-за Paxos round (prepare → propose → commit). Использовать точечно (regional uniqueness, idempotency на критичных операциях).