Every row in Cassandra belongs to exactly one partition, and every partition lives on replicas chosen by a hash of its key. That mapping decides which nodes serve a query, how evenly load spreads and how much work repair, compaction and streaming do. When a cluster misbehaves, the cause is often a partition that is too large, too hot, or both.
This article follows one partition key from the application to its replicas: what the partitioner computes, how a token becomes a replica set, how drivers route with it, how to scan by token range, and what partition size and traffic do to the system. Primary-key design and bucketing arithmetic are covered in Cassandra data modeling basics, and token allocation in Cassandra vnodes.
What a partition is
A CQL primary key has two parts. The partition key, the first element or the parenthesized group, decides placement. The clustering columns that follow decide the order of rows inside the partition. All rows that share a partition key are stored together, sorted by clustering key, on every replica that owns that key.
The partition is therefore the unit of three things at once. It is the unit of placement: a partition is never split across nodes, so it can never outgrow one node. It is the unit of locality: a query that names one partition key reads from one replica set and, within each SSTable, from one contiguous region. And it is the unit of atomicity and isolation for writes: writes to one partition are applied atomically and in isolation on each replica, while writes that span partitions have weaker guarantees.
From key to token: the partitioner
The partitioner is a function from partition key bytes to a token, a position on a ring. The default since Cassandra 1.2, and the right choice for new clusters, is Murmur3Partitioner. It computes a 128-bit MurmurHash3 of the key bytes and keeps the first 64 bits as a signed long, so tokens lie between -2^63 and 2^63 - 1.
The bytes that are hashed depend on the key's shape. For a single-column partition key, the input is simply the serialized value of that column. For a composite partition key, Cassandra builds one byte string from all components: for each one, a two-byte length, the value bytes, and a single 0x00 end-of-component byte. That is why token(sensor_id, day) must name every component in order, and why changing a column's type, even between types that look compatible, changes every token in the table.
Two details matter for tooling. The hash value -2^63 is normalized to 2^63 - 1, so -2^63 never appears as a partition token and serves as the ring's starting sentinel. And Cassandra's Murmur3 differs from some stock libraries for certain inputs, so ask the cluster with token() rather than trusting a generic Murmur3 package.
| Partitioner | Token space | Status |
|---|---|---|
Murmur3Partitioner | signed 64-bit, -2^63 to 2^63 - 1 | default, use for new clusters |
RandomPartitioner | MD5, 0 to 2^127 - 1 | legacy; slower hash, kept for old clusters |
ByteOrderedPartitioner | raw key bytes, ordered | legacy; enables key-range scans but creates hot spots |
The partitioner is fixed for the life of a cluster; changing it moves every partition, so the only migration is a copy into a new cluster. Hashing gives up key order in exchange for a near-uniform spread, and Cassandra recovers ordering where you need it through clustering columns inside a partition.
From token to replicas: the ring
Each node owns one or more tokens. Sorted, those tokens cut the ring into ranges, and a node token T owns the range from the previous node token, exclusive, up to T, inclusive. The primary replica for a partition with token t is therefore the first node token that is greater than or equal to t, wrapping to the lowest token if t is beyond the highest one.
The keyspace's replication strategy then chooses the remaining replicas by walking clockwise from that point. SimpleStrategy takes the next distinct nodes until it has RF. NetworkTopologyStrategy, which production keyspaces should use, walks separately per data centre and prefers racks it has not yet used, with membership supplied by the snitch. The sketch below models the logic; it is not Cassandra's code.
import bisect
MIN_TOKEN = -2**63 # sentinel: the start of the ring, never a real partition token
def primary_index(ring_tokens, token):
"""ring_tokens is sorted. The owner is the first node token >= token, wrapping to 0."""
i = bisect.bisect_left(ring_tokens, token)
return i % len(ring_tokens)
def simple_strategy(ring, token, rf):
"""ring: sorted list of (node_token, node). Walk clockwise collecting distinct nodes."""
tokens = [t for t, _ in ring]
distinct = len({n for _, n in ring})
if rf > distinct:
raise ValueError("rf exceeds the number of nodes")
i, out = primary_index(tokens, token), []
while len(out) < rf:
node = ring[i % len(ring)][1]
if node not in out: # with vnodes the same node appears many times
out.append(node)
i += 1
return out
ring = [(-4611686018427387904, "N1"), (0, "N2"),
(4611686018427387904, "N3"), (9223372036854775807, "N4")]
print(simple_strategy(ring, 1523456789012345678, 3)) # ['N3', 'N4', 'N1']
print(simple_strategy(ring, -9000000000000000000, 3)) # ['N1', 'N2', 'N3']With vnodes each node holds many tokens, so the same node appears many times around the ring and the walk skips nodes it has already chosen. How many tokens to use, and how they are allocated, is covered in the vnodes article.
Worked example: tracing one partition
Take a sensor platform that stores readings per sensor per day, so the partition key is the pair (sensor_id, day). The table and the two commands that trace a partition to its replicas look like this:
CREATE TABLE iot.readings (
sensor_id text,
day date,
ts timestamp,
value double,
PRIMARY KEY ((sensor_id, day), ts)
) WITH CLUSTERING ORDER BY (ts DESC);
-- Which token does this partition hash to?
SELECT token(sensor_id, day) FROM iot.readings
WHERE sensor_id = 'sensor-17' AND day = '2026-10-01' LIMIT 1;
-- Which nodes hold it? (run on any node; composite key parts are joined with ':')
-- $ nodetool getendpoints iot readings sensor-17:2026-10-01Both answers are deterministic and the same on every node. Suppose, for illustration, the token comes back as 1,523,456,789,012,345,678 on the four-node, single-token ring from the sketch above. The first node token at or above it is N3's, at 2^62, so N3 is the primary. With SimpleStrategy and RF 3 the walk continues to N4 and wraps to N1. The partition for the same sensor on the next day has a completely different token, and very likely a different replica set.
That is the key design consequence. Adding day to the partition key spreads one sensor's history across the cluster and keeps each partition bounded, but a query for a week of data must read seven partitions from up to seven replica sets.
Routing: coordinators and token-aware drivers
Any node can accept any request. The node that receives it becomes the coordinator: it computes the token, finds the replicas, sends the request to as many of them as the consistency level requires, and waits. If the coordinator is not itself a replica, every request costs an extra network hop and extra coordinator CPU.
Drivers avoid that by caching the token map and running the partitioner themselves, so a token-aware policy sends each request to a replica. It only works when the driver knows the partition key bytes, which means prepared statements with the partition key bound: preparing returns the key's position among the bound values. A plain string query carries no routing key, so the driver falls back to its child policy.
from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT
from cassandra.policies import TokenAwarePolicy, DCAwareRoundRobinPolicy
profile = ExecutionProfile(
load_balancing_policy=TokenAwarePolicy(DCAwareRoundRobinPolicy(local_dc="dc1")))
session = Cluster(["10.0.0.11"], execution_profiles={EXEC_PROFILE_DEFAULT: profile}).connect()
# Prepared: the driver knows which bound values form the partition key, computes the token
# itself and sends the request straight to a replica in dc1.
insert = session.prepare(
"INSERT INTO iot.readings (sensor_id, day, ts, value) VALUES (?, ?, ?, ?)")
session.execute(insert, ("sensor-17", day, ts, 21.4))Check your driver's configuration rather than assuming token awareness is on. Note also that a batch across many partitions has no single routing key, so the coordinator fans it out, and that token awareness makes a hot partition a hot coordinator for exactly RF nodes.
Scanning a table by token range
Because the partitioner hashes keys, there is no efficient way to ask for a range of partition keys. You can, however, ask for a range of tokens, and every partition lies in exactly one token range. That is how analytics connectors, export tools and custom jobs read a whole table without an unbounded scan: they cut the ring into many ranges, query each with a token() predicate, and page within each range.
MIN_T, MAX_T = -2**63, 2**63 - 1
def splits(n):
"""Cut the whole ring into n contiguous (start, end] ranges."""
step = (MAX_T - MIN_T) // n
bounds = [MIN_T + i * step for i in range(n)] + [MAX_T]
return list(zip(bounds[:-1], bounds[1:]))
scan = session.prepare(
"SELECT sensor_id, day, ts, value FROM iot.readings "
"WHERE token(sensor_id, day) > ? AND token(sensor_id, day) <= ?")
scan.fetch_size = 1000 # page within each range
for start, end in splits(256): # run these concurrently, a few per node
for row in session.execute(scan, (start, end)):
handle(row)The ranges are half-open, greater-than the start and less-than-or-equal-to the end. The first range starts at the sentinel -2^63, which no partition can have, so a strict greater-than loses nothing. For locality, align splits with the token ranges nodes actually own, which driver metadata exposes. Keep concurrency to a few ranges per node, and remember that a scan at LOCAL_ONE is cheap but can miss writes that have not reached the replica it reads.
What partition size does to every path
A large partition is harmless to hashing and costly almost everywhere else:
- Reads. SSTables index where each partition starts and, for wide partitions, sample rows within it, so a slice query can skip to the right block. The index of a very large partition must still be read, and a whole-partition query reads every fragment from every SSTable.
- Compaction. A huge partition makes one long, heavy merge that competes with foreground work. Cassandra logs a warning when it writes a partition above a configurable size; the setting has been renamed between versions, so grep the logs for large-partition warnings.
- Repair and streaming. Repair streams data for token ranges whose Merkle tree hashes differ, so a mismatch inside a huge partition can stream all of it, and bootstrap moves it as one piece.
- Tombstones. Deletes inside a partition leave tombstones that reads must scan past until they are purged, so queue-like access patterns inside one partition degrade over time. The mechanics are in Cassandra tombstones.
- Heap and GC. Large partitions allocate heavily on the JVM heap, and the GC pauses slow every request on the node.
Measure rather than guess. nodetool tablehistograms keyspace table reports partition size and cell count percentiles, and nodetool tablestats reports the largest compacted partition. A common working guideline is to keep partitions under about a hundred megabytes and well under a hundred thousand rows, but treat that as a warning line to design against, not a guarantee.
Hot partitions: load, not size
A small partition can still take down its replicas. If one key, a celebrity user or a tenant configuration row, receives a large share of traffic, all of it lands on RF nodes, and adding nodes does not help. The symptom is a few nodes with high CPU, pending tasks and latency while the rest idle.
The fixes split the key or keep traffic away from it. Write sharding adds a shard number, such as (item_id, shard) with shard chosen at random from 0 to 15, and readers merge all shards. Time bucketing bounds how long a key stays hot. Application caching absorbs read-heavy keys. And a single global sequence number simply does not belong in a partitioned store.
Failure modes and trade-offs
- Unbounded partitions. A partition key that grows forever, such as a user ID for an activity feed, is fine for months and then fails slowly. Add a time bucket before launch, not after.
- Low-cardinality keys. A key like
countryorstatuscreates a handful of enormous partitions, and no amount of cluster growth spreads them. - Unprepared statements. String queries lose token awareness and add a hop to every request, quietly raising latency and coordinator load.
- Multi-partition queries. Large
INlists make one coordinator wait on many replica sets. Prefer concurrent single-partition requests. - The trade-off itself. Small partitions spread load and keep every path cheap, but make range queries touch many replicas. Large partitions make range queries local but concentrate size and load. Good designs choose the partition as the unit a query reads in one request, bounded in both size and traffic.
What to do next
- For each table, write down the partition key and estimate its largest and busiest partition from real traffic.
- Run
nodetool tablehistogramsandnodetool tablestatson a production node and find the largest partitions you already have. - Pick one known key and trace it with
SELECT token(...)andnodetool getendpointsso the mapping stops being abstract. - Confirm that every hot path uses prepared statements with the partition key bound, and that your driver's load-balancing policy is token-aware and data-centre aware.
- Add a time bucket or shard suffix to any key that grows without limit or takes a disproportionate share of traffic.
- Replace full-table queries in batch jobs with token-range scans that page within each range, and check the consistency level they read at.