Topic 11.3
Cassandra: Query-First Modeling, Consistency Levels and Tombstones
In one line
Cassandra distributes rows by partition key over a token ring, replicates each partition to RF nodes, and orders rows inside a partition by clustering columns. You design one table per query, denormalising freely, keeping partitions bounded (~100 MB, ~100K rows is a common guide). Tunable consistency (ONE, QUORUM, LOCAL_QUORUM) sets the latency–consistency trade per request; deletes create tombstones that can slow reads.
Think of it like this
A post office sorting letters by PIN code into bags (partition key), each bag sorted by house number (clustering key). Delivering to one street is fast; finding every letter mentioning "birthday" means opening every bag.
Key ideas
- 01
Primary key = ((partition key columns), clustering columns). Queries must specify the full partition key; clustering columns support equality and range in order, and ORDER BY only in clustering order. No joins;
ALLOW FILTERINGis a warning sign. - 02
Write path: commit log + memtable → flushed to immutable SSTables → compaction (size-tiered, leveled, time-window). Writes are cheap; there's no read-before-write. Updates and inserts are upserts.
- 03
Read path: memtable + SSTables with Bloom filters, partition index and summary; results merged by timestamp (last write wins per cell).
- 04
Consistency: RF = 3 per DC; write and read at LOCAL_QUORUM gives overlap within a DC. Hinted handoff, read repair and
nodetool repair(anti-entropy) restore convergence; run repairs withingc_grace_seconds(10 days default) or deleted data can resurrect. - 05
Tombstones: deletes and TTL expiries write markers kept until gc_grace passes and compaction purges them. Queue-like patterns (insert, then delete) leave partitions full of tombstones, and reads that scan them time out. Use time-bucketed partitions and TWCS so whole SSTables expire.
Code & diagrams
CREATE KEYSPACE chat WITH replication =
{'class': 'NetworkTopologyStrategy', 'dc_mumbai': 3, 'dc_singapore': 3};
-- query: newest messages in a conversation, paged
-- bucket by month to keep partitions bounded
CREATE TABLE chat.messages_by_conversation (
conversation_id uuid,
month text, -- '2026-09'
message_id timeuuid,
sender_id uuid,
body text,
PRIMARY KEY ((conversation_id, month), message_id)
) WITH CLUSTERING ORDER BY (message_id DESC)
AND compaction = {'class': 'TimeWindowCompactionStrategy',
'compaction_window_unit': 'DAYS', 'compaction_window_size': 1};
CONSISTENCY LOCAL_QUORUM;
SELECT message_id, sender_id, body FROM chat.messages_by_conversation
WHERE conversation_id = 5b6962dd-3f90-4c93-8f61-eabfa4a803e2 AND month = '2026-09'
LIMIT 50;
-- a second query needs a second table: messages by user for a "my messages" page
CREATE TABLE chat.messages_by_sender (
sender_id uuid, month text, message_id timeuuid, conversation_id uuid, body text,
PRIMARY KEY ((sender_id, month), message_id)
) WITH CLUSTERING ORDER BY (message_id DESC);Interview problem
The problem
Model IoT readings in Cassandra
200K devices each send a reading every 10 s (20K writes/sec). Queries: latest readings for a device, readings for a device in a time range, and daily averages per device. Keep 180 days. Design the tables.
When it breaks
Unbounded partition (all messages of a group chat in one partition)
What you see
The partition grows to gigabytes; compaction and repair struggle, reads time out, and one node pair becomes hot.
Fix & prevent
Add a time bucket to the partition key (month or day) and page across buckets in the application.
Repairs not run within gc_grace_seconds
What you see
A replica that missed a delete keeps the old value after the tombstone is purged elsewhere; the deleted data comes back (zombie data).
Fix & prevent
Schedule repairs (Cassandra Reaper, or incremental repair) more often than gc_grace_seconds.
Explain it without notes
Why does Cassandra modeling start from queries, not entities?
Practice
Compute partition size for 1 message/sec in a conversation, 300 bytes each, bucketed by month.
Trade-offs
- ↔
Linear write scale and multi-DC availability, paid for with query rigidity, denormalisation, tombstones and repair operations.
Done when you can
I can design Cassandra tables per query with bounded partitions and choose consistency levels and compaction.