Not exactly. Sharding is useful for both relational and non-relational databases, but it is generally more complex to implement in relational databases. MongoDB is well known for built-in sharding.
🔥 Ramesh Style
Think of sharding as:
One huge database → split data across multiple database servers.
Suppose we have 1 billion customers.
Without sharding
Application
|
↓
┌───────────┐
│ DB Server │
│ 1 Billion │
│ customers │
└───────────┘Problems:
CPU → 🔴
RAM → 🔴
Storage → 🔴
I/O → 🔴With sharding
Application
|
Sharding Router
|
┌─────────────┼─────────────┐
↓ ↓ ↓
DB Shard 1 DB Shard 2 DB Shard 3
Customer Customer Customer
1-10M 10M-20M 20M-30MNow the workload is distributed.
Relational DB — Can we shard?
Yes. Absolutely.
For example, PostgreSQL, MySQL, Oracle, etc. can be sharded using different approaches/tools.
But the challenge is relational relationships.
Suppose:
CUSTOMER
---------
customer_id
name
ORDER
---------
order_id
customer_id
amountIf we shard by customer_id:
Shard 1
Customer 1-1M
Orders for those customers
Shard 2
Customer 1M-2M
Orders for those customersGood if most queries are:
SELECT *
FROM orders
WHERE customer_id = 123;But consider:
SELECT *
FROM customer c
JOIN orders o
ON c.customer_id = o.customer_id
WHERE ...If related data is spread across shards, you may need:
Application
↓
Shard 1 ─────┐
├── Distributed JOIN
Shard 2 ─────┘That becomes expensive and complicated.
MongoDB — Why sharding feels more natural
MongoDB was designed around documents, rather than relational tables and joins.
Example:
{
"customerId": 101,
"name": "Ramesh",
"orders": [
{"orderId": 1, "amount": 500},
{"orderId": 2, "amount": 800}
]
}You can shard using a shard key, for example:
customerIdThen:
customerId 1-1M
↓
Shard 1
customerId 1M-2M
↓
Shard 2
customerId 2M-3M
↓
Shard 3MongoDB provides built-in distributed sharding infrastructure.
⭐ Why MongoDB is commonly associated with sharding
Because MongoDB provides:
MongoDB
|
├── Shard
├── Shard
├── Shard
|
├── Config Servers
|
└── mongos routersConceptually:
Application
|
↓
mongos
|
├────────→ Shard 1
├────────→ Shard 2
└────────→ Shard 3The router determines which shard owns the relevant data.
🔥 Important distinction
Don't say in an interview:
❌ "Sharding is not useful for relational databases."
Say:
✅ "Sharding is applicable to both relational and NoSQL databases. However, relational databases can make distributed joins, transactions, foreign keys and cross-shard consistency more complicated. MongoDB has built-in sharding capabilities and its document model can make certain sharding patterns more natural."
That's a much stronger Senior Architect answer.
Sharding vs Replication
Another very important interview distinction:
Replication
Same data copied to multiple nodes:
Master
/ \
↓ ↓
Replica 1 Replica 2
Same dataPurpose:
HA
Read scaling
Disaster recoverySharding
Different data on different nodes:
Shard 1 → Customers 1-1M
Shard 2 → Customers 1M-2M
Shard 3 → Customers 2M-3MPurpose:
Horizontal data scaling
Horizontal write scaling
Storage scalingEasy memory trick
Replication = copy the data.
Sharding = split the data.
And in real production systems, you often combine them:
Cluster
|
┌───────────┼───────────┐
↓ ↓ ↓
Shard 1 Shard 2 Shard 3
/ \ / \ / \
R1 R2 R1 R2 R1 R2Sharding + Replication = scale + high availability.
-------------------------
Design challengs - data access/data access pattern challenges in sharding.
📚 Sharding — Class Notes | Ramesh Style
1. What is Sharding?
Sharding means horizontally partitioning data across multiple database servers (shards).
Instead of keeping all data in one DB:
Application
|
↓
Single Database
1 Billion Recordswe distribute it:
Application
|
Shard Router
┌───────────┼───────────┐
↓ ↓ ↓
Shard 1 Shard 2 Shard 3
Users 1-1M Users 1M-2M Users 2M-3M🎯 Main purpose
Sharding
↓
Horizontal Scaling
↓
More CPU + RAM + Storage + I/O
↓
Handle larger data and traffic2. Sharding Key
A sharding key determines which shard stores a particular record.
Example:
customerIdSuppose:
customerId = 101Routing logic:
customerId
↓
Shard Key
↓
Shard Router
↓
Shard 1Another customer:
customerId = 2,500,000
↓
Shard 3Important
Choosing a good shard key is one of the most critical decisions in sharding.
A bad shard key can create a hot shard.
3. Complexity Introduced by Sharding
Without sharding:
Application
|
↓
DBWith sharding:
Application
|
↓
Shard Router
|
┌───┼────┐
↓ ↓ ↓
S1 S2 S3Now the application/system must know:
Which shard contains the data?
How should requests be routed?
What happens when a shard fails?
How are connections managed?
How are transactions handled across shards?
Therefore:
Sharding improves scalability but increases architectural complexity.
4. Shard Routing
Shard routing means determining the correct shard for a request.
Example:
GET customer 101
↓
customerId = 101
↓
Shard Key
↓
Shard Router
↓
Shard 1For:
GET customer 2500000the router might send it to:
Shard 3Flow
Request
↓
Extract Sharding Key
↓
Calculate/lookup shard
↓
Route request
↓
Correct shard
↓
Response5. Limited Data Model
Sharding can influence how you design your data model.
Suppose:
Customer
|
└── OrdersIf both are stored on the same shard:
Shard 1
├── Customer 101
└── Orders of Customer 101queries are relatively easy.
But if they are on different shards:
Shard 1 Shard 2
Customer Orders
| |
└──────── JOIN ─────────┘Now we have a cross-shard operation.
This can be expensive.
6. Cross-Shard Join
Consider:
SELECT *
FROM customer c
JOIN orders o
ON c.customer_id = o.customer_id;Without sharding:
Application
↓
DB
↓
Customer JOIN OrdersWith sharding:
Query
|
┌───────┴───────┐
↓ ↓
Shard 1 Shard 2
Customer A Orders B
\ /
\ /
Cross-shard
JOINThis can cause:
More network calls
↓
More latency
↓
More complexityInterview point
Cross-shard joins are generally expensive and should be minimized through data modeling, co-location, denormalization, or other design techniques.
7. Limited Data Access Patterns
Suppose we shard by:
userIdThen this query is excellent:
SELECT *
FROM orders
WHERE user_id = 101;Because the system knows:
userId = 101
↓
Shard 1Only one shard needs to be queried.
This is called a targeted query.
But consider:
SELECT *
FROM orders
WHERE amount > 10000;There is no userId.
Which shard contains the matching records?
The system may need:
Query
|
┌───────┼───────┐
↓ ↓ ↓
Shard1 Shard2 Shard3
↓ ↓ ↓
Result Result Result
└───────┼───────┘
↓
MergeThis is commonly called a scatter-gather query.
Result:
Multiple shards
↓
More network calls
↓
More processing
↓
Higher latency8. Denormalization
One solution to expensive cross-shard queries is denormalization.
Instead of:
Customer
|
↓
Ordersrequiring a join, we may store commonly needed information together.
Example:
{
"orderId": 101,
"customerId": 500,
"customerName": "Ramesh",
"amount": 5000
}Now the order query may not need to contact the Customer shard.
Trade-off
Denormalization
↓
Faster reads
+
Less cross-shard querying
↓
But
↓
Duplicate data
+
Consistency/update complexity9. Caching
Another solution:
Application
|
↓
Cache
|
↓
If not found
|
↓
Multiple shardsFor frequently requested cross-shard information:
Redis
↓
Cached resultcan reduce expensive distributed queries.
10. Operational Challenges
Sharding doesn't end after deployment.
You now have:
Shard 1
Shard 2
Shard 3
Shard 4
...
Shard NEach needs operational management.
Backup
Without sharding:
Backup DBWith sharding:
Backup Shard 1
Backup Shard 2
Backup Shard 3
...11. Recovery
Suppose:
Shard 2 ❌You need to recover:
Shard 2
↓
Replica / Backup
↓
Restore
↓
Rejoin systemThe recovery strategy becomes more complicated as the number of shards increases.
12. Rebalancing
🔥 Very important concept.
Suppose initially:
Shard 1 → 30%
Shard 2 → 30%
Shard 3 → 40%Application grows.
Now:
Shard 1 → 20%
Shard 2 → 20%
Shard 3 → 60% 🔴Shard 3 becomes overloaded.
We need rebalancing.
Before:
Shard 1 → 30%
Shard 2 → 30%
Shard 3 → 40%
↓ Rebalance
After:
Shard 1 → 33%
Shard 2 → 33%
Shard 3 → 34%This can involve moving large amounts of data.
13. Hot Shard
Bad shard-key selection can create a hot shard.
Example:
Shard key = countrySuppose:
India → 70% of users
USA → 10%
UK → 5%
Others → 15%Then:
Shard India
████████████████████ 70% 🔥
Shard USA
██ 10%
Shard UK
█ 5%One shard gets most of the traffic.
That's called a hotspot/hot shard.
Lesson
Good shard-key selection should distribute both data and workload.
14. Sharding + Replication
In real systems, we commonly combine both.
Cluster
|
┌────────────┼────────────┐
↓ ↓ ↓
Shard 1 Shard 2 Shard 3
/ \ / \ / \
↓ ↓ ↓ ↓ ↓ ↓
Master Replica Master Replica Master ReplicaSharding gives:
ScalabilityReplication gives:
High AvailabilityTherefore:
Sharding = split the data.
Replication = copy the data.
15. Key Limitations — Interview Table
| Problem | Why? | Typical solution |
|---|---|---|
| More complexity | Multiple DB nodes | Sharding middleware/router |
| Cross-shard joins | Data distributed | Co-locate/denormalize |
| Scatter-gather queries | Query doesn't target one shard | Better shard key/indexing |
| Hot shard | Poor shard-key distribution | Choose better shard key |
| Rebalancing | Data grows unevenly | Automated rebalancing |
| Backup complexity | Multiple shards | Centralized backup strategy |
| Recovery complexity | Multiple failure domains | Replicas + automation |
| Distributed transactions | Data spans shards | Avoid where possible / use appropriate transaction design |
⭐ Most Important Interview Concept: Shard Key
If interviewer asks:
"What is the biggest design consideration in sharding?"
Answer:
"Choosing the shard key."
A good shard key should ideally provide:
High cardinality
+
Even distribution
+
Good query targeting
+
Low hotspot riskFor example:
userId
customerId
accountIdcan be good candidates depending on the workload.
But a field like:
countrymay create hotspots if one country dominates the traffic.
🧠 Complete Sharding Flow
APPLICATION
|
↓
Shard Router
|
Sharding Key
|
┌────────┼────────┐
↓ ↓ ↓
Shard 1 Shard 2 Shard 3
| | |
↓ ↓ ↓
Data Data DataFor a good targeted query:
Request
↓
userId = 101
↓
Shard Router
↓
Shard 1
↓
ResponseFor a bad/non-targeted query:
Request
↓
No shard key
↓
Shard 1 ──┐
Shard 2 ──┼──→ Scatter
Shard 3 ──┘
↓
Gather
↓
Merge
↓
Response🎯 Ramesh Interview Summary
Remember these 5 points:
1. Sharding
↓
Split data horizontally
2. Shard Key
↓
Determines where data goes
3. Cross-Shard Operations
↓
Expensive and complex
4. Rebalancing
↓
Required when data/workload becomes uneven
5. Sharding + Replication
↓
Scalability + High Availability🔥 One-line interview answer
"Sharding provides horizontal scalability by distributing data across multiple nodes, but it increases application and operational complexity. The biggest challenges are choosing the right shard key, avoiding hotspots, minimizing cross-shard joins and scatter-gather queries, and managing rebalancing, backup, recovery, and distributed transactions."
------
Read Replication will handle more read traffic, Sharding - will handle more write traffic, Now we are going to talk on consisent hashing topic:
📚 Consistent Hashing & Replication — Class Notes
Ramesh Style | Distributed Systems
These two concepts are often used together in distributed databases, caches, and systems like Redis.
The easiest way to remember:
Consistent Hashing → Where should the data go?
Replication → How many copies of the data should we keep?
1. Consistent Hashing
Definition
Consistent hashing is a technique for distributing data across multiple nodes while minimizing data movement when nodes are added or removed.
Normal hashing:
hash(key) % numberOfNodeslooks simple, but it has a major problem.
2. Problem with Normal Hashing
Suppose we have 3 servers:
Server 0
Server 1
Server 2We use:
hash(key) % 3Suppose:
User A → hash = 10
User B → hash = 20
User C → hash = 30Now add Server 3:
Server 0
Server 1
Server 2
Server 3Now we calculate:
hash(key) % 4Many keys get a different result.
Before:
Key A → Server 0
Key B → Server 1
Key C → Server 2
After:
Key A → Server 2
Key B → Server 0
Key C → Server 3💥 A huge amount of data needs to move.
3. Consistent Hashing Solution
Instead of directly doing:
hash(key) % Nwe create a hash ring.
Think of it like a clock:
0
┌───────────┐
330 / \ 30
/ \
300 | | 60
| |
270 | | 90
\ /
240 \ / 120
└───────────┘
210 180 150The entire hash space forms a circle.
4. Put Servers on the Ring
Suppose:
Node A → hash position 50
Node B → hash position 150
Node C → hash position 270Conceptually:
0
|
Node A(50)
|
|
Node B(150)
|
|
Node C(270)
|
|
back to 05. Where does data go?
Hash the key.
Suppose:
hash("user:101") = 80Find position 80 on the ring.
Move clockwise until you find a node.
Node A
50
|
| ← key = 80
|
↓
Node B
150Therefore:
user:101 → Node BMemory trick:
Hash the key → put it on the ring → walk clockwise → first node owns the key.
6. What happens when a node is added?
Suppose:
Node A
Node B
Node CNow add:
Node DNode D gets a position on the ring.
A
|
D ← new node
|
B
|
COnly the keys in the region affected by Node D need to move.
Before
Some keys → Node BAfter
Some of those keys → Node DOther keys remain where they were.
Therefore:
Only a relatively small portion of data needs to move.
This is the biggest advantage of consistent hashing.
7. Node Removal
Suppose:
A
B
CNode B crashes:
A
B ❌
CKeys previously assigned to B can move to the next node according to the ring.
Keys belonging to B
↓
Next node
↓
CAgain, we avoid reshuffling the entire dataset.
8. Virtual Nodes
🔥 Important interview concept.
If each physical server gets only one position:
A ---------------- B ---------------- Cdistribution might not be perfectly balanced.
So we create virtual nodes.
Instead of:
Server A → one positionwe use:
Server A
├── A1
├── A2
├── A3
├── A4
└── A5Similarly:
Server B → B1 B2 B3 B4 B5
Server C → C1 C2 C3 C4 C5Ring:
A1 → B1 → C1 → A2 → B2 → C2 → A3 → ...This improves load distribution.
Why?
Because each physical server owns multiple small regions instead of one large region.
9. Benefits of Consistent Hashing
① Dynamic scalability
Add nodes:
Node A
Node B
Node C
↓ Add Node D
Node A
Node B
Node C
Node DOnly part of the data moves.
② Less data movement
Normal hashing:
Node count changes
↓
Many keys remapped
↓
Huge data movementConsistent hashing:
Node count changes
↓
Limited keys remapped
↓
Less data movement③ Better availability
When combined with replication:
Consistent Hashing
+
Replication
↓
Scalable + Fault Tolerant10. Replication
Now let's understand the second concept.
Replication means maintaining multiple copies of the same data on different nodes.
Example:
Data
|
┌─────┴─────┐
↓ ↓
Node A Node B
Copy 1 Copy 2If Node A fails:
Node A ❌
Node B ✅
Copy availableThis improves availability and fault tolerance.
11. Why do we need Replication?
Without replication:
User Data
↓
Node A
↓
Node A crashes
↓
💥 Data unavailableWith replication:
User Data
|
┌─┴────┐
↓ ↓
A B
Copy CopyIf A crashes:
A ❌
B ✅The application can continue using B depending on the system's architecture.
12. Synchronous Replication
Suppose:
Client
|
| WRITE
↓
Primary
|
├────────→ Replica 1
|
└────────→ Replica 2The write is considered successful only after required replicas acknowledge it.
Client
↓
Primary
↓
Replica 1 ✅
↓
Replica 2 ✅
↓
ACK to ClientAdvantage
Stronger consistency.
Disadvantage
Higher latency.
If a replica is slow:
Primary
↓
Replica slow 🐌
↓
Client waits13. Asynchronous Replication
Here:
Client
|
↓
Primary
|
↓
ACK immediately
|
↓
Client continues
Meanwhile:
Primary → ReplicaFlow:
Client
↓
Primary
↓
ACK ✅
↓
Client
Primary
↓
Replica
↓
Replication laterAdvantage
Lower latency and better availability.
Risk
If primary crashes before replication:
Primary
|
| New data
↓
❌ Crash
Replica
|
↓
Doesn't have latest dataPotential data loss.
14. Synchronous vs Asynchronous
| Feature | Synchronous | Asynchronous |
|---|---|---|
| Write acknowledgement | After required replicas confirm | Primary can acknowledge first |
| Consistency | Stronger | Eventual |
| Latency | Higher | Lower |
| Availability | Can be reduced if replicas unavailable | Generally better |
| Risk of lost recent writes | Lower | Higher |
Memory trick:
Sync = Wait for replicas.
Async = Don't wait for replicas.
15. Master-Slave / Primary-Replica
Common architecture:
Primary
|
┌────────┼────────┐
↓ ↓ ↓
Replica 1 Replica 2 Replica 3Writes:
Client
↓
PrimaryReads may be served by replicas depending on consistency requirements:
Client
↓
ReplicaIf primary fails:
Primary ❌
↓
Failover
↓
Replica promoted
↓
New PrimaryThis connects directly to your Redis Sentinel topic.
16. Consistent Hashing + Replication Together
🔥 This is the most important architecture to understand.
Suppose:
Hash Ring
┌────────────────────────┐
│ │
Node A Node B
Primary Primary
│ │
A-Replica B-Replica
│ │
└────────────────────────┘Consistent hashing decides:
Which node owns this key?
Replication decides:
Where should copies of this key exist?
So:
DATA
|
┌─────────┴─────────┐
↓ ↓
Consistent Hashing Replication
↓ ↓
Which node? How many copies?
↓ ↓
Partition data Fault tolerance17. Failure Example
Suppose:
Key = user:101
Hash
↓
Node BReplication factor = 3:
user:101
|
├── Node B
├── Node C
└── Node DNow:
Node B ❌Copies remain:
Node C ✅
Node D ✅So data remains available.
18. Node Addition Example
Initially:
A
B
CAdd D:
A
B
C
DConsistent hashing:
Only affected hash ranges
↓
Move some data
↓
DReplication then ensures the moved data continues to have the required number of copies.
19. Challenges
A. Node Failure
Node A ❌
↓
Detect failure
↓
Promote replica / route elsewhereNeed:
Failure detection
Failover
Data recovery
Replica synchronization
B. Replication Overhead
Suppose:
Original data = 1 TB
Replication factor = 3Approximate storage requirement:
1 TB × 3
=
3 TBPlus network traffic for replication.
So replication improves reliability but costs:
Storage
+
Network bandwidth
+
CPU20. CAP Connection
This connects directly to what we just studied.
Replication
↓
Multiple copies
↓
Network partition
↓
Replicas may disagree
↓
Consistency vs Availability trade-offFor example:
Primary
X
ReplicaIf they cannot communicate:
Choose availability
Continue serving requests:
Replica → accepts operationsBut state may diverge.
Choose consistency
Stop/reject unsafe operations:
Replica → rejectAvailability is reduced.
🧠 Complete Class Diagram
DISTRIBUTED DATABASE
|
┌──────────┴──────────┐
↓ ↓
CONSISTENT HASHING REPLICATION
↓ ↓
Where data goes? How many copies?
↓ ↓
Hash Ring Primary + Replicas
↓ ↓
Shard selection Fault tolerance
↓ ↓
Add/remove node Node failure
↓ ↓
Minimal data movement High availability
|
↓
Virtual Nodes
↓
Better load balance🎯 Interview Questions
Q1. What problem does consistent hashing solve?
It minimizes data redistribution when nodes are added or removed from a distributed system.
Q2. How does consistent hashing work?
Keys and nodes are mapped onto a hash ring. A key is assigned to the next node clockwise in the ring. When nodes change, only the affected portion of the keyspace needs to be remapped.
Q3. Why use virtual nodes?
Virtual nodes assign multiple positions to each physical node, improving load distribution and reducing hotspots.
Q4. What is replication?
Maintaining multiple copies of data across different nodes to improve availability and fault tolerance.
Q5. Synchronous vs asynchronous replication?
Synchronous replication waits for required replica acknowledgements before confirming the write, providing stronger consistency but higher latency. Asynchronous replication acknowledges earlier and replicates later, providing lower latency but a greater possibility of losing recent writes.
Q6. What happens when a node fails?
Node fails
↓
Failure detection
↓
Replica available?
↓
Yes
↓
Failover / rerouting
↓
Continue serving⭐ Final Ramesh Memory Trick
Remember these four questions:
Consistent Hashing
↓
"WHERE?"
↓
Which node should store my data?
Replication
↓
"HOW MANY?"
↓
How many copies should I maintain?
Sharding
↓
"HOW TO SPLIT?"
↓
How do I distribute the dataset?
Failover
↓
"WHAT IF NODE DIES?"
↓
How do I continue serving?🔥 One-line interview answer
"Consistent hashing determines where data should be placed with minimal redistribution when nodes change, while replication maintains multiple copies of that data for fault tolerance and availability. Together they provide a foundation for scalable and highly available distributed systems."
No comments:
Post a Comment