Every row in Cassandra lives on a small set of nodes, its replicas. The replication factor says how many there are; the replication strategy says which ones. That second decision is easy to ignore because it happens silently, inside every coordinator, for every read and write. It is also the decision that determines whether losing a rack loses data, whether one node quietly carries twice the load of its neighbours, and whether a quorum read can still succeed during a zone outage.

This article is about the placement algorithm itself. It walks through how SimpleStrategy and NetworkTopologyStrategy choose replicas, step by step, including the detail most explanations skip: what happens to nodes that are passed over because their rack is already used. It gives a small Python model you can run against your own topology, a worked example that shows how an uneven rack layout creates a hot spot, and the commands that show where a given key really lives. Replication factors, write fan-out and the procedure for changing them are covered in the replication overview; multi-datacenter design is covered in the multi-DC topology guide.

Starting from the ring

Placement starts from the token ring. Each node owns one or more tokens, which are positions on a ring of 64-bit values (with the default Murmur3 partitioner). A partition key is hashed to a token, and the node whose token is the first one at or after the key's token, moving clockwise, is the first candidate. Everything else follows from walking the ring from that point. How tokens are assigned, and why vnodes exist, is covered in the token ring and vnodes article; here, all that matters is that the walk produces an ordered list of distinct nodes.

SimpleStrategy takes the first RF distinct nodes from that walk. It does not look at racks or datacenters at all. With RF 3 and the ring order n1, n2, n3, it picks n1, n2 and n3, even if all three sit in the same rack, behind the same switch and the same power feed. That is why SimpleStrategy is only appropriate for a single-node development box or a throwaway test cluster. In production it gives you three copies that can fail together.

NetworkTopologyStrategy (NTS) is configured per datacenter: {'class': 'NetworkTopologyStrategy', 'dc1': 3, 'dc2': 3}. It runs the walk once per datacenter, considering only that datacenter's nodes, and within each datacenter it tries to put every replica in a different rack. The datacenter and rack of each node come from the snitch, normally GossipingPropertyFileSnitch reading cassandra-rackdc.properties; see the snitch article for how that information spreads through the cluster.

How NetworkTopologyStrategy walks a datacenter

Within one datacenter, the NTS walk follows these rules. They are a plain-language version of the per-datacenter logic in current Cassandra source (the DatacenterEndpoints class); 3.x expressed the same idea with a list of skipped nodes, and gives the same replicas.

  1. Walk the ring clockwise from the key's token, looking only at nodes in this datacenter, and skip any node already chosen.
  2. If the node's rack has not been used yet, take it as a replica and mark the rack as used.
  3. If the rack has already been used, take the node only while there is an allowance of rack repeats left. The allowance starts at RF minus the number of racks in the datacenter, so it is zero whenever there are at least RF racks. Otherwise pass over the node.
  4. Stop when RF replicas have been chosen (or every node in the datacenter, if there are fewer nodes than RF).

Two consequences follow directly. When the datacenter has at least RF racks, every replica lands in a different rack, and losing one rack costs at most one replica of each row. When it has fewer racks than RF, some rack must hold two replicas, and the walk decides which nodes those are: the first repeat-rack nodes met clockwise from the key's token, not a random choice. The walk is also deterministic. Every coordinator that knows the same topology computes the same replica list, without any coordination.

One DC, 6 nodes, 3 racks, RF 3: walk clockwise from the key's tokenn1rack an2rack an3rack bn4rack bn5rack cn6rack ctoken herereplica 1n1, rack a newpassed overn2, rack a seenreplica 2n3, rack b newpassed overn4, rack b seenreplica 3n5, rack c newnot neededRF reachedWith fewer racks than RF, the first RF minus racks repeat-rack nodes met are taken as well.SimpleStrategy would have taken n1, n2, n3: two replicas in rack a.
NetworkTopologyStrategy within one datacenter. Nodes whose rack is already represented are passed over, unless the datacenter has fewer racks than RF.

A placement model you can run

The rules are short enough to model. This function computes the NTS replica set for one datacenter. The ring is a list of (token, node, rack) tuples sorted by token, with one token per node to keep it readable; the real implementation handles vnodes the same way, because it walks tokens and skips nodes it has already chosen. Checked against a skip-list version of the 3.x logic on thousands of random topologies, it returns the same replica sets.

def nts_replicas(ring, start, rf):
    """ring: [(token, node, rack)] sorted by token, one datacenter."""
    rf_left = min(rf, len(ring))
    repeats = rf - len({rack for _, _, rack in ring})  # rule 3 allowance
    i = next((k for k, (t, _, _) in enumerate(ring) if t >= start), 0)
    replicas, seen_racks = [], set()
    for step in range(len(ring)):
        _, node, rack = ring[(i + step) % len(ring)]
        if node in replicas:                    # rule 1: vnodes repeat nodes
            continue
        if rack not in seen_racks:              # rule 2: new rack, always take
            seen_racks.add(rack)
        elif repeats > 0:                       # rule 3: allowed repeat
            repeats -= 1
        else:
            continue                            # rule 3: pass over
        replicas.append(node)
        if len(replicas) == rf_left:            # rule 4
            break
    return replicas

To see how much data each node is responsible for, sum the width of every token range a node replicates:

from collections import Counter

def ownership(ring, rf, ring_size=1000):
    share = Counter()
    for k, (token, _, _) in enumerate(ring):
        width = (token - ring[k - 1][0]) % ring_size or ring_size
        for node in nts_replicas(ring, token, rf):
            share[node] += width
    return {n: round(100 * v / ring_size) for n, v in sorted(share.items())}

Feeding in your real topology (from nodetool ring) is a cheap way to check a planned expansion before you run it.

Worked example: the uneven rack

Take a datacenter with six nodes split evenly over three racks, a, b and c, two nodes per rack, with evenly spaced tokens and RF 3. The model gives every node exactly 50 percent effective ownership: each rack holds one full copy, split between its two nodes. Losing any whole rack leaves two replicas of every row, so LOCAL_QUORUM still succeeds everywhere.

Now suppose one node in rack c is decommissioned to save money, leaving five nodes: n1 and n4 in rack a, n2 and n5 in rack b, and n3 alone in rack c, with tokens at 0, 200, 400, 600 and 800. Running the model gives:

NodeRackEffective ownership (RF 3)
n1a40%
n2b40%
n3c100%
n4a60%
n5b60%

Node n3 holds a full copy of the entire datacenter's data. The rule is mechanical: every replica set must include one node from each of the three racks, and rack c only has one node to offer. That node now receives a write for every mutation in the datacenter, carries 2.5 times the data of n1, compacts 2.5 times as much, and is the slowest to repair and to stream when replaced. Meanwhile n1 and n2 carry less than before, so adding capacity elsewhere does nothing for the hot node.

The fix is a balanced layout, not a bigger n3: either add a node to rack c, or rebuild with racks that each hold the same number of nodes. The same arithmetic explains the general rule. With RF equal to the rack count, every rack holds exactly one copy of everything, so each rack's nodes share 1/RF of the replica load among themselves, and the smallest rack sets the per-node maximum.

Planning rules that follow from the algorithm

From the algorithm, a few planning rules follow:

  • Racks per datacenter should equal RF, or be a multiple of it. Three racks for RF 3 is the common choice, and in clouds the natural mapping is one rack per availability zone. With fewer racks than RF, some rows have two replicas in one rack, so a single rack failure can take down a quorum for those rows.
  • Keep the same number of nodes in every rack. Grow and shrink in multiples of the rack count: add three nodes, one per rack, rather than one at a time.
  • One rack per datacenter is acceptable, because NTS then degenerates to a plain ring walk inside the datacenter, which balances fine. What hurts is a few racks of unequal size.
  • Do not create racks with only one node unless every rack has one node. A single-node rack is a single node that owns 1/RF of every copy.
  • Never change a running node's rack or datacenter. Cassandra checks the stored values at startup and refuses to start if they differ. The -Dcassandra.ignore_rack=true and -Dcassandra.ignore_dc=true overrides exist, but using them silently changes which nodes own which ranges without moving any data. The safe way to move a node is to decommission it and bootstrap it again with the new location.
  • Use NTS from the start, even with one datacenter. Moving from SimpleStrategy later means altering the keyspace and running a full repair.

Configuring location and replication

A node's location comes from cassandra-rackdc.properties when the snitch is GossipingPropertyFileSnitch:

# conf/cassandra-rackdc.properties on a node in zone eu-west-1b
dc=eu-west
rack=eu-west-1b

The keyspace then names each datacenter exactly as the snitch reports it. Datacenter names are case-sensitive and a mismatch is the most common placement bug. Cassandra 4.0 and later reject an unknown datacenter name when you create or alter a keyspace; older versions accepted it and stored zero replicas in the datacenter you meant.

CREATE KEYSPACE orders WITH replication = {
  'class': 'NetworkTopologyStrategy',
  'eu-west': 3,
  'us-east': 3
};

From 4.0, NTS also accepts 'replication_factor': 3 at creation time and expands it to every datacenter that exists at that moment. It is convenient for scripts, but listing datacenters explicitly makes the intent obvious in code review and protects you from a datacenter that happened to exist when the script ran. To keep token ownership even when nodes join, Cassandra 4.0 added allocate_tokens_for_local_replication_factor in cassandra.yaml; set it to the datacenter's RF so the allocator picks tokens that balance replicated ownership rather than primary ownership.

Verifying where data actually lives

Three commands answer most placement questions:

# Which nodes hold this partition key? (keyspace, table, key)
nodetool getendpoints orders order_items 'a1f3c9e2-2b7d-4c1e-9f3a-0c5d6e7f8a9b'

# Effective ownership per node for this keyspace's replication settings
nodetool status orders

# Token ranges with their replica endpoints
nodetool describering orders

nodetool status without a keyspace argument shows ownership only when every keyspace shares the same replication, and otherwise warns that the figures are not meaningful. Always pass the keyspace you care about. In the output, compare the Owns (effective) column across nodes in the same datacenter: on a balanced cluster the values sit within a few percent of each other, and a node far above its neighbours is the symptom from the worked example.

Check rack spread too. Take a sample of keys, run getendpoints for each, map the returned addresses to racks using nodetool status, and confirm no replica set contains two nodes from one rack.

Failure modes

The failures that come from placement, rather than from the nodes themselves, are predictable:

  • SimpleStrategy in production. Replicas ignore racks, so a top-of-rack switch failure can make some partitions unavailable at QUORUM. Fix: alter to NTS with the snitch's datacenter name, then run a full repair so the new replicas actually hold the data.
  • Uneven racks. A hot node with much higher load, disk and compaction than its peers. Fix: add nodes to the small rack or rebalance rack membership through decommission and bootstrap.
  • Fewer racks than RF. Two replicas share a rack for some ranges, so a rack outage can drop LOCAL_QUORUM for exactly those ranges, and only those, which makes the incident look random.
  • Datacenter name typo on an older version. Writes acknowledge at LOCAL_ONE in the correct datacenter while the intended second copy never exists.
  • Changed rack in the properties file. The node refuses to start; forcing it with ignore_rack leaves its data on the wrong ranges until a repair and cleanup are run.
  • Missing repair after a placement change. Altering replication changes the replica set immediately for new requests, but existing data is not copied. Reads at ONE can return nothing from a new replica. Run repair, as described in the repair guide, before relying on it.

Trade-offs

Rack awareness is not free. Placing replicas in different availability zones adds cross-zone latency to every quorum operation and, in most clouds, cross-zone transfer charges for every write and every repair stream. The alternative, putting all replicas in one zone, is cheaper and faster until the zone fails. For almost every production workload, one rack per zone with RF 3 is the right default. Single-zone placement is reasonable only for data you can rebuild, such as caches or derived tables.

Higher replication factors buy availability and read capacity at the cost of disk and write amplification: each extra replica is another full copy and another write per mutation. RF 5 with five racks survives two rack failures at QUORUM, but costs two-thirds more storage than RF 3.

What to do next

  1. Run DESCRIBE KEYSPACES and check every application keyspace uses NetworkTopologyStrategy with correctly spelled datacenter names.
  2. Run nodetool status <keyspace> for your largest keyspace and compare effective ownership within each datacenter.
  3. Count nodes per rack. If they are unequal, plan additions that make them equal before the next growth step.
  4. Load your real ring into the Python model and test the next planned expansion or removal before you touch the cluster.
  5. Set allocate_tokens_for_local_replication_factor to the datacenter RF on new nodes.
  6. Add a post-change check that samples keys with nodetool getendpoints and fails if two replicas share a rack.
Key takeaway: NetworkTopologyStrategy walks the ring per datacenter, takes one node per new rack, and repeats a rack only when RF exceeds the rack count. Use it everywhere, keep racks equal in size and equal to or a multiple of RF, never change a node's rack in place, and verify placement with nodetool status and getendpoints after every topology change.