Data Systems
Data Partitioning
When a single node cannot store or process all data, partitioning splits data across multiple nodes. The partitioning strategy determines how data is distributed and what queries are efficient.
- Horizontal Partitioning â Splitting rows across nodes
- Range Partitioning â Contiguous key ranges per partition
- Hash Partitioning â Hash function determines partition
Partitioning is the key to scaling data storage and query throughput beyond a single machine.
Why Partition?
A single database has hard limits: disk space, CPU, memory, and I/O bandwidth. Partitioning (sharding) distributes data across multiple machines.
Horizontal vs Vertical Partitioning
| Strategy | Description | Use Case |
|---|---|---|
| Vertical | Split columns into different tables/databases | Rarely used in modern systems |
| Horizontal | Split rows across multiple databases | Standard sharding approach |
Vertical Partitioning
Splitting a table by columns into separate tables or databases.
Horizontal Partitioning
Splitting rows across multiple database instances.
Range Partitioning
Each partition owns a contiguous range of keys.
Range Partitioning Trade-offs
Hash Partitioning
A hash function determines which partition owns a key.
Consistent Hashing
When N changes, naive hash partitioning remaps all keys. Consistent hashing minimizes key movement.
Partition Key Selection
The partition key determines data distribution and query efficiency.
Rebalancing
When partitions become uneven or nodes are added/removed, data must be rebalanced.
Cross-Partition Queries
Queries spanning multiple partitions are inherently more expensive.
Scatter-Gather
A query that touches all partitions must scatter (send to all) and gather (combine results).
Partitioning in Practice
| System | Partitioning Strategy | Notes |
|---|---|---|
| MySQL (Vitess) | Hash or range | Configurable per table |
| MongoDB | Hash or range | Auto-balancing |
| Cassandra | Hash (consistent) | Virtual nodes for balance |
| DynamoDB | Hash (consistent) | Automatic partition management |
| CockroachDB | Range | Automatic split/merge |
Practice Exercises
-
Design: Design a partitioning strategy for a ride-sharing app with 100M rides/day. Queries: rides by user, rides by city, rides by time range. What partition key supports the most common query?
-
Analysis: Compare range and hash partitioning for an e-commerce order system where 80% of queries are "my orders" (by user_id) and 20% are "all orders today" (by timestamp).
-
Rebalancing: A system has 4 partitions. Adding a 5th node with consistent hashing moves ~20% of keys. With fixed partitioning and 16 partitions, how many partitions move?
-
Hotspot Mitigation: A time-series database receives 1M writes/second, all with current timestamps. Design a partitioning strategy that avoids hotspots while maintaining time-range query efficiency.
What to Learn Next
-> Consistent Hashing Hash rings, virtual nodes, and minimal key redistribution.
-> Data Replication Leader-follower, multi-leader, and conflict resolution.
-> Databases SQL vs NoSQL, indexing, replication, and sharding.
-> Database Indexing B-trees, LSM trees, and query optimization.
-> CAP Theorem Consistency models, availability, and partition tolerance.
-> Scalability Fundamentals Vertical vs horizontal scaling, load balancing, and capacity planning.