Data modelling for Cassandra and DynamoDB
How to design tables for wide-column and key-value stores: query-first modelling, partition keys and sort keys, one table per query versus single-table design, time-bucketed partitions, secondary indexes, avoiding hot and huge partitions, and consistency options.
Reading is half of it. See this used in a real interview: walk through Design WhatsApp →
Cassandra, ScyllaDB, DynamoDB, Bigtable and HBase scale writes and storage almost without limit, which is why they appear in chat, feed, metrics and URL shortener designs. The price is that you cannot query them freely: there are no joins and limited filtering, and every efficient query must name a partition. Modelling for them works backwards from the queries. Getting the partition key right is the difference between a design that scales and one that falls over on its first hot key.
The model
- A table's primary key has two parts:
- the partition key, hashed to choose which nodes store the row; - the clustering or sort key, which orders rows inside the partition.
- Efficient reads: get one row by full key, or range-scan one partition by sort key ("the latest 50 messages in conversation 42").
- Inefficient or impossible: queries that do not specify the partition key, joins, and arbitrary filters.
In DynamoDB terms: partition key and sort key. In Cassandra: PRIMARY KEY ((partition columns), clustering columns).
Query-first design
- List the access patterns with their frequencies: "get messages for a conversation, newest first", "get conversations for a user ordered by last activity".
- For each, design a table (or an item type) whose partition key is what the query looks up by and whose sort key is the order it needs.
- Denormalise: store the same data in several tables shaped for different queries, written together.
Example for chat:
| Table | Partition key | Sort key | Serves |
|---|---|---|---|
| messages | conversation_id | message_id (time-ordered) | messages in a conversation, paged |
| user_conversations | user_id | last_activity, conversation_id | a user's inbox |
| conversations | conversation_id | (none) | conversation metadata |
See Design WhatsApp and data modelling and denormalisation.
Partition sizing and time buckets
Partitions must stay bounded. A conversation with millions of messages, or a sensor writing every second forever, creates a partition that grows without limit: slow to read, hard to repair, and stuck on a few nodes. Fix with bucketing: add a time bucket to the partition key.
- Metrics: partition by
(metric_id, day), sort by timestamp. - Busy chats: partition by
(conversation_id, month).
Queries then read the current bucket and walk back to older ones as needed. Aim for partitions in the tens to low hundreds of megabytes at most. See time-series data and Design a Monitoring System.
Hot partitions
Each partition is served by a limited set of nodes with limited throughput (DynamoDB enforces per-partition limits). A celebrity's timeline, a global counter or a date-only partition key concentrates traffic. Use high-cardinality keys, add a random suffix for write-heavy keys, and cache hot reads. See hot keys and skew.
Secondary indexes
- DynamoDB global secondary indexes are asynchronously maintained copies with a different key; great for alternative access patterns, eventually consistent.
- Cassandra secondary indexes are local to each node and query every node unless combined with the partition key; usually prefer a separate denormalised table.
- Materialised views exist but have operational caveats; many teams maintain their own tables.
Single-table design (DynamoDB)
A DynamoDB style where many entity types share one table with generic keys (PK = USER#42, SK = ORDER#2026-10-05#981), so related items live in one partition and one query fetches a user and their recent orders. It minimises requests and suits stable, well-known access patterns, at the cost of readability and flexibility. Multiple tables are fine when patterns are simpler or still changing.
Writes, updates and deletes
- Writes are upserts and very fast (LSM storage). See storage engines.
- Deletes create tombstones; many deletes in one partition (queues modelled as tables) slow reads badly. Avoid queue-like patterns; use TTLs for expiry.
- Lightweight transactions (Cassandra) and conditional writes or transactions (DynamoDB) exist for compare-and-set, at higher cost.
Consistency
- Cassandra offers tunable consistency per query:
QUORUMreads and writes (with replication factor 3) give read-your-writes;ONEis faster and weaker. UseLOCAL_QUORUMin multi-region setups. See replication and consistency. - DynamoDB reads are eventually consistent by default, with strongly consistent reads available on the base table.
When not to use them
Ad-hoc queries, complex relationships, multi-row transactions across entities and frequently changing access patterns fit a relational database better. Many systems combine both: relational for core entities and money, wide-column for high-volume append-heavy data. See SQL vs NoSQL.
Checklist
- Access patterns listed first; a table or item type per pattern.
- Partition key for lookup, sort key for order.
- Time buckets to bound partition size.
- High-cardinality keys; hot keys mitigated.
- Denormalised tables over secondary indexes for important queries.
- No queue-like delete patterns; TTLs for expiry.
- Consistency level chosen per query.