Summary
Complete guide to data partitioning strategies for system design interviews. Covers horizontal vs vertical partitioning, sharding techniques, consistent hashing, and handling distributed data challenges. Critical for designing databases that scale to billions of records.
What is Data Partitioning?
Definition: Splitting a large dataset across multiple machines/nodes to improve performance, scalability, and availability.
Also Known As: Sharding, Data Distribution, Horizontal Partitioning
Why Data Partitioning?
- Scalability: Handle datasets larger than single machine capacity
- Performance: Parallel processing, reduced query scope
- Availability: Failure isolation (one partition down ≠ entire system down)
- Geographic Distribution: Data closer to users (reduced latency)
Key Concepts
Partition Key
- The attribute used to determine which partition stores the data
- Examples: user_id, timestamp, geographic_region
Partition Function
- Algorithm that maps partition key → partition number
- Must be deterministic and uniform
Partitioning Strategies
1. Range-Based Partitioning
How: Data divided based on ranges of partition key
Partition 1: A-F
Partition 2: G-M
Partition 3: N-S
Partition 4: T-Z
Pros:
- Simple implementation
- Range queries efficient
- Natural ordering preserved
Cons:
- Uneven distribution (hotspots)
- Manual rebalancing needed
Use Case: Time-series data, alphabetical user directories
2. Hash-Based Partitioning
How: Hash function on partition key determines partition
partition = hash(user_id) % num_partitions
Pros:
- Even data distribution
- No hotspots (with good hash function)
- Automatic load balancing
Cons:
- Range queries inefficient
- Resharding complex (requires rehashing)
Use Case: User data, session storage
3. List-Based Partitioning
How: Explicit mapping of values to partitions
Partition_US: ['USA', 'Canada', 'Mexico']
Partition_EU: ['UK', 'France', 'Germany']
Partition_ASIA: ['India', 'China', 'Japan']
Pros:
- Business logic friendly
- Predictable data location
Cons:
- Manual maintenance
- Limited flexibility
Use Case: Geographic data, category-based systems
4. Composite Partitioning
How: Combination of multiple strategies
Level 1: Hash by user_id
Level 2: Range by timestamp
Use Case: Large-scale systems needing both distribution and query efficiency
Consistent Hashing
Problem: Adding/removing nodes in hash-based systems requires resharding all data
Solution: Consistent hashing minimizes data movement
# Only K/n keys need remapping when adding/removing nodes
# K = total keys, n = number of nodes
Key Features:
- Nodes and keys mapped to same hash ring
- Each key assigned to nearest node clockwise
- Virtual nodes for better distribution
Architecture Patterns
1. Shared Nothing Architecture
- Each node independent
- No shared memory/disk
- Horizontal scaling
2. Master-Slave Replication
Partition 1: Master1 → Slave1a, Slave1b
Partition 2: Master2 → Slave2a, Slave2b
3. Peer-to-Peer
- All nodes equal
- No single point of failure
- Example: Cassandra, DynamoDB
⚖️ Trade-offs & Challenges
1. Cross-Partition Operations
Challenge: JOINs, transactions across partitions
Solutions:
- Denormalization
- Application-level joins
- Distributed transactions (2PC, Saga)
2. Rebalancing
When Needed:
- Uneven data distribution
- Adding/removing nodes
- Hotspot detection
Strategies:
- Dynamic partitioning
- Consistent hashing
- Pre-splitting
3. Partition Size
Too Large:
- Performance degradation
- Difficult to move/backup
Too Small:
- Overhead of managing many partitions
- Metadata explosion
Rule of Thumb: 1-10GB per partition (varies by system)
Real-World Examples
MongoDB
// Range-based sharding
sh.shardCollection("mydb.users", { "country": 1, "user_id": 1 })
// Hash-based sharding
sh.shardCollection("mydb.posts", { "post_id": "hashed" })
Cassandra
CREATE TABLE users (
user_id UUID,
country TEXT,
name TEXT,
PRIMARY KEY ((country), user_id) -- country as partition key
);
DynamoDB
# Partition key determines physical partition
table = dynamodb.create_table(
TableName='Users',
KeySchema=[
{'AttributeName': 'user_id', 'KeyType': 'HASH'}, # Partition key
{'AttributeName': 'timestamp', 'KeyType': 'RANGE'} # Sort key
]
)
Interview Tips
Common Questions
"How would you partition a social media feed?"
- Hash by user_id for even distribution
- Consider time-based secondary partitioning for recent posts
"Design a URL shortener with partitioning"
- Range-based: Pre-assign ranges to servers
- Hash-based: Hash(short_url) to determine partition
"How to handle hot partitions?"
- Monitor partition metrics
- Split hot partitions
- Add caching layer
- Use composite keys to spread load
Key Metrics to Discuss
- Throughput: Requests/second per partition
- Storage: Data size per partition
- Latency: Query response time
- Availability: Partition failure impact
Best Practices
Choose partition key carefully
- High cardinality
- Even distribution
- Aligns with access patterns
Plan for growth
- Start with more partitions than needed
- Design for 10x current scale
Monitor continuously
- Track partition sizes
- Identify hotspots
- Plan rebalancing
Quick Decision Framework
Need range queries? → Range-based partitioning
Need even distribution? → Hash-based partitioning
Have categorical data? → List-based partitioning
Complex requirements? → Composite partitioning
⚡ Performance Considerations
Write Performance
- More partitions = higher write throughput
- Consider write patterns when choosing strategy
Read Performance
- Partition pruning: Query only relevant partitions
- Parallel queries across partitions
- Caching frequently accessed partitions
Network Overhead
- Cross-partition queries require network calls
- Minimize partition-spanning operations
- Consider data locality
🚨 Common Pitfalls
Choosing wrong partition key
- Low cardinality → hotspots
- Timestamp only → all writes to one partition
Ignoring data skew
- Monitor partition sizes
- Plan for rebalancing
Over-partitioning
- Too many small partitions
- Increased metadata overhead
Under-partitioning
- Partitions too large
- Cannot scale horizontally
📝 Summary Checklist
✅ Understand data access patterns
✅ Choose appropriate partition strategy
✅ Select high-cardinality partition key
✅ Plan for rebalancing and growth
✅ Consider cross-partition operations
✅ Monitor partition health
✅ Handle failures gracefully
✅ Test with realistic data distributions
Remember: The best partitioning strategy depends on your specific use case, data characteristics, and access patterns. Always consider trade-offs!