LearnThatStack Ace your next interview

Data Partitioning Cheat Sheet - System Design Interviews.

Quick reference for Data Partitioning Cheat Sheet - System Design Interviews - sectioned for fast scanning. Skim the part you're shaky on, walk in confident.

System Design Concepts 14-section reference ~4 min read

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?

  1. Scalability: Handle datasets larger than single machine capacity
  2. Performance: Parallel processing, reduced query scope
  3. Availability: Failure isolation (one partition down ≠ entire system down)
  4. 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

  1. "How would you partition a social media feed?"

    • Hash by user_id for even distribution
    • Consider time-based secondary partitioning for recent posts
  2. "Design a URL shortener with partitioning"

    • Range-based: Pre-assign ranges to servers
    • Hash-based: Hash(short_url) to determine partition
  3. "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

  1. Choose partition key carefully

    • High cardinality
    • Even distribution
    • Aligns with access patterns
  2. Plan for growth

    • Start with more partitions than needed
    • Design for 10x current scale
  3. 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

  1. Choosing wrong partition key

    • Low cardinality → hotspots
    • Timestamp only → all writes to one partition
  2. Ignoring data skew

    • Monitor partition sizes
    • Plan for rebalancing
  3. Over-partitioning

    • Too many small partitions
    • Increased metadata overhead
  4. 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!

Found this useful? Pass it on.
Pro · $10/mo

The sheet is free. Pro goes deeper.

Pro opens the full question library behind every sheet, every refresher and a monthly AI allowance. One subscription, all formats.

Full question library All refreshers Cancel anytime