Advanced System design concept · Distributed Systems & Data · 45 mins read
Data Partitioning & Replication
Split data across machines so it scales, and copy it so it survives failures — and understand what each choice costs.
Partitioning Strategies
Range, hash and directory partitioning: how each maps keys to partitions, and what it does to queries and rebalancing.
Intuition
Once data is split across machines, every query must first work out which machine to ask. If the split matches the query pattern, a request touches one partition. If it does not, every query fans out to every machine and the cluster is slower than a single server.
The partitioning scheme is one of the hardest decisions to change later, because changing it means moving all the data. Interviewers expect you to pick a strategy and a key and explain why.
Mental Model
Range partitioning gives each partition a contiguous key range, like A–F, G–M. Hash partitioning gives each partition a range of hash values, which scatters neighbouring keys. Directory partitioning keeps an explicit lookup table from key to partition. Whatever the scheme, the partition key should be the field most queries filter on. Think of it like: Range partitioning is a dictionary split into volumes by letter — easy to find a run of words, but the S volume is huge. Hash partitioning deals words into volumes by a shuffle — even sizes, but neighbouring words end up in different volumes.
Building Blocks
- Partition key: The field that decides where a row lives, such as user_id or tenant_id.
- Range partitioning: Each partition owns a contiguous key range; efficient for range scans, risky for sequential keys.
- Hash partitioning: Each partition owns a range of hash(key); even load, but range queries must ask every partition.
- Directory (lookup) partitioning: A mapping service records where each key or tenant lives; flexible, but the directory must be fast and highly available.
- Compound keys: Hash on the first part, sort on the second — for example (user_id, timestamp) keeps one user's events together and ordered.
Definitions
- Partition (shard)
-
A subset of the data stored and served together.
- A machine usually hosts many partitions.
- Scatter-gather query
-
A query sent to every partition with results merged afterwards.
- Its latency is set by the slowest partition.
- Local vs global secondary index
-
A local index covers only its own partition, so index queries must ask every partition; a global index is itself partitioned by the indexed field, so reads hit one partition but writes update another.
- Global indexes are usually updated asynchronously.
Patterns
- Hash on entity, sort within it — When queries fetch one entity's items in order.
- Tenant-based directory — Multi-tenant SaaS where tenants vary hugely in size.
Strategies
- Choose the key from the access pattern When: Before creating any partitioned table. How: List the top queries; pick a key that most of them filter on, so they hit one partition, and that has many distinct, evenly used values. Example: Orders are partitioned by customer_id because 'show my orders' is the top query.
- Split partitions dynamically When: When data volume per key range is hard to predict. How: Start with a few partitions and split one into two when it passes a size or load threshold, moving half to another node. Example: HBase regions and DynamoDB partitions split automatically as they grow.
Why secondary indexes are hard on partitioned data
Partitioning by user_id makes 'get this user's orders' a single-partition read. But 'find all orders with status = shipped' does not include the partition key, so with local indexes every partition must be searched and the answers merged.
A global secondary index partitions the index by status instead, so the read hits one index partition. The cost moves to writes: updating one order now touches its data partition and a different index partition, which usually means the index is updated asynchronously and can lag behind. DynamoDB's global secondary indexes are eventually consistent for exactly this reason.
Tradeoffs
| Decision | Upside | Downside |
|---|---|---|
| Range vs hash | Range supports efficient scans and ordered reads; hash spreads writes evenly. | Range concentrates sequential writes on one partition; hash turns range queries into scatter-gather. |
| Directory vs computed placement | A directory can move any key or tenant anywhere. | The directory is extra infrastructure on the request path and must never be the bottleneck. |
Real World
| System | How it's used |
|---|---|
| Amazon DynamoDB | Hashes the partition key to choose a partition and sorts items by the sort key inside it. |
| HBase and Bigtable | Range-partition rows into tablets or regions that split automatically as they grow. |
Interview
Questions interviewers ask
- How would you partition a table of chat messages?
- What is the downside of hash partitioning?
- How do secondary indexes work across partitions?
What a strong answer covers
Justify a partition key from the main queries, compare range and hash, and explain local vs global indexes.
Common traps
- Partitioning by timestamp alone, which sends all new writes to one partition.
- Ignoring scatter-gather cost for queries that do not include the partition key.
Quiz
Which strategy makes range scans over the key cheapest?
- Hash partitioning
- Range partitioning
- Random placement
- Round-robin
Range partitioning keeps neighbouring keys together, so a scan touches few partitions.
Why is partitioning only by timestamp risky?
- Timestamps cannot be hashed
- All new writes go to the latest partition
- It breaks indexes
- It needs a directory
Sequential keys concentrate current writes on one partition, creating a hot spot.
A query does not include the partition key. What usually happens?
- It fails
- It is sent to every partition and the results merged
- It uses the directory
- It hits the leader only
Without the key the system cannot know which partition holds the data, so it scatters and gathers.
What does a global secondary index trade off?
- Faster writes for slower reads
- Single-partition index reads for more complex, often asynchronous writes
- Nothing
- Consistency for durability
The index is partitioned by the indexed field, so writes touch another partition and the index often lags.
Key (user_id, created_at) with hash on user_id gives…
- Random order for each user
- A user's rows together, sorted by time
- All users on one partition
- Global time ordering
The hash part picks the partition and the sort part orders rows inside it.
Hot Partitions
When one key or range gets far more traffic than the rest, and how to spread it without breaking queries.
This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.
Unlock the full lessonReplication Topologies
Single-leader, multi-leader and leaderless replication: who accepts writes, and how conflicts are prevented or resolved.
This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.
Unlock the full lessonSynchronous vs Asynchronous Replication
When a write is acknowledged, what can be lost on failover, and the read anomalies caused by replication lag.
This section is part of the full PRISM roadmap, with worked examples, trade-off tables, interview questions and a quiz.
Unlock the full lessonPractice data partitioning & replication in PRISM
Concepts stick when you watch them fail. Build an architecture that depends on data partitioning & replication, push traffic through it in the PRISM simulator, and see the latency and error rates change as you adjust the design.