You can use Cassandra for years treating it as a fast key-value store with a SQL-like syntax. Then a deleted row reappears, or a balance goes wrong after a datacentre hiccup, or one partition melts a node, and you find that the database was teaching distributed systems all along. This page is the summary of those lessons: what Cassandra makes you learn about time, deletion, consistency, failure detection, data placement, retries and load. Each lesson is tied to the mechanism that causes it, the failure you see in production, and a drill you can run to prove your system handles it.

It is a capstone rather than an introduction. The arithmetic of consistency levels is covered in CAP in practice, and Cassandra's ancestry in the Dynamo paper's influence. Here the question is: what should a team that runs Cassandra, or any leaderless replicated store, have internalised?

The lessons at a glance

The seven lessons, at a glance:

LessonMechanismFailure you seeDrill
Time is an input to every writelast-write-wins by cell timestampan update silently lostskew one app host's clock in staging
Deletes are writestombstones, gc_grace_secondsdeleted data resurrectsskip repair past gc_grace on a test table
Quorum overlap is not linearizabilityQUORUM reads/writes, read repaira failed write becomes visible laterfail writes mid-flight and read back
Failure detection is a guessgossip, phi accrual, hintsslow node treated as up, or flappinginject latency, not just crashes
The partition is the unit of everythingtoken ring, partition keyone hot or huge partitionload test with the real key distribution
Retries need idempotenceclient retries, countersdouble-counted incrementskill connections during writes
Load shedding beats queuingcoordinators, speculative retrylatency collapse under backlogoverload one node and watch p99

Lesson 1: time is an input to every write

Cassandra resolves conflicting writes to the same cell by timestamp: the highest timestamp wins, and on an exact tie a tombstone beats a value. The timestamp is microseconds, supplied by the client driver or, if the client does not set one, by the coordinator. There is no version vector and no merge. That makes writes cheap, and it means the wall clock of whichever machine stamped the write is part of your data.

Worked example. Two app servers update a user's email. Server X's clock is 300 ms fast. At real time 10:00:00.000, X writes a@example.com stamped 10:00:00.300. At real time 10:00:00.100, server Y writes b@example.com stamped 10:00:00.100. Y's write happened later, was acknowledged, and loses forever, because X's timestamp is higher. Nothing logs an error.

The lessons generalise. Keep clocks tightly synchronised and alert on offset. Do not use read-then-write patterns to implement updates whose order matters; if two writers can race on one cell and the order matters, you need a lightweight transaction or a data model where each writer appends its own row. And treat explicit USING TIMESTAMP as a sharp tool: it is useful for idempotent replays, and a single bad value can shadow a cell for years. Clocks in distributed systems explains why physical time cannot order events across machines.

Lesson 2: deletes are writes

In a system where any replica may miss a write, a delete cannot simply remove data: the replica that missed the delete would later hand the old value back. So a delete writes a tombstone, a marker with a timestamp that shadows older data. Tombstones are kept for gc_grace_seconds (default 864000, ten days) and only then become eligible to be purged by compaction.

The diagram shows what that window is for. Repair is the process that carries the tombstone to the replica that missed it. If a full repair of that token range does not complete within gc_grace_seconds, the other replicas can purge the tombstone, and the next repair or read repair finds the stale value on the lagging replica and treats it as live. Deleted data comes back.

Lesson 2 in one picture: a delete forgotten by one replica comes backReplica AReplica BReplica Cwrite v1ts=100write v1ts=100write v1ts=100tombstonets=200tombstonets=200C was downnever got itcompactionafter gc_grace: purgedcompactionafter gc_grace: purgedstill has v1no repair ranread / repairv1 is newest leftv1 copied backTimeline left to right. The tombstone lives gc_grace_seconds (default 864000 s, 10 days) so repair can carry it to C.If no full repair reaches C inside that window, A and B forget the delete, and anti-entropy treats C's v1 as live data.Fix: complete a repair of every token range more often than gc_grace_seconds, and alert when it does not.
Zombie data: replicas A and B purge a tombstone that replica C never received, and repair copies C's old value back.

Two more costs follow from tombstones. Reads must scan past them, so a queue-like table that inserts and deletes heavily degrades until reads hit the tombstone failure threshold. And range deletes and TTLs create tombstones too. The operational lessons: run repair on a schedule shorter than gc_grace_seconds and alert when a range misses it; do not lower gc_grace_seconds unless repair is reliably faster; and model data so that whole partitions expire or are dropped together. Tombstones and repair have the details.

Lesson 3: quorum overlap is not linearizability

With replication factor 3, writing and reading at QUORUM means each touches 2 replicas, and since 2 + 2 is greater than 3, every read set overlaps every successful write set. Teams often stop there and call the system strongly consistent. It is not linearizable, for two reasons.

Failed writes are not rolled back. If a QUORUM write reaches one replica and then times out, the client gets an error, but that replica keeps the value. A later read that includes that replica returns the new value, and blocking read repair writes it to the others. An operation the client was told failed takes effect, possibly after a different client has already read the old value.

Concurrent reads can disagree. While a write is in flight, one reader can see the new value and a slightly later reader the old one, depending on which replicas each contacted. Read repair narrows this window; it does not close it. Note that since Cassandra 4.0 the probabilistic read_repair_chance options are gone; blocking read repair at quorum levels remains.

When you genuinely need compare-and-set, use a lightweight transaction (IF NOT EXISTS, IF col = ?), which runs Paxos on a single partition. It costs several extra round trips and serialises contention on that partition. Cassandra 4.1 added a reworked Paxos implementation (paxos_variant: v2) that reduces the cost; measure it on your version rather than trusting a number. Multi-partition ACID transactions (the Accord work, CEP-15) are under development and were not part of the 5.0 release. See lightweight transactions.

Lesson 4: failure detection is a guess

Nodes learn about each other through gossip, and each node decides whether a peer is down with a phi accrual failure detector: instead of a fixed timeout, it computes a suspicion level from the history of heartbeat intervals and convicts when it crosses phi_convict_threshold (default 8). The lesson is that a failure detector produces an opinion, not a fact, and different nodes can hold different opinions at the same moment.

While a replica is considered down, coordinators store hints for it and replay them when it returns, up to max_hint_window (default 3 hours). Hints are an optimisation, not durability: they live on the coordinator, they stop being written after the window, and a coordinator that dies loses them. A node down longer than the window needs repair.

The production failure is rarely a clean crash. It is a node that is slow: a GC pause, a failing disk, a saturated NIC. It keeps answering gossip, so it is not convicted, and it drags every request that waits on it. Your drills should inject latency and packet loss, not only kill processes.

Worked example. In a six-node cluster with replication factor 3, one node starts taking two-second old-generation GC pauses every minute. Between pauses its heartbeats arrive normally, so its phi value spikes briefly and falls back under 8; it is never marked down, so coordinators keep choosing it. With token-aware routing, a third of all partitions have it as a replica, and a QUORUM read that picks it as one of its two replicas waits out the pause. Cluster-wide p99 read latency jumps from 8 ms to 2 s while every dashboard shows six nodes up. What helps: dynamic snitching, which steers reads towards replicas with better recent latency; speculative retry, covered in lesson 7; and alerting on per-node GC pause time and per-node coordinator latency, not just node status. The fix is on the node (heap sizing, collector choice), but the lesson is about detection: liveness and health are different signals, and you need both.

Lesson 5: the partition is the unit of everything

The partition key is hashed to a token, and the token picks the replicas. Everything follows from that: a partition lives entirely on its replicas, is read and compacted as a unit, and cannot be split. So a partition that is hot gets the throughput of 3 nodes no matter how many you add, and a partition that grows without bound (all events for a tenant, forever) makes compaction, repair and reads slower until it becomes an incident. Keep partitions bounded, commonly well under about 100 MB, by adding a time bucket or a shard suffix to the key:

-- unbounded: one partition per sensor, forever
CREATE TABLE readings_bad (sensor_id text, ts timestamp, value double,
  PRIMARY KEY (sensor_id, ts));

-- bounded: one partition per sensor per day
CREATE TABLE readings (sensor_id text, day date, ts timestamp, value double,
  PRIMARY KEY ((sensor_id, day), ts))
  WITH CLUSTERING ORDER BY (ts DESC)
   AND default_time_to_live = 2592000;   -- 30 days; whole partitions expire together

The deeper lesson is that you design tables for queries, not entities: one table per access pattern, denormalised, because the only efficient read is one partition. Load test with your real key distribution, because a uniform synthetic key will never show the hot partition production has.

Lesson 6: retries need idempotence

Drivers retry. A timeout does not tell you whether the write applied; it often did. For most Cassandra writes that is harmless, because an INSERT or UPDATE that sets values is idempotent: applying it twice with the same timestamp leaves the same state. Make sure the timestamp is fixed on the client before the first attempt, so a retry does not create a newer version that shadows an intervening write.

Some operations are not idempotent. Counter increments, list appends and LWTs can apply twice, and drivers should be configured not to retry them speculatively. If you need an exact count, write one row per event keyed by an event ID and aggregate, rather than incrementing a counter. The general lesson, which applies to every system on this site, is that at-least-once delivery plus idempotent operations is the only retry design that survives timeouts.

Worked example. A loyalty service awards points with UPDATE points SET balance = balance + 50 WHERE user_id = ? on a counter column. A coordinator accepts the write, the replicas apply it, and the response is lost when a network switch flaps. The driver's retry policy, or the application's own retry loop, sends it again, and the user now has 100 extra points. Nothing in the cluster is wrong; the client asked twice. The idempotent design writes INSERT INTO point_events (user_id, event_id, delta) VALUES (?, ?, 50) with event_id taken from the order that earned the points, so a retry rewrites the same row, and the balance is the sum of the user's rows, computed on read or by a periodic job. You trade a slightly more expensive read for a value you can audit and rebuild.

Lesson 7: shed load instead of queuing it

When a node falls behind, requests queue on its coordinators and latency climbs for everyone who touches its ranges. Two mechanisms help. Speculative retry sends a read to an extra replica if the first is slower than a table-level threshold (a percentile such as 99p by default), trading a little extra load for a much shorter tail. Shedding is the other half: a request that cannot finish before the client gives up should be dropped, not executed. Cassandra drops messages that exceed their timeout rather than processing them late; your clients should enforce deadlines, cap in-flight requests per host, and back off instead of retrying immediately into an overloaded cluster.

A useful game day for all seven lessons:

# staging cluster, RF=3, one keyspace under synthetic load
nodetool status                       # baseline: all UN
# lesson 1, on ONE APP HOST (drivers stamp writes client-side by default):
sudo date -s "+2 seconds"             #   then race two writers on the same cell
# the rest run on one Cassandra node
nodetool disablegossip; sleep 30      # lesson 4: node looks down; watch hints accumulate
nodetool enablegossip
nodetool repair -pr my_ks             # lesson 2: confirm repair completes inside gc_grace
tc qdisc add dev eth0 root netem delay 200ms   # lessons 4 and 7: slow, not dead
nodetool tpstats                      # look for dropped messages and pending queues
tc qdisc del dev eth0 root

Record what the application saw at each step: errors, latency, and any value that changed when it should not have. Restore normal time sync after the skew test.

What to do next

  1. Alert on clock offset on every app server and every Cassandra node.
  2. Confirm a full repair of every range completes in less than gc_grace_seconds, and alert when it does not.
  3. List every place you assume QUORUM means linearizable; move true compare-and-set logic to LWTs.
  4. Find unbounded partitions with nodetool tablehistograms and add a bucket to the key.
  5. Set client-side timestamps and mark non-idempotent statements so the driver does not retry them speculatively.
  6. Enforce client deadlines and per-host in-flight limits.
  7. Run the game day above in staging once a quarter.
Key takeaway: Cassandra makes distributed-systems theory concrete: write timestamps decide conflicts, so clocks are data; deletes are tombstones that only repair can spread; quorum overlap does not make failed writes disappear; failure detection is probabilistic; the partition bounds throughput and size; retries are safe only for idempotent operations; and overload must be shed. Tie each lesson to an alert and a drill, and the surprises become planned events.