Key Technology: Cassandra, DynamoDB and Wide-Column Stores
When a design needs very high write throughput, horizontal scale without a single primary, and predictable low latency for known access patterns, the answer is often a wide-column or partitioned key-value store. Apache Cassandra and Amazon DynamoDB are the two names interviewers use. They differ in operation and details, and they share a design philosophy that you can reason about from first principles: partition by key, replicate across nodes, and model the data around the queries.
Product limits and defaults change, so treat specific numbers as things to confirm in current documentation. The principles here are stable.
1. The model
Data is organised into tables, but the table is not a general relational table. Each row is addressed by a primary key with two parts:
- The partition key decides which node (or nodes) store the row. It is hashed, and the hash picks a position in the cluster.
- The clustering key (called the sort key in DynamoDB) orders rows within a partition on disk.
So all rows sharing a partition key live together, sorted by the clustering key. That layout makes one kind of query extremely fast: fetch a slice of one partition, such as the latest 50 messages in a conversation, in a single sequential read.
What these stores deliberately do not give you:
- No joins. Related data is duplicated into tables shaped for each query.
- Limited ad hoc queries. You query by the partition key, and a range or order on the clustering key. Anything else needs another table or index.
- Limited multi-row transactions. Operations are cheap within a partition, and costly or restricted across partitions.
In exchange you get scale with predictable performance: a lookup by key costs the same at ten gigabytes and at ten terabytes.
2. Data modelling: start from the query
This is the idea that distinguishes a prepared candidate. In a relational design you model the entities and then write queries. Here you list the queries first, and design one table per query.
<!--fig:model-->Worked example: a chat application. The main query is "the latest 50 messages in a conversation". Design the table with the conversation identifier as the partition key and the message sequence number as the clustering key, sorted descending. The query is one partition and one sequential read. A second query, "all conversations for a user", is a different access pattern, so it gets its own table keyed by user. When a message is sent, the application writes to both tables. The duplication is intentional: writes are cheap, and reads become single-partition lookups.
Choosing the partition key. It must:
- appear in the query, so reads hit one partition;
- have high cardinality, so data and load spread across nodes;
- keep each partition a bounded size. An unbounded partition, such as all events ever for one device, grows without limit and becomes slow and hard to manage. Add a time bucket to the key (device and month) to cap partition size.
Hot partitions. If one key receives far more traffic than others, the node that owns it becomes the bottleneck, whatever the cluster size. Examples are a celebrity account or a global counter. Fix it by changing the key, adding a salt that spreads one logical key across several partitions (and merging at read time), or caching the hot item.
3. How data is distributed and replicated (Cassandra)
Nodes sit on a token ring (see the consistent hashing discussion in the storage chapter). A row's partition key hashes to a position, and the node owning that position stores it. Each node owns many ranges (virtual nodes) to even the distribution. With replication factor 3, the row is also stored on the next two nodes clockwise, ideally in different racks or zones.
<!--fig:ring-->There is no leader. Any replica can accept a write. A client connects to any node, which acts as the coordinator for that request and forwards it to the replicas.
Tunable consistency. For each request you choose how many replicas must respond:
| Level | Meaning |
|---|---|
| One | One replica responds. Fastest, may read stale data |
| Quorum | A majority of replicas (2 of 3). Strong when used for both reads and writes |
| All | Every replica. Strongest, and fails if any is down |
| Local quorum | A majority within the local data centre, avoiding cross-region latency |
The overlap rule from the consistency chapter applies: if the replicas written plus the replicas read exceed the replication factor, a read sees the latest write. With replication factor 3, quorum writes and quorum reads (2 + 2 > 3) give that property while tolerating one node down.
4. Write path, and why writes are fast
Cassandra uses a log-structured merge tree:
- A write is appended to a commit log (sequential, durable) and applied to an in-memory memtable.
- When the memtable fills, it is flushed to an immutable sorted file on disk.
- In the background, compaction merges files and discards overwritten data.
Writes are sequential and in memory, so they are very fast and do not need a read first. Reads may consult several files, which is why Bloom filters are used to skip files that cannot hold the key.
Consequences to mention:
- Deletes are writes. A delete writes a marker called a tombstone. Reading across many tombstones is slow, so deleting heavily and then scanning is an anti-pattern. Prefer a time-to-live and time-bucketed partitions you can drop whole.
- Reads are more expensive than writes, and a read can touch several files.
- Updates are not in place. The latest version by timestamp wins, so concurrent writers use "last write wins", and clock differences matter.
5. Failure handling and repair
Because any replica can take a write, the system has to heal differences that arise when a replica is down or slow.
- Hinted handoff. If a replica is down, the coordinator stores a hint and replays it when the replica returns, for a limited time.
- Read repair. A read that finds replicas disagree updates the stale ones with the newest value.
- Anti-entropy repair. A periodic, background comparison of replicas (using hash trees of the data) fixes divergence that hints and read repair missed. It must be run regularly, or deleted data can come back and replicas can drift.
- Gossip. Nodes exchange state with each other to learn who is alive.
This is an eventually consistent, highly available design. In terms of CAP, it favours availability under partition, and you can strengthen consistency per request by raising the consistency level, paying with latency and availability.
6. DynamoDB specifics
DynamoDB is a managed service with the same partition-key and sort-key model. The differences you should know:
- Managed scale. You do not run nodes. You choose a capacity mode: provisioned throughput, or on-demand, which scales with traffic. Throughput is spread across partitions, so a hot key can be throttled even when the table as a whole has spare capacity.
- Read consistency choices. Reads are eventually consistent by default, and you can request a strongly consistent read from the leader for a single partition, at higher cost.
- Secondary indexes. A global secondary index lets you query by a different partition key and sort key. It is maintained asynchronously, so it is eventually consistent, and it has its own capacity. A local secondary index shares the partition key and offers a different sort key, and is subject to size constraints per partition.
- Item size limit. Items are limited in size (hundreds of kilobytes). Store large objects in object storage and keep a reference.
- Transactions. It supports transactions across items, at a higher cost than single-item operations. Use them sparingly.
- Streams and TTL. A change stream feeds other systems, and a time-to-live attribute expires items automatically.
- Single-table design. A common advanced pattern stores several entity types in one table with generic key names, so that related items sit in one partition and can be fetched together. It trades readability for efficiency. In an interview, mention it as an option and explain the access-pattern reasoning.
Say that you would check the current limits and pricing model in the provider documentation, since both change.
7. When to choose it, and when not to
Choose a wide-column or partitioned key-value store when:
- the access patterns are known and few, and mostly "get by key" or "range within a partition";
- you need very high write throughput, such as messages, events, sensor readings and activity feeds;
- you need linear horizontal scale and high availability across zones or regions;
- the data is time-ordered, and old data can expire.
Avoid it when:
- you need ad hoc queries, joins or analytics over the same data. Pair it with a search engine or a warehouse fed from it;
- you need multi-row transactions with strong guarantees. Use a relational database;
- the data is small and the queries are varied, where a relational database is simpler and cheaper;
- the team cannot afford the modelling discipline, or the operational burden of running Cassandra themselves, when a managed alternative would serve.
8. Interview questions and model answers
Q: How do you choose a partition key? From the dominant query: it must appear in the query so reads touch one partition, have high cardinality so load spreads, and keep partitions bounded in size. I add a time bucket when a partition would otherwise grow forever.
Q: What is a hot partition and how do you fix it? One key receiving disproportionate traffic, so the node owning it saturates. I change the key, spread one logical key over several partitions with a salt and merge on read, or cache the hot item.
Q: How does Cassandra achieve high availability? Data is replicated to several nodes, any replica can accept writes, and the consistency level is chosen per request. Hinted handoff, read repair and periodic repair converge replicas after failures.
Q: What do quorum reads and writes give you? With replication factor 3, writing to 2 and reading from 2 means every read overlaps the latest write, so reads see it, while one node can be down.
Q: Why are Cassandra writes fast? They are appended to a commit log and applied in memory, then flushed as sorted files and merged later. There is no read before write and no random disk access.
Q: When would you not use it? For ad hoc queries, joins and multi-row transactions. A relational database fits those, and a search or analytics store can be fed from the wide-column store for flexible queries.
9. Common mistakes
- Designing tables like a relational schema and then wanting joins.
- A partition key with low cardinality, or one that is not in the query.
- Unbounded partitions that grow forever.
- Heavy deletes followed by scans, leaving many tombstones.
- Using the highest consistency level everywhere and losing the availability benefit.
- Skipping regular repair.
- Treating a secondary index as free, when it is a second copy with its own cost and lag.