Concept lesson · Foundations
Data partitioning and sharding
Start here
Definition
Data partitioning divides a dataset into smaller parts. Sharding is horizontal partitioning across separately managed storage groups: each shard owns different records, while replicas hold copies of the same records.
Why it matters: One database may run out of storage or processing capacity. Dividing ownership lets different groups handle different records, at the cost of routing and operations that cross those groups.
The example uses customerNumber modulo two to choose one owner. A global report still fans out or needs a separate read model.
Read the diagram step by step
- C12 and C44 are even, so their orders O1/O2 and O5/O6 belong to A. C27 is odd, so O3/O4 belong to B.
- C27 last ten orders routes directly to B, then uses a local ordered index.
- All orders today crosses customer owners and needs fanout or an analytical model.
- When moving ownership, copy and replay, fence old writes, and switch a versioned routing epoch.
Worked example
Using customerNumber mod 2, C12’s orders O1/O2 go to shard A and C27’s O3/O4 go to shard B. A request for C27’s history contacts B; a report across all customers needs both owners.
Key takeaways
You will learn to
- Distinguish partitioning from replication and indexing.
- Choose a key from concrete queries and related writes.
- Explain skew, cross-shard work, and an online shard move.
Practice in this chapter
8 interview questions with model answers and follow-ups.
Go to interview practiceUseful foundations: Databases, data models, and ACID transactions · Database indexes: B-trees, composite keys and query access
Workload and timing examples are interview assumptions.
Dotted concept links open the relevant explanation in a new tab.
01Partitioning and sharding: definitions
Data partitioning divides a dataset into smaller parts. Sharding is horizontal partitioning across separately managed storage groups. Horizontal means dividing records rather than splitting fields of each record. Each shard owns a subset of rows or records; ownership means that a designated storage group is responsible for those records and decides which writes it accepts. Replication instead creates copies of the same records. You can shard orders by customer and also replicate each shard; one choice divides ownership and the other protects each owner's data.
Suppose a shop stores one billion orders and its single database cannot meet the required storage or throughput. Before splitting, inspect inefficient queries and unused indexes; sharding adds routing, movement, and cross-shard complexity. When splitting is justified, choose a partition key: a field or combination of fields used to determine ownership.
Our example has customers C12, C27, and C44, and orders O1 through O6. Most screens ask for one customer's recent orders. That query suggests keeping a customer's records together rather than scattering every order randomly.
02Horizontal, vertical, functional and directory partitioning
Dividing a dataset involves two choices: what to separate, and how to find each part. Horizontal, vertical and functional partitioning describe what is separated. A directory describes how requests find the owner of a part.
Horizontal partitioning divides rows of the same kind: some customers' orders on shard A and others on B. Range partitioning, which groups rows by intervals of a key, is one horizontal method, not another name for all horizontal partitioning.
The cells show the same small dataset. In the vertical split, both partitions keep the ID needed to join the fields.
Remember: Horizontal cuts between records; vertical cuts between fields.
Try from memoryWhy does U1 appear in both vertical partitions?
It is the shared identity used to join the fields back into one logical record.
Vertical partitioning divides columns or attributes. A frequently read account profile might be separate from large optional biography data. Functional partitioning separates different responsibilities or datasets, such as orders, catalog, and billing. Both can reduce unnecessary work, but a user operation that needs separated data must combine it somewhere.
Directory-based placement keeps a lookup from a logical group to its physical owner. For example, a directory can map tenant T7, a customer organization sharing the service, to shard B, allowing T7 to move later without changing its identity. The directory becomes important routing metadata: cache it carefully, version it, and make stale routes detectable rather than treating it as an infallible box.
03Worked example: place and query six orders
Now apply horizontal partitioning to the orders example: keep every order for one customer on the same shard, so that customer’s order history can be read locally. The router needs a rule that turns a customer number into a shard destination.
For a two-shard example, use customerNumber mod 2. The modulo operation returns the remainder after division by two: even customers go to A, odd customers to B. This simple function is for demonstrating placement, not the final resharding scheme.
| Record | Customer | Calculation | Shard |
|---|---|---|---|
| O1 | C12 | 12 mod 2 = 0 | A |
| O2 | C12 | 12 mod 2 = 0 | A |
| O3 | C27 | 27 mod 2 = 1 | B |
| O4 | C27 | 27 mod 2 = 1 | B |
| O5 | C44 | 44 mod 2 = 0 | A |
| O6 | C44 | 44 mod 2 = 0 | A |
A query for C27’s last ten orders computes shard B, then uses a local index on (customerId, createdAt, orderId). It contacts one owner. A report over all customers’ orders cannot identify one shard from that key; it requires fanout, meaning subqueries to the relevant shards followed by a merge, or a separate analytical/read model organized for reporting.
- 1 → 21. Read customer C27Request: C27 order history → Router: customerNumber mod 2
- 2 → 42. 27 mod 2 = 1: route to BRouter: customerNumber mod 2 → Shard B: C27 O3/O4
- 2 → 3Other customer: even keys route to ARouter: customerNumber mod 2 → Shard A: C12 O1/O2, C44 O5/O6
- 4 → 5Replication copies B; it does not split BShard B: C27 O3/O4 → Replica of B: same O3/O4
04Choose a shard key and a placement rule
The previous example chose customerNumber as the shard key and used mod 2 as the placement rule. The key supplies the value used for routing; the rule determines its destination. Choosing a rule affects which records stay together and which queries must contact several shards.
Range, hash and list partitioning choose a destination from key values. Round-robin placement cycles through destinations for new records. All four assign whole records to partitions, so they are approaches to horizontal partitioning. Each row below is a separate placement example.
| Placement method | How records are assigned | When it helps, and the cost |
|---|---|---|
| Range | Assign intervals of a key to partitions: customers 1–999 on A, 1000–1999 on B. | Nearby key values stay together for range queries; a popular or growing range can overload one owner. |
| Hash | Apply a hash function, which maps the key to a repeatable numeric value, then map that value to a partition or logical bucket. | Spreads many distinct keys; adjacent original values usually scatter, so range scans contact several owners. |
| List | Explicitly name the key values assigned to each partition, such as selected countries in one group. | Gives direct control over placement; the lists and each group's capacity need maintenance. |
| Round robin | Assign successive new rows to A, then B, then A again. | Spreads insert counts, but a later key lookup needs stored location metadata or a search across partitions. Equal row counts need not mean equal load. |
A composite shard key combines fields, such as (tenantId, customerId); it is a choice of key, not a fifth placement algorithm. A system can apply range or hash placement to that combined key. It can also combine rules in stages: choose a tenant's shard group, then hash the customer ID within that group. This gives control over tenant placement while distributing its customers; routing to one customer needs both dimensions.
A further question is how placement changes when machines are added or removed. Separating logical groups of records from physical servers makes those moves easier to manage.
A logical bucket is a named group of keys independent of a physical server. Use many logical buckets and a versioned bucket-to-machine map when machines must change. Directly applying key mod numberOfMachines changes many assignments when the machine count changes. Consistent hashing is another way to reduce membership-related movement, but still requires actual data migration and hot-key handling.
Also inspect cardinality, the number of distinct key values, and frequency, how often each value occurs. Hashing a two-value status field still leaves only two groups; it does not manufacture independently movable keys. Check whether a key grows monotonically, whether one value dominates bytes or traffic, and whether its value can change. Updating a customer’s shard-key value can require moving its records rather than changing one local field. Prefer a stable key when it fits the access patterns.
05Cross-shard joins, transactions and denormalization
A shard key that makes one query local can separate records needed by another operation. This affects both reading related data (joins) and updating related data together (transactions); the examples below show where extra coordination or a stored copy becomes necessary.
Suppose O3 and its order items share C27's partition. A local transaction can update them together on B. If an order also changes globally shared inventory, the customer key does not co-locate that inventory. You now need an explicit transaction or workflow across owners, or a different ownership design.
A foreign key requires a referenced record to exist, such as an order referring to an existing customer. Foreign keys enforce relationships inside the database scope that supports them; do not assume an arbitrary cross-shard reference gets the same automatic enforcement. A deleted customer and retained order may require a clear retention and deletion workflow.
Denormalization stores a useful copy of related data, such as the product name at purchase time. That can avoid a cross-shard catalog join and may correctly preserve the historical receipt. For a field that must reflect the latest value, however, copied data needs updates or a freshness contract. Explain why that copy's meaning is suitable, instead of adding denormalization to every design by reflex.
06Resharding: copy, catch up and transfer ownership
Resharding changes how records are distributed among shards, for example to add capacity or relieve an overloaded owner. For the bucket-based scheme above, a move has two jobs: transfer the data and transfer permission to accept writes. Clients may still use an old route during the change, so the handover needs an explicit protocol.
Imagine bucket 17, containing C27, must move from B to C. A safe outline is:
- Copy a consistent snapshot from B to C while B remains the write owner. Bind the snapshot to a committed change-log position L0 and retain all subsequent changes, so there is no gap between snapshot contents and replay.
- Replay subsequent changes so C catches up. Verify record counts/checksums appropriate to the storage model.
- Briefly coordinate the ownership cutover, fencing the old owner so it cannot keep accepting writes after transfer. Fencing means the storage owner rejects commands whose authority is obsolete; merely updating clients does not stop a paused old writer. Publish routing epoch 9, a numbered ownership version, pointing bucket 17 to C.
- A client with epoch 8 reaches B. B rejects or redirects the stale route. The client refreshes metadata and retries the same logical operation safely.
- Retain the old copy until the recovery and stale-client window is closed, then reclaim it.
To switch bucket 17 from B to C, first stop new writes at B and finish or reject writes already running. Record B’s final committed log position; C must apply all changes through it before routing version (epoch) 9 permits writes at C. B then rejects writes using old epoch 8. If the coordinator cannot prove B has stopped accepting writes, it must not enable C. Briefly pausing writes avoids two conflicting histories. A database may use its own consensus or transfer protocol to enforce this handover.
07Interview example: defend a customer shard key
Interviewer: “Why shard orders by customer?”
Candidate: “Most interactive requests list one customer’s orders, so I keep those orders on one shard and route by customer ID. Global reports must query several shards or use an analytical copy. One large customer can still overload a shard, so I measure customer traffic and can split that customer’s data or give it dedicated capacity. To move data, I copy it, apply changes made during copying, then switch write ownership using a new routing version.”
This answer explains placement, the read path, an unfavorable query, and how the system evolves. Merely saying “hash the key” leaves all four undecided.
Choose implementation scope deliberately. PostgreSQL declarative table partitioning can improve pruning and retention management within a database; it does not by itself create a cluster of independently writable servers. If document workloads justify distributed sharding, MongoDB provides mongos routing, configuration metadata and replica-set shards. Its range/hashed placement and zones are product features, while the epoch cutover above is a conceptual protocol to explain ownership—not a claim that MongoDB implements those exact steps or exposes those epoch numbers. Verify supported transactions and constraints for the selected deployment. See PostgreSQL partitioning and MongoDB sharding.
Practise the interview questions
Say your answer aloud before opening the model answer. Then answer the follow-up and compare the reasoning.
How is a shard different from a replica?
Reveal a model answer
“A shard owns a different subset of records; a replica is another copy of the same records. A and B split customers, while A1 and A2 could be copies of shard A. I need separate rules for routing to an owner and for keeping that owner’s copies consistent.”
Interviewer follow-up
Does adding replicas always increase write throughput?
Reveal the follow-up answer
Not for a single-authority write path. Replicas improve resilience and may serve reads, but coordinating them can add write work.
What the answer must demonstrate: Draw ownership and copies separately.
Why choose customer ID for order partitioning?
Reveal a model answer
“The dominant query asks for one customer’s orders. Keeping those records together allows one routed query and local updates of related order data. I would verify the customer traffic distribution and identify global queries that this choice makes more expensive.”
Interviewer follow-up
Would order ID be equally good?
Reveal the follow-up answer
It can spread individual orders better, but listing a customer’s orders needs a secondary location/index path or fanout. The best key depends on the required queries.
What the answer must demonstrate: Connect the key to an actual query.
When does range partitioning help?
Reveal a model answer
“Queries over adjacent keys can target a small set of contiguous ranges. It is useful when the range matches the query, such as a time slice. The risk is skew: always appending to the newest timestamp range can concentrate writes.”
Interviewer follow-up
Is every horizontal partition a range partition?
Reveal the follow-up answer
No. Horizontal means splitting records; hash, list, and other placement rules are alternative ways to do that.
What the answer must demonstrate: Explain the category and the method.
Hashing is uniform. Why is one shard still overloaded?
Reveal a model answer
“Uniform placement distributes keys, not necessarily requests. One customer may account for half the work, or one key may be exceptionally large. I inspect traffic and bytes by key, then consider splitting that workload, replicating reads, or allocating dedicated capacity.”
Interviewer follow-up
Can you split a customer without cost?
Reveal the follow-up answer
It can turn a formerly local order listing or transaction into cross-partition work. I explain that cost and preserve the required ordering or atomicity explicitly.
What the answer must demonstrate: Do not promise hashing eliminates hot keys.
What happens to a join between orders and products?
Reveal a model answer
“If they live on different owners, a local SQL join may no longer cover them. I can perform bounded application lookups, co-locate relevant data, or keep a suitable read copy. For receipts, recording product name and price at purchase time is often the correct historical data.”
Interviewer follow-up
Does denormalization mean every copy must stay current?
Reveal the follow-up answer
No. A historical purchase snapshot should remain historical; a current product description needs an update policy. Similarly, a local index proves uniqueness only in its own scope: a global email claim or order ID requires an explicit cross-shard constraint or single claim owner.
What the answer must demonstrate: Distinguish historical facts from current replicas.
How do you move a shard without losing writes?
Reveal a model answer
“Keep B accepting writes while copying a consistent snapshot tied to log position L0. Apply later logged changes at C. To switch, stop B’s writes and make C apply through B’s final committed position. Then enable C under a new routing version and reject writes using B’s old version. Stale clients refresh their routes and retry the same operation. If I cannot prove B can no longer commit writes, I do not enable C.”
Interviewer follow-up
Why not switch the directory halfway through copying?
Reveal the follow-up answer
C may lack records or writes that arrived after the snapshot. The directory must not send writes to C until C has the required data and B can no longer accept conflicting writes.
What the answer must demonstrate: Separate data catch-up and ownership transfer.
What if the shard directory is unavailable?
Reveal a model answer
“Clients can use a cached version only while the ownership protocol makes stale routes safe. Owners validate epochs and reject invalid writes. For metadata changes I need a durable authoritative directory; guessing a new owner can create conflicting histories.”
Interviewer follow-up
Does caching the directory remove the dependency?
Reveal the follow-up answer
It reduces steady-state lookups but does not remove the need for safe membership updates and recovery.
What the answer must demonstrate: Explain how stale metadata is detected.
How do you support a report for all orders today?
Reveal a model answer
“Customer-based sharding does not localize a global time query. I can fan out bounded queries and merge results for modest needs, or stream order changes into an analytical store partitioned for reporting. I state the reporting freshness delay and avoid making every checkout wait for analytics.”
Interviewer follow-up
What if the report must be an exact cross-shard snapshot?
Reveal the follow-up answer
That needs a defined consistent snapshot or coordinated read protocol. Independently querying owners at different times does not automatically represent one instant.
What the answer must demonstrate: Name the cost of a query the key does not serve.
Blank-page exercise · 20 minutes
Build the answer yourself
Place six customer orders on two shards, add a global-report query, then move one customer’s bucket safely.
- Show an explicit key-to-shard table.
- Explain one query that becomes expensive.
- Separate replication from partitioning.
- Describe catch-up, routing epochs, and old-owner rejection.
Check that each component and design decision follows from your requirements and workload.
Recall the key ideas
Answer from memory before opening each card. Explain why the choice works and what it costs. Revisit missed cards tomorrow.
Data partitioning and shardingWhy choose customer ID as the shard key for order history?Recall first, then reveal
It keeps one customer’s orders together, so one shard can answer the query. Check whether large customers create uneven storage or traffic.
Query together → store together → check imbalance.
Return to lessonData partitioning and shardingShard versus indexRecall first, then reveal
The shard key finds the owner; the index finds records within it.
Which machine, then which record.
Return to lessonData partitioning and shardingOnline movementRecall first, then reveal
Copy, catch up, transfer ownership safely, retire the old copy later.
Copy is not cutover.
Return to lessonFinal revision
Summary and interview notes
Sharding divides record ownership so independent groups can store and serve different parts of a workload. A useful shard key keeps records needed by common queries and correctness rules together; routing, cross-shard work, skew and safe ownership transfer are the costs.
Remember these points
- Horizontal partitioning divides records, replication copies them, and indexing locates records within a query path.
- Ask which queries stay on one shard, how many distinct key values exist, which repeat often, whether they change, and how uneven their data sizes are.
- Uniform hashes distribute distinct keys; they do not split one hot key.
- Joins, global uniqueness and transactions require an explicit supported scope after sharding.
- Copy a consistent snapshot, apply every later change through the final write, then let only the new owner accept writes.
Interview tips
- Show one fast query and one expensive query under the proposed key.
- Ask which related writes and uniqueness claims must remain atomic before choosing a partition boundary.
- During resharding, identify who may write before, during and after cutover.
Important qualifications
- Native table partitioning inside one database is not automatically distributed sharding.
- A cached directory is safe only when storage rejects obsolete ownership.
- The example sacrifices brief cutover availability; a production database may supply a different verified transfer protocol.
Technical references
- MongoDB sharding overviewPrimary reference for ranged/hashed placement and query routing.
- MongoDB zone placementAn implementation example of policy-based range placement; not a universal database guarantee.
- PostgreSQL declarative table partitioningPartition pruning and partition management are distinct from a distributed shard cluster; checked 2026-09-23.
Practice marks stay in this browser.