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 levelacks (RF = 3)use
ONE1fastest; risk of stale reads
QUORUM2 (across all DCs)strong reads and writes together
LOCAL_QUORUM2 in the local DCthe usual production choice in multi-DC
ALL3rarely: one node down blocks the request
CASSANDRA: RING, COORDINATOR, QUORUM
a write at QUORUM with replication factor 3
clientdriver, token-awarecoordinatorany nodereplica 1token range ownerreplica 2replica 3 (down)hintstored for r3
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