In Cassandra the partition key is the most consequential decision you make about a table, and the only one you cannot change later with an ALTER statement. It decides which nodes store each row, which queries are cheap, which are impossible, how evenly load spreads across the cluster, and how large each unit of storage grows over years. Most production Cassandra problems that look like operations problems, such as a node pinned at 100 percent CPU, slow reads on one table, or compactions that never finish, trace back to a partition key chosen in an afternoon.

This article is a method for choosing one. It assumes you know what a partition is; if not, start with how Cassandra partitions data. Here you will learn the four properties every candidate key must be scored on, the three families of keys, a small script that simulates distribution before you create anything, how to measure a live table, and how to migrate when the key you have is wrong.

Architecture at a glance

Scoring a candidate partition keyAccess patternsevery query, written downCandidate keysnatural, composite, shardedSimulatesizes and skew from sample data1. Query alignmentone partition per read2. Cardinalitymany more keys than nodes3. Skewhottest key vs median4. Growthsize bounded over timePick the key that passes all fourthen add bucketing only if growth failsWrong key in production?new table, dual write, backfill, cut over
Every candidate key is scored on four properties against real access patterns and sample data; bucketing and sharding are added only when a property fails, and a wrong key is fixed by migration, not ALTER.

From first principles: what the partition key decides

A Cassandra primary key has two parts. The partition key is hashed (by the default Murmur3 partitioner) to a token, and the token decides which replicas own the row. The clustering columns sort rows inside the partition. Everything else follows from one fact: a read is efficient when it touches one partition and a contiguous slice of its clustering order. A query that names the full partition key goes to exactly the replicas that hold it. A query that does not must either be rejected or scan every node.

So the partition key is simultaneously your unit of locality (rows read together should live together), your unit of distribution (different keys spread across nodes), and your unit of growth (all rows for a key accumulate in one place forever, unless TTL or deletes remove them). Those three roles pull in opposite directions. A key coarse enough to answer a query in one partition may be so coarse that one partition holds a million rows; a key fine enough to spread load may scatter the rows a query needs.

The four properties of a good key

Score every candidate on four properties, in this order. A key that fails the first is wrong no matter how well it scores on the rest.

PropertyQuestion to askHow to measureRed flag
Query alignmentDoes every frequent query supply the full partition key?Write each query next to the keyNeeds ALLOW FILTERING, IN over many keys, or a scan
CardinalityHow many distinct key values will exist?Count distinct values in a data sampleFewer than about 100 values per node, e.g. status or country
SkewHow much bigger and busier is the hottest key than the median?p99 and max rows per key, plus request shareOne customer, tenant or date carries a large share of traffic
GrowthIs the partition size bounded as time passes?Rows per key per day times retentionUnbounded keys such as sensor_id with no time component

There is no hard maximum partition size, but long-standing community rules of thumb are to keep partitions under roughly 100 MB and well under a few hundred thousand rows. Larger partitions still work, and recent versions handle them better than older ones, but they make compaction, repair, streaming and reads of the partition slower, put heap pressure on the nodes that own them, and concentrate load. Treat the numbers as a design budget, not a database limit.

Natural, composite and synthetic shard keys

Candidate keys come in three families. Most good designs start with the first and move down only when a property fails.

Natural keys use an identifier the application already has: user_id, order_id, device_id. They align with lookups by that identifier and usually have excellent cardinality. They fail on growth when the entity accumulates rows forever, and on skew when a few entities are enormous, like a marketplace's largest seller.

Composite keys combine an identifier with a bounding component, usually time: PRIMARY KEY ((device_id, day), ts). The extra component caps growth, because each partition holds one day, and it raises cardinality. The cost is that the query must know the bucket, so a read across a week becomes seven partition reads issued in parallel. Pick the bucket size from arithmetic, so a busy key fills a bucket to a fraction of the size budget; time-series modeling walks through that calculation.

Synthetic shard keys add a component that exists only to spread load: PRIMARY KEY ((tenant_id, shard), event_id) where shard = hash(event_id) % 8. They fix skew for very hot keys at the price of fan-out on read: to read a tenant you query all eight shards and merge. Use them only for the keys that need them, and keep the shard count fixed or derivable, because a reader must know every shard that might hold data.

Simulating a key before you create the table

You can score cardinality, skew and growth before creating a table by replaying a realistic sample through each candidate key. The script below counts rows per key, projects growth for keys without a time bound, and reports the largest partition and the skew ratio. It is deliberately simple; its job is to catch an obviously bad key in minutes rather than months.

# Simulate partition sizes and skew for candidate keys from a data sample.
# Input: a CSV export or log sample with one row per future Cassandra row.
import csv, collections, statistics, hashlib

ROW_BYTES = 220            # measured average row size incl. overhead (estimate yours)
RETENTION_DAYS = 365
SAMPLE_DAYS = 7            # how many days the sample covers

CANDIDATES = {
    "tenant":            lambda r: (r["tenant_id"],),
    "tenant_day":        lambda r: (r["tenant_id"], r["ts"][:10]),
    "tenant_shard8":     lambda r: (r["tenant_id"],
                                    int(hashlib.md5(r["event_id"].encode()).hexdigest(), 16) % 8),
}

def score(rows, keyfn, time_bounded):
    counts = collections.Counter(keyfn(r) for r in rows)
    sizes = sorted(counts.values())
    scale = 1 if time_bounded else RETENTION_DAYS / SAMPLE_DAYS   # projected growth
    p50 = statistics.median(sizes) * scale
    mx = sizes[-1] * scale
    return {
        "partitions": len(sizes),
        "p50_rows": int(p50),
        "max_rows": int(mx),
        "max_mb": round(mx * ROW_BYTES / 1e6, 1),
        "skew_max_over_p50": round(mx / max(p50, 1), 1),
    }

rows = list(csv.DictReader(open("events_sample.csv")))
for name, fn in CANDIDATES.items():
    print(name, score(rows, fn, time_bounded=name.endswith("_day")))

Two cautions. First, the sample must include your largest customers and busiest periods, or skew will look better than it is. Second, row size must come from measurement (for example, write a few thousand rows to a test table and divide the on-disk size by the row count), because rows carry per-cell overhead and estimates from column types are usually low.

Worked example: a multi-tenant audit log

A SaaS product stores an audit log: every user action, per tenant, kept for a year. The query that matters is "show the latest events for tenant T, newest first, paged", plus a compliance export of one tenant for a date range. There are about 20,000 tenants. A seven-day sample shows a median tenant writing about 300 events a day, but the largest tenant writes about 1.5 million a day, and the top 10 tenants produce a third of all writes. Rows average about 400 bytes.

Candidate A, (tenant_id). Alignment is perfect. Growth fails badly: the largest tenant would accumulate about 550 million rows a year, which is over 200 GB in one partition on three replicas. Skew fails too, because every write for that tenant lands on the same replicas.

Candidate B, (tenant_id, day). Alignment holds, since the latest events come from today's partition and a range export iterates days. Growth is bounded: the median partition is about 120 KB, but the largest tenant still produces 1.5 million rows and about 600 MB in a single day partition, and that partition is a write hotspot all day.

Chosen design, (tenant_id, day, shard). Normal tenants always use shard 0, so their reads stay single-partition. Tenants on a configured hot list spread writes over eight shards by hashing the event id, so the largest tenant's day partitions are about 75 MB each and its write load spreads over up to eight replica sets. Reads for hot tenants issue eight queries in parallel and merge by event id, which the application already does for paging. The hot list lives in configuration that readers and writers share; adding a tenant to it is safe at a day boundary, because older days remain on shard 0.

-- Candidate A: fails growth and skew
CREATE TABLE audit.events_by_tenant (
  tenant_id  uuid,
  event_id   timeuuid,
  actor      text,
  action     text,
  payload    text,
  PRIMARY KEY ((tenant_id), event_id)
) WITH CLUSTERING ORDER BY (event_id DESC);

-- Chosen: day bucket for growth, fixed shard count for the few hot tenants
CREATE TABLE audit.events_by_tenant_day (
  tenant_id  uuid,
  day        date,
  shard      tinyint,          -- 0 for normal tenants, 0..7 for tenants on the hot list
  event_id   timeuuid,
  actor      text,
  action     text,
  payload    text,
  PRIMARY KEY ((tenant_id, day, shard), event_id)
) WITH CLUSTERING ORDER BY (event_id DESC)
  AND default_time_to_live = 31536000;   -- 365 days

-- Read the last day for a hot tenant: 8 parallel single-partition queries
SELECT * FROM audit.events_by_tenant_day
 WHERE tenant_id = ? AND day = ? AND shard = ? LIMIT 100;

Measuring partitions in a live cluster

After launch, measure rather than trust the simulation. Run these on several nodes, because each reports only the partitions it owns:

# Partition size and cell count percentiles for one table, on one node
nodetool tablehistograms audit events_by_tenant_day

# Per-table summary, including compacted partition minimum, mean and maximum bytes
nodetool tablestats audit.events_by_tenant_day

# Sample the hottest partitions by reads and writes for 10 seconds (duration in ms)
nodetool toppartitions audit events_by_tenant_day 10000

Read the histograms for the gap between the 99th percentile and the maximum partition size. A maximum many times the p99 means a handful of outlier keys, which is a skew problem, not a bucketing problem. Also watch the system log for the large-partition warnings that compaction emits once a partition exceeds the configured threshold, and compare per-node read and write rates for the table: one node far busier than its peers for one table is the signature of a hot key.

Migrating to a new partition key

A partition key cannot be altered, because changing it changes where every row lives. Migration is a data copy, done online in five steps.

  1. Create the new table with the new key alongside the old one.
  2. Dual write. Change the application to write every new row to both tables. Writes to the new table must be idempotent, so the same row written twice is harmless; Cassandra upserts make this natural if the full primary key is deterministic.
  3. Backfill. Copy historical rows with a token-range scan of the old table, using Spark with the Cassandra connector, DSBulk unload and load, or a custom job that pages through token ranges. Throttle it, because it competes with production traffic.
  4. Verify. Compare row counts per sampled key between tables, and shadow-read: serve from the old table, read from the new one too, and log mismatches.
  5. Cut over and clean up. Switch reads to the new table, keep dual writes for a rollback window, then stop writing the old table and drop it after a snapshot.

Order matters: start dual writes before the backfill begins, or rows written during the backfill are missed. If the old table has TTLs, carry the remaining TTL across rather than resetting it, or expired data comes back to life in the new table. Likewise copy each cell's original write time, read with WRITETIME() and written with USING TIMESTAMP, so a stale backfilled row cannot overwrite a newer dual-written one.

Failure modes

  • Low-cardinality keys. Keying by status, country or type creates a few enormous partitions that a handful of nodes must serve.
  • Unbounded growth. A key with no time component grows forever; reads slow down year by year until compaction and repair struggle.
  • Buckets the reader cannot compute. If the bucket depends on data the reader does not have, such as a random shard not stored anywhere, the rows are effectively lost to queries.
  • Time buckets that put all writes on one partition. Keying only by day sends all of today's writes to one replica set. Combine time with an identifier.
  • Designing for one query and patching the rest with secondary indexes. Each extra access pattern usually deserves its own table; see query-first data modeling.

Trade-offs

Coarser keys give cheaper reads and worse distribution; finer keys give better distribution and more fan-out. Time buckets bound growth at the cost of multi-partition range reads. Synthetic shards fix hot keys at the cost of merging results and keeping the shard rule consistent between writers and readers. Denormalising into several tables, one per query, is normal in Cassandra and costs storage and write amplification, which is usually cheaper than a key that serves every query badly. Decide which query must be fastest and let the others pay.

What to do next

  1. List every query the table must serve, with frequency and latency target, before writing any CQL.
  2. Write two or three candidate keys and score each on alignment, cardinality, skew and growth.
  3. Run the simulation script on a sample that includes your largest tenants and peak days, with a measured row size.
  4. Choose the simplest key that passes all four; add a time bucket only if growth fails and shards only for keys that are proven hot.
  5. After launch, schedule a monthly check of nodetool tablehistograms and toppartitions on several nodes, and alert on large-partition warnings.
  6. If a key is already wrong, plan the dual-write migration now, while the backfill is still small.
Key takeaway: The partition key decides locality, distribution and growth at once, and it cannot be altered later. Score every candidate on query alignment first, then cardinality, skew and growth, using a simulation on real sample data with a measured row size. Start with a natural key, add a time bucket only to bound growth and a fixed shard count only for proven hot keys. Measure live tables with nodetool tablehistograms, tablestats and toppartitions, and fix a wrong key with a new table, dual writes, a throttled backfill, verification and a cut-over.