System designby Learnastra

Concept lesson · Foundations

Data partitioning and sharding

By Anup Rai

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 visual modelHash partitioning and shard-local queries

The example uses customerNumber modulo two to choose one owner. A global report still fans out or needs a separate read model.

Hash partitioning and shard-local queriesThe example uses customerNumber modulo two to choose one owner. A global report still fans out or needs a separate read model. 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.Toy placement: customerNumber mod 2C12: evenO1, O2C27: oddO3, O4C44: evenO5, O6Owner AC12 + C44 ordersOwner BC27 ordersC27 history: query B only. All customers today: fan out.The shard key selects the owner; the index finds rows inside that owner.
Read the diagram step by step
  1. 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.
  2. C27 last ten orders routes directly to B, then uses a local ordered index.
  3. All orders today crosses customer owners and needs fanout or an analytical model.
  4. 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

  • A shard divides records; a replica copies them.
  • Choose a shard key from queries, related writes and traffic distribution.
  • Moving bytes is not enough: transfer write ownership safely.

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 practice

Useful foundations: Databases, data models, and ACID transactions · Database indexes: B-trees, composite keys and query access

Workload and timing examples are interview assumptions.

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.

Concept in focusSplit rows or split fields?

The cells show the same small dataset. In the vertical split, both partitions keep the ID needed to join the fields.

Split rows or split fields?The cells show the same small dataset. In the vertical split, both partitions keep the ID needed to join the fields. Compare the orientation of the split using the same two records. Horizontal: U1 belongs to shard A and U2 to shard B. Vertical: names are in one partition and regions in another; both keep U1 and U2.Horizontal: divide rowsU1AdaUSU2BoUKshard Ashard BVertical: divide fields; retain the joining keyU1AdaU1USU2BoU2UKProfile partitionRegion partition

Remember: Horizontal cuts between records; vertical cuts between fields.

Read the diagram
  1. Compare the orientation of the split using the same two records.
  2. Horizontal: U1 belongs to shard A and U2 to shard B.
  3. Vertical: names are in one partition and regions in another; both keep U1 and U2.
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.

Worked example diagramPartitioning puts different customer records on A and B. Replication puts another copy of B’s records on its replica. A local index then finds C27’s rows inside B.
Data partitioning and sharding: architecture diagram1. Request: C27 order history to 2. Router: customerNumber mod 2: 1. Read customer C27; 2. Router: customerNumber mod 2 to 4. Shard B: C27 O3/O4: 2. 27 mod 2 = 1: route to B; 2. Router: customerNumber mod 2 to 3. Shard A: C12 O1/O2, C44 O5/O6: Other customer: even keys route to A; 4. Shard B: C27 O3/O4 to 5. Replica of B: same O3/O4: Replication copies B; it does not split B1 → 2: 1. Read customer C272 → 4: 2. 27 mod 2 = 1: route to B2 → 3: Other customer: even keys route to A4 → 5: Replication copies B; it does not split B01Request: C27 orderhistory02Router:customerNumber mod 203Shard A: C12 O1/O2,C44 O5/O604Shard B: C27 O3/O405Replica of B: sameO3/O4
  1. 1 → 21. Read customer C27Request: C27 order history → Router: customerNumber mod 2
  2. 2 → 42. 27 mod 2 = 1: route to BRouter: customerNumber mod 2 → Shard B: C27 O3/O4
  3. 2 → 3Other customer: even keys route to ARouter: customerNumber mod 2 → Shard A: C12 O1/O2, C44 O5/O6
  4. 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:

  1. 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.
  2. Replay subsequent changes so C catches up. Verify record counts/checksums appropriate to the storage model.
  3. 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.
  4. 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.
  5. 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.

Foundation · Question 1

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.”

What the answer must demonstrate: Draw ownership and copies separately.

Applied · Question 2

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.”

What the answer must demonstrate: Connect the key to an actual query.

Foundation · Question 3

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.”

What the answer must demonstrate: Explain the category and the method.

Applied · Question 4

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.”

What the answer must demonstrate: Do not promise hashing eliminates hot keys.

Applied · Question 5

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.”

What the answer must demonstrate: Distinguish historical facts from current replicas.

Applied · Question 6

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.”

What the answer must demonstrate: Separate data catch-up and ownership transfer.

Applied · Question 7

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.”

What the answer must demonstrate: Explain how stale metadata is detected.

Applied · Question 8

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.”

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 lesson
Data 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 lesson
Data partitioning and shardingOnline movementRecall first, then reveal

Copy, catch up, transfer ownership safely, retire the old copy later.

Copy is not cutover.

Return to lesson

Final 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

Practice marks stay in this browser.