Every row in Apache Cassandra lives on several nodes at once. Replication is what makes a node failure, a rack failure or even a lost datacenter a non-event, and it is also the source of most of the surprising behaviour people meet in production: reads that miss recent writes, a rack outage that takes down quorum, or a replication change that quietly makes data unreadable until repair finishes.
This article explains replication as three connected decisions. Placement: which nodes hold the copies of a partition. Fan-out: how a write reaches those copies, including across datacenters. Convergence: how copies that missed a write catch up. It ends with the operational runbooks for changing the replication factor and adding or removing a datacenter. Consistency levels are touched on only as far as replication needs them; the full quorum arithmetic is in Cassandra consistency levels.
The replication factor and the strategy
Replication is configured per keyspace with two things: a strategy class, which decides placement, and a replication factor (RF), the number of copies. With NetworkTopologyStrategy (NTS) you give an RF per datacenter. SimpleStrategy takes a single RF and ignores datacenters and racks completely.
-- Two datacenters, three replicas in each. DC names must match what the snitch reports.
CREATE KEYSPACE orders
WITH replication = {'class': 'NetworkTopologyStrategy', 'dc1': 3, 'dc2': 3};
-- Where does one partition live? Run from any node.
-- nodetool getendpoints orders order_items 'cust-4711'
-- Development only: SimpleStrategy ignores datacenters and racks.
CREATE KEYSPACE scratch
WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 1};Use NTS for every keyspace that holds real data, even in a single-datacenter cluster. SimpleStrategy places copies on consecutive ring positions with no regard for racks, so two of three copies can land on the same rack, and moving to a second datacenter later means changing strategy on live data. Since Cassandra 4.0, NTS also accepts a replication_factor option that expands to the same RF in every datacenter, which is convenient but easy to misread when DCs need different values; spell out DC names in production.
The datacenter and rack names come from the snitch, which every node uses to describe itself and its peers. A typo between the snitch configuration and the keyspace definition means zero replicas in the misspelled DC, and writes at LOCAL_QUORUM there fail. The snitch article covers how nodes learn topology.
Placement: walking the ring
Each partition key is hashed by the partitioner to a token. Each node owns several tokens on a ring (16 per node by default since Cassandra 4.0). To find the replicas, Cassandra starts at the first node token at or after the partition's token and walks clockwise. SimpleStrategy takes the next RF distinct nodes. NTS walks the same ring but, for each datacenter, prefers nodes on racks it has not used yet. When it meets a node on a rack already holding a copy, it skips it and remembers it; only once every rack in that DC has a copy does it fall back to the skipped nodes, in ring order. The path from key to token is covered in Cassandra partitioning.
This small simulation reproduces the rule for one or more datacenters:
from bisect import bisect_left
def nts_replicas(ring, key_token, rf_by_dc):
# ring: list of (token, node, dc, rack) sorted by token
tokens = [t for t, *_ in ring]
start = bisect_left(tokens, key_token) % len(ring) # first node at or after the token
racks_in = {}
for _, _, dc, rack in ring:
racks_in.setdefault(dc, set()).add(rack)
chosen = {dc: [] for dc in rf_by_dc}
seen_racks = {dc: set() for dc in rf_by_dc}
skipped = {dc: [] for dc in rf_by_dc}
for i in range(len(ring)):
_, node, dc, rack = ring[(start + i) % len(ring)]
if dc not in rf_by_dc or node in chosen[dc] or len(chosen[dc]) >= rf_by_dc[dc]:
continue
if rack not in seen_racks[dc]:
chosen[dc].append(node)
seen_racks[dc].add(rack)
if seen_racks[dc] == racks_in[dc]: # every rack used: take skipped nodes
while skipped[dc] and len(chosen[dc]) < rf_by_dc[dc]:
chosen[dc].append(skipped[dc].pop(0))
elif seen_racks[dc] == racks_in[dc]:
chosen[dc].append(node)
else:
skipped[dc].append(node) # same rack again; hold for later
return chosen
ring = [(10, "n1", "dc1", "a"), (20, "n2", "dc1", "a"), (30, "n3", "dc1", "b"),
(40, "n4", "dc1", "b"), (50, "n5", "dc1", "c"), (60, "n6", "dc1", "c")]
print(nts_replicas(ring, 15, {"dc1": 3})) # {'dc1': ['n2', 'n3', 'n5']}Worked example. Six nodes in dc1, two per rack, with tokens 10 to 60 and racks a, a, b, b, c, c. A key hashes to token 15. The walk starts at n2 (rack a) and takes it. n3 (rack b) is a new rack and is taken. n4 is rack b again, so it is skipped. n5 (rack c) is taken, giving n2, n3 and n5: one copy per rack. SimpleStrategy on the same ring would take n2, n3 and n4, two of them on rack b. Lose rack b and that partition has one live copy; LOCAL_QUORUM needs two, so reads and writes for it fail. With NTS the same outage leaves two copies and quorum holds.
Rack awareness has a cost. If racks have unequal node counts, nodes on the small rack own a larger share of the data, because every partition needs a copy on that rack. Keep the number of racks equal to RF or a multiple of it, with the same number of nodes in each. With vnodes, allocate_tokens_for_local_replication_factor in cassandra.yaml makes new nodes pick tokens that balance ownership for the given RF; set it to your DC's RF before bootstrapping.
Fan-out: how a write reaches every replica
The client sends a write to any node, which becomes the coordinator for that request. A token-aware driver picks a node that is itself a replica, saving one hop. The coordinator sends the mutation to all replicas, in every datacenter, regardless of the consistency level. The consistency level only decides how many acknowledgements it waits for before answering the client.
For remote datacenters the coordinator does not send RF separate messages over the WAN. It sends one message to a single replica in each remote DC, along with the list of the other remote replicas, and that node forwards the mutation inside its own DC. This keeps WAN traffic to one copy per write per remote DC, which matters when you size inter-DC links.
Reads differ. The coordinator asks the closest replica, as ranked by the dynamic snitch, for the full data and asks enough other replicas for a digest, a hash of their answer, to satisfy the consistency level. If the digests disagree, it fetches full data, resolves by timestamp, and writes the newest version back to the stale replicas before answering. That is blocking read repair.
Consistency levels in one paragraph
A quorum is floor(RF / 2) + 1 copies. With RF 3 in the local DC, LOCAL_QUORUM waits for 2, so it tolerates one replica down; reads and writes both at LOCAL_QUORUM overlap on at least one replica, so a read sees the latest acknowledged write in that DC. RF 2 is a trap: quorum is 2, so losing one node makes LOCAL_QUORUM fail, and it gives no more availability than RF 1 at quorum while costing twice the disk. ONE waits for any single replica and is fast but can read stale data. EACH_QUORUM needs a quorum in every DC for writes and so fails during a WAN partition.
| RF per DC | LOCAL_QUORUM needs | Replicas that can be down | Typical use |
|---|---|---|---|
| 1 | 1 | 0 | Dev, or data you can rebuild |
| 2 | 2 | 0 at quorum, 1 at ONE | Rarely a good choice |
| 3 | 2 | 1 | The production default |
| 5 | 3 | 2 | Very large clusters, or two simultaneous failures in one DC |
Convergence: how missed writes catch up
Replicas miss writes all the time: a node restarts, a GC pause outlasts the write timeout, a WAN link flaps. Three mechanisms repair the gaps, each covering a different window.
- Hinted handoff. When a replica does not acknowledge, the coordinator stores a hint and replays it when the replica returns, but only for outages shorter than the hint window, three hours by default. Hinted handoff covers sizing and replay.
- Read repair. Fixes the partitions that are read, at the time they are read. Partitions nobody reads stay divergent.
- Anti-entropy repair. Compares data between replicas range by range and streams the differences. It is the only mechanism that guarantees convergence, and it must complete on every range within
gc_grace_seconds(10 days by default), or deleted data can come back. See Cassandra repair.
The practical consequence: replication gives you copies, but the copies are only identical if repair runs on schedule. A cluster that never repairs is relying on luck for any partition that is written during an outage longer than the hint window.
Runbooks: changing RF and datacenters
Changing replication settings is a metadata change that takes effect immediately, while the data movement it implies does not happen on its own. That gap is where outages come from.
# Raise RF in dc1 from 2 to 3.
cqlsh -e "ALTER KEYSPACE orders WITH replication =
{'class': 'NetworkTopologyStrategy', 'dc1': 3, 'dc2': 3};"
# New replicas own ranges they have no data for until repair streams it to them.
for h in $(dc1_hosts); do ssh "$h" nodetool repair --full orders; done
# Lower RF in dc1 from 3 to 2: alter, then drop data nodes no longer own.
cqlsh -e "ALTER KEYSPACE orders WITH replication =
{'class': 'NetworkTopologyStrategy', 'dc1': 2, 'dc2': 3};"
for h in $(dc1_hosts); do ssh "$h" nodetool cleanup orders; done
# Add dc3: start its nodes empty, add it to every keyspace's replication,
# then stream existing data into each new node from an existing DC.
cqlsh -e "ALTER KEYSPACE orders WITH replication =
{'class': 'NetworkTopologyStrategy', 'dc1': 3, 'dc2': 3, 'dc3': 3};"
for h in $(dc3_hosts); do ssh "$h" nodetool rebuild -- dc1; done- Raising RF. After the ALTER, the new replica for each range is responsible for data it does not have. Reads at
ONEorLOCAL_ONEthat land there return nothing until repair streams the data. Raise RF during quiet hours, read at LOCAL_QUORUM until the full repair finishes, and raise one step at a time. - Lowering RF. Safe for availability. Nodes keep data they no longer own until
nodetool cleanupremoves it, so disk use does not drop until then. - Adding a datacenter. Bring up the new nodes with the correct DC name in the snitch configuration, without them bootstrapping data, then alter every keyspace, including system keyspaces such as
system_authandsystem_distributed, and runnodetool rebuildon each new node with an existing DC as the source. Keep clients pinned to the old DCs until rebuild finishes. - Removing a datacenter. Move clients away, alter every keyspace to drop the DC, then decommission its nodes; with RF 0 there, they have nothing to stream.
- The auth keyspace.
system_authdefaults to a low RF with SimpleStrategy. In a multi-node cluster, change it to NTS with several replicas per DC and repair it; otherwise one node down can lock everyone out.
Transient replication and other variations
Cassandra 4.0 added transient replication as an experimental feature, disabled by default. It lets some of a keyspace's replicas hold data only until incremental repair has copied it to the full replicas, reducing disk use, but it has significant restrictions, for example on materialized views and secondary indexes. Treat it as something to test, not to rely on, and check the documentation for your exact version.
Replication does not protect against logical mistakes. A DROP TABLE, a bad batch of updates or a buggy client is replicated faithfully to every copy within milliseconds. Snapshots and backups are the defence; see multi-DC topology for how DC layout interacts with that.
Failure modes
- Rack count not matching RF. Two racks with RF 3 means one rack holds two copies of every partition; losing it loses quorum everywhere.
- DC name mismatch. The keyspace says
DC1, the snitch saysdc1, and LOCAL_QUORUM writes fail with too few replicas. - RF raised without repair. Reads at LOCAL_ONE return empty results for existing rows.
- New DC added without rebuild. The new DC answers queries for data it never received.
- Repair skipped beyond gc_grace_seconds. Deleted rows reappear because a replica that missed the tombstone still holds the old value.
- EACH_QUORUM on writes across a WAN. Any DC partition becomes a write outage.
What to do next
- List every keyspace with
DESCRIBE KEYSPACESand its replication settings, and convert any SimpleStrategy keyspace holding real data to NTS. - Check that each DC has racks equal to RF or a multiple of it, with equal nodes per rack.
- Run
nodetool getendpointsfor a few hot keys and confirm their replicas sit on different racks. - Set
allocate_tokens_for_local_replication_factorto your RF before adding nodes. - Raise
system_authto several replicas per DC and repair it. - Schedule anti-entropy repair to complete within gc_grace_seconds, and alert when it does not.
- Write the RF-change and add-DC runbooks above into your operations docs and rehearse them on staging.