People coming to Cassandra from relational databases usually start by drawing entities and relationships, normalising them into tables and then writing queries. In Cassandra that order produces tables that cannot answer the queries, or answer them by scanning the cluster. Cassandra has no joins, no general secondary access paths that are free, and a storage layout fixed by the primary key. The model has to start from the queries.

This page builds that skill from first principles: what each part of a primary key controls, the query-first method, a complete worked model for order history, arithmetic for keeping partitions a healthy size, how to keep several denormalised tables in step, which WHERE clauses Cassandra will and will not run, and how deletes shape a model. It ends with a checklist you can apply to your own schema.

Advertisement

Why Cassandra modeling is different

Cassandra distributes data by hashing the partition key to a token and placing the partition on the replicas that own that token. A partition is the unit of distribution and the unit of efficient reading: all rows in one partition live together on each replica, sorted by their clustering columns. A query that names one partition goes to one set of replicas and reads one contiguous, sorted region. A query that does not name a partition has to ask every node.

Three consequences follow. First, every efficient query must supply the full partition key. Second, rows come back in clustering order, so sorting has to be designed into the table rather than requested at query time. Third, if two queries need data organised differently, you write the data twice, in two tables. Denormalisation is not a compromise here; it is the normal design. Storage is cheap, writes are cheap, and cross-node reads are expensive.

Anatomy of a primary key

PRIMARY KEY ((customer_id, month), placed_at, order_id): what each part decidesPartition key: (customer_id, month)hashed by Murmur3 to a token: WHERE the data livesClustering: placed_at DESC, order_idsort order on disk: HOW rows are laid out insidetokenToken ring, RF = 3node Anode Bnode Cone partition = one contiguous read on each replicano partition key in WHERE = scan every nodesortedPartition (c42, 2026-09)2026-09-29 18:02 | o-981 | total, status2026-09-14 09:15 | o-944 | total, status2026-09-02 11:40 | o-902 | total, statusQuery: WHERE customer_id = ? AND month = ? LIMIT 20one replica set, one partition, first 20 rows in stored order: no sort, no scanPick the partition key for distribution and bounded size; pick clustering columns for the order your query reads.
The partition key decides placement and the size of each partition; the clustering columns decide on-disk order within it. A well-shaped query reads the first rows of a single partition.

In PRIMARY KEY ((customer_id, month), placed_at, order_id) the inner parentheses are the partition key, here composite, and the rest are clustering columns. The full key identifies a row; writing an existing key overwrites it (an upsert), which is why order_id is included even though two orders rarely share a millisecond.

Partition keys are compared only for equality. Clustering columns support equality and ranges, in declared order: order_id can be restricted only after placed_at is fixed. CLUSTERING ORDER BY (placed_at DESC) stores newest first, so "latest 20" reads 20 rows and stops. static columns are stored once per partition, which suits header data next to detail rows.

Advertisement

The query-first method

  1. List the access patterns. Write each query the application will run, with its parameters, expected result size, sort order and frequency. Screens and API endpoints are a good source.
  2. Design one table per access pattern (or per group of patterns that share a partition and order). The partition key is what the query supplies by equality; clustering columns give the order and any range.
  3. Check partition size for each table using the arithmetic below, and add a bucket to the partition key if a partition can grow without bound.
  4. Decide how each table is written: which event writes which tables, and how they are kept consistent.
  5. Plan deletes and expiry so tombstones do not accumulate where reads scan.

The entity-relationship diagram is still useful as a description of the domain and of which attributes belong together. It just stops being the physical schema.

Worked example: order history

An online shop needs three queries. Q1: show a customer's recent orders, newest first, 20 per page, with total and status. Q2: show one order with its lines. Q3: for operations, list today's orders in a given status, such as all orders still PLACED, newest first. The customer query runs constantly; the operations query runs a few times a minute.

-- Q1: a customer's recent orders, newest first, paged
CREATE TABLE shop.orders_by_customer (
    customer_id uuid,
    month       text,          -- bucket, e.g. '2026-09'
    placed_at   timestamp,
    order_id    uuid,
    total_cents bigint,
    status      text,
    item_count  int,
    PRIMARY KEY ((customer_id, month), placed_at, order_id)
) WITH CLUSTERING ORDER BY (placed_at DESC, order_id ASC);

-- Q2: one order by id, with its lines
CREATE TABLE shop.order_lines_by_order (
    order_id     uuid,
    line_no      int,
    customer_id  uuid static,  -- stored once per partition
    placed_at    timestamp static,
    status       text static,
    sku          text,
    qty          int,
    price_cents  bigint,
    PRIMARY KEY ((order_id), line_no)
);

-- Q3: operations view: orders in a given status for a day, spread over 16 buckets
CREATE TABLE shop.orders_by_status_day (
    status      text,
    day         date,
    bucket      tinyint,       -- hash(order_id) % 16
    placed_at   timestamp,
    order_id    uuid,
    customer_id uuid,
    PRIMARY KEY ((status, day, bucket), placed_at, order_id)
) WITH CLUSTERING ORDER BY (placed_at DESC, order_id ASC);

Q1's month bucket stops a long-lived customer building an unbounded partition; the application reads the current month and steps back when a page runs out. Q2 keeps an order's lines in one partition with order-level fields static. Q3 is the dangerous one: partitioning by status alone would put every PLACED order in one hot partition. The day bounds its size and a 16-way bucket spreads writes; the operations screen reads the 16 partitions in parallel and merges them.

Notice what is duplicated: placed_at, status and customer_id appear in all three tables. When an order ships, all three must change, and Q3 needs the row removed from the PLACED partition and added to SHIPPED. That is the real cost of denormalisation, and it is dealt with below.

Partition size arithmetic and bucketing

Partitions have no hard size limit short of very large values, but large partitions hurt: they are read and compacted as a unit, repair and streaming move them whole, and their index entries cost memory. A common guideline is to keep partitions under roughly 100 MB and a few hundred thousand rows, and much smaller where latency matters.

Estimate rows per partition and bytes per row. For Q1, a heavy customer places at most around 300 orders a month. Each row holds two uuids (16 bytes each), a timestamp (8), a bigint (8), a short status string, an int, plus per-cell and per-row overhead of a few tens of bytes: call it 150 bytes. That is 300 x 150 = 45 KB per partition, comfortably small. Without the month bucket, a ten-year customer would reach 36,000 rows and about 5 MB: still fine, but a business account placing 50,000 orders a month would reach 90 MB in a single month's bucket. If such accounts exist, bucket them by day, or pick the bucket width per account and store it.

For Q3, suppose 200,000 orders a day, all passing through PLACED. With 16 buckets, each partition holds about 12,500 rows of about 120 bytes, or 1.5 MB, and the write load spreads across 16 token ranges. The formula generalises: rows per partition equals events per bucket period divided by the number of buckets, and size equals rows times bytes per row. Choose the bucket so the worst realistic case, not the average, stays within budget.

Keeping several tables in step

A logged batch guarantees that if any statement is applied, all eventually are, even if the coordinator fails, because the batch is first written to a batch log on other nodes. It gives no isolation, and it costs extra writes, so use it for atomicity across a few partitions, never for bulk loading. The trade-offs are covered in Cassandra batches.

Status changes that move a row between partitions, as Q3 requires, are a delete plus an insert. Put both, plus the updates to Q1 and Q2, in one logged batch. Use the same client-supplied timestamps or values on retry so that repeating the batch is harmless: inserts in Cassandra are idempotent upserts, which makes retries safe as long as the values do not change between attempts. Where a transition must happen at most once, for example PLACED to SHIPPED only if still PLACED, a lightweight transaction (UPDATE ... IF status = 'PLACED') gives compare-and-set on one partition at the price of a Paxos round; see lightweight transactions.

Materialized views and indexes can maintain some of these tables for you. Materialized views are marked experimental and are disabled by default in recent versions because views can drift from their base table; Storage-Attached Indexes (SAI), added in Cassandra 5.0, are the better option for filtering within or across partitions at moderate selectivity, as explained in SAI. For the core access paths, explicit tables remain the most predictable choice.

Driver code: prepared statements, batches and paging

Application code should prepare statements once, bind values per request and let the driver route each request to a replica that owns the partition (token-aware routing is the default in the DataStax Python driver). The sketch below writes a new order to all three tables and pages through Q1.

from cassandra.cluster import Cluster
from cassandra.query import BatchStatement, BatchType, ConsistencyLevel

session = Cluster(["10.0.0.11", "10.0.0.12"]).connect("shop")

ins_cust = session.prepare(
    "INSERT INTO orders_by_customer (customer_id, month, placed_at, order_id, total_cents, status, item_count) "
    "VALUES (?, ?, ?, ?, ?, ?, ?)")
ins_line = session.prepare(
    "INSERT INTO order_lines_by_order (order_id, line_no, customer_id, placed_at, status, sku, qty, price_cents) "
    "VALUES (?, ?, ?, ?, ?, ?, ?, ?)")
ins_stat = session.prepare(
    "INSERT INTO orders_by_status_day (status, day, bucket, placed_at, order_id, customer_id) "
    "VALUES (?, ?, ?, ?, ?, ?)")

def place_order(o):
    # Logged batch: all three tables eventually get the order, even if the
    # coordinator dies mid-way. Not isolation: readers may briefly see one table only.
    b = BatchStatement(batch_type=BatchType.LOGGED, consistency_level=ConsistencyLevel.LOCAL_QUORUM)
    month = o.placed_at.strftime("%Y-%m")
    b.add(ins_cust, (o.customer_id, month, o.placed_at, o.order_id, o.total_cents, "PLACED", len(o.lines)))
    for i, ln in enumerate(o.lines):
        b.add(ins_line, (o.order_id, i, o.customer_id, o.placed_at, "PLACED", ln.sku, ln.qty, ln.price_cents))
    b.add(ins_stat, ("PLACED", o.placed_at.date(), o.order_id.int % 16, o.placed_at, o.order_id, o.customer_id))
    session.execute(b)      # same values on retry -> idempotent upserts

recent = session.prepare(
    "SELECT order_id, placed_at, total_cents, status FROM orders_by_customer "
    "WHERE customer_id = ? AND month = ?")
recent.fetch_size = 20      # driver pages; pass paging_state to continue

Use LOCAL_QUORUM for both writes and reads when read-your-writes matters within a data center: with a replication factor of 3, two replicas must acknowledge each, so every read overlaps every acknowledged write. Paging with a fetch size bounds memory on both sides; see Cassandra paging.

What CQL will and will not run

QueryResultWhy
WHERE customer_id = ? AND month = ?runs, one partitionfull partition key by equality
... AND placed_at > ?runs, range slicefirst clustering column, range
... AND order_id = ? without placed_atrejectedclustering columns must be restricted in order
WHERE customer_id = ? alonerejectedpartial partition key cannot be hashed
WHERE status = 'PLACED' on Q1rejected unless ALLOW FILTERINGwould scan every partition on every node
ORDER BY placed_at ASCrunsreverse of stored order is allowed within a partition

ALLOW FILTERING tells Cassandra to read and discard rows. Within a single, small partition that can be acceptable. Without a partition key it is a full-cluster scan whose cost grows with the data, and it is the most common cause of timeouts in young Cassandra applications.

Design for deletes: tombstones

Deletes in Cassandra write a tombstone, a marker that shadows older data until compaction removes both after gc_grace_seconds (ten days by default) and after repair has had a chance to spread the tombstone. Reads must skip over tombstones they encounter, so a model that deletes many rows from the front of a partition and then reads from the front, such as a queue, gets slower with every delete until reads hit the tombstone failure threshold.

Q3 is exposed to this: every order leaves the PLACED partition when it ships. Bucketing by day limits the damage, because each day's partitions stop being written and read after the day ends, and a TTL on the table lets old buckets expire without explicit deletes. Using a time window compaction strategy for time-bucketed, TTL-ed tables lets whole SSTables be dropped at once. Never model a work queue as "insert, read the oldest, delete it" in one partition. The mechanics are explained in Cassandra tombstones.

Anti-patterns and failure modes

  • Unbounded partitions keyed by something that grows forever (a user's events, a sensor's readings) without a time bucket.
  • Low-cardinality partition keys such as status, country or boolean flags, which create a few huge, hot partitions.
  • Relational habits: normalised tables joined in application code with one query per row, multiplying latency.
  • ALLOW FILTERING in production paths, and secondary indexes on high-cardinality columns queried without a partition key.
  • Read-before-write to decide what to write, which doubles latency and races with concurrent writers; prefer upserts or lightweight transactions.
  • Large logged batches used for bulk loading, which overload coordinators and trigger batch size warnings.
  • Queue-shaped tables that delete from the head of a partition and read from it, accumulating tombstones.

What to do next

  1. Write down every query with its parameters, result size, sort order and frequency.
  2. Derive one table per query and check that each query supplies its full partition key by equality.
  3. Compute worst-case rows and bytes per partition, and add a time or hash bucket wherever the worst case exceeds your budget.
  4. List which events write which tables, and wrap multi-table changes in logged batches with stable, retry-safe values.
  5. Search the code for every delete and every ALLOW FILTERING, and redesign the ones on hot paths.
  6. Load-test with production-like key distributions, and watch partition-size and tombstone warnings before launch.
Key takeaway: Cassandra stores each partition together and sorted by its clustering columns, so the model starts from the queries: one table per access pattern, a partition key the query supplies by equality, clustering columns in the order the query reads, and a bucket wherever a partition could grow without bound. Duplicate data freely, keep the copies in step with logged batches and idempotent writes, use lightweight transactions only for true compare-and-set, avoid ALLOW FILTERING outside a single partition, and shape deletes so tombstones expire with their buckets.