Part 2 · 1 chapters · ~8 min
Cassandra Architecture
The ring, tokens and virtual nodes, the partitioner, gossip and the phi accrual failure detector, snitches, racks and data centres, the coordinator, tunable consistency levels with quorum maths, hinted handoff, read repair and anti-entropy with Merkle trees, lightweight transactions with Paxos, and the commit log.
5
Leaderless, tunable, repaired
code
-- consistency is chosen per query
CONSISTENCY LOCAL_QUORUM; -- quorum within the local data centre: low latency, strong-ish
INSERT INTO events (account_id, ts, type) VALUES (?, ?, ?);
-- lightweight transaction: compare-and-set using Paxos (4 round trips: use sparingly)
INSERT INTO usernames (name, user_id) VALUES ('ada', ?) IF NOT EXISTS;
UPDATE accounts SET limit_kobo = 30000000 WHERE id = ? IF limit_kobo = 20000000;| consistency level | acks (RF = 3) | use |
|---|---|---|
| ONE | 1 | fastest; risk of stale reads |
| QUORUM | 2 (across all DCs) | strong reads and writes together |
| LOCAL_QUORUM | 2 in the local DC | the usual production choice in multi-DC |
| ALL | 3 | rarely: one node down blocks the request |
CASSANDRA: RING, COORDINATOR, QUORUM
a write at QUORUM with replication factor 3
swipe the figure sideways, or tap expand for full screen
1/5
the ring
Nodes own token ranges on a ring (with virtual nodes); a partition key hashes (Murmur3) to a token, and the replication strategy picks RF nodes clockwise, across racks and data centres (Distributed Systems P12).
partition key → token → RF replicasNetworkTopologyStrategy across racks