Search IOCombats

Search challenges, guides, questions and articles

System DesignTopic 11 of 12BeginnerApr 5, 2026

Database Sharding

Learn how database sharding distributes data across multiple databases to handle large scale traffic and improve performance.

system-designdatabasesscalability

Why It Matters

A single database has limits.

Example:

  • Max capacity: 5,000 writes per second
  • Traffic grows to 20,000 writes per second

Result:

  • high latency
  • write failures

Replication helps reads, but not writes.

Sharding solves this by distributing data.

What This Concept Actually Means

Sharding means splitting data across multiple databases.

Each database stores a subset of data.

Think like this:

  • Instead of one big database
  • You have multiple smaller databases handling different data

How It Works (With Example)

Let’s take a user system with 10 million users.

Without Sharding

  • 1 database stores all users
  • Writes: 10,000 per second
  • DB limit: 5,000 per second

System fails.

With Sharding

  • Split users into 4 shards

Example strategy:

  • Shard 1: user_id 1 to 2.5M
  • Shard 2: user_id 2.5M to 5M
  • Shard 3: user_id 5M to 7.5M
  • Shard 4: user_id 7.5M to 10M

Now:

  • Each shard handles around 2,500 writes per second
  • System handles total 10,000 writes per second

Flow:

Diagram
100%
flowchart LR User --> API API --> Shard1 API --> Shard2 API --> Shard3 API --> Shard4
visualized byIOCombats

Key Techniques / Variations

1. Range Based Sharding

Split data based on value ranges.

Example:

  • user_id ranges

Problem:

  • uneven distribution if data is skewed

2. Hash Based Sharding

Use hash function to distribute data.

Example:

  • user_id mod 4

Benefit:

  • even distribution

3. Directory Based Sharding

Use lookup service to find shard.

Example:

  • user_id mapped to shard via metadata service

Benefit:

  • flexible

4. Geo Based Sharding

Split data by region.

Example:

  • India users in one shard
  • US users in another

Trade-offs and Design Decisions

When to use sharding

  • high write traffic
  • large datasets (millions or billions of rows)

Pros

  • scales writes
  • distributes load
  • improves performance

Cons

  • complex queries across shards
  • joins become difficult
  • rebalancing shards is hard

Real-world example

  • Single DB limit: 5,000 writes per second
  • 4 shards: total capacity becomes around 20,000 writes per second

Cost example:

  • 1 DB server: 20k per month
  • 4 DB servers: 80k per month

But system scales 4 times

Architecture / Flow Diagram

Diagram
100%
flowchart LR Client --> API API --> Router Router --> Shard1 Router --> Shard2 Router --> Shard3 Router --> Shard4
visualized byIOCombats

Failure Modes To Watch

  • uneven data distribution
  • hotspot shards
  • cross shard queries slowing system
  • shard rebalancing complexity

Design Checklist

Ask yourself:

  • What is data size? (e.g., millions or billions of records)
  • What is write load? (e.g., 10k writes per second)
  • How will I choose shard key?
  • Will data be evenly distributed?
  • Do I need cross shard queries?

Summary

  • Sharding splits data across multiple databases
  • It helps scale write-heavy systems
  • Choose shard key carefully to avoid hotspots
  • Adds complexity but enables massive scale
  • Often used with replication for full scalability

Advertisement

Related Videos
Watch these videos to reinforce the topic after reading the guide.
Frequently Asked Questions

When should I use sharding?

When a single database cannot handle the data size or traffic, especially for write-heavy systems.

Does sharding improve read performance?

Yes, but it is mainly used to scale writes and distribute data load.