0%
Data-Intensive Applications
Foundations of Data Systems
Encoding and Evolution
Batch Processing
Stream Processing
Data Quality and Governance
Operational Patterns
Partitioning and Sharding
A single database server has a ceiling. It can only store so much data on disk, handle so many queries per second, and index so many rows before performance collapses. Partitioning is how you break through that ceiling: you split a dataset across multiple machines so that each machine is responsible for a subset of the data. Every read and write hits only the machine that owns the relevant slice, not the entire dataset.
Partitioning is rarely the first tool to reach for, because it forecloses options. Once data is split, cross-partition joins and transactions become distributed problems, with the costs described in the trouble with distributed systems and consistency and consensus. Exhaust the cheaper answers first: an index that removes a scan, in index design and query optimization, a cache that removes the read, in caching strategies in depth, and read replicas, in replication strategies.
The reason this matters is not just storage capacity. It is throughput. If your users table has 500 million rows on one server, every query competes for the same CPU, memory, and I/O bandwidth. Split that table across 10 partitions and each partition handles approximately 50 million rows with approximately one-tenth of the total query load. Queries that would have taken 200ms on a single overloaded server now take 20ms on a partition that is handling one-tenth the traffic.
Vertical Scaling vs. Horizontal Partitioning
Before partitioning, teams typically try vertical scaling: buy a bigger server with more RAM, faster CPUs, and larger SSDs. This works until you hit hardware limits. The largest single-server databases top out at a few terabytes of RAM and a few hundred thousand IOPS. Beyond that, there is no bigger machine to buy. Even if there were, the cost curve is not linear. Doubling a server's capacity often costs 3-4x more. Partitioning across commodity machines gives you linear cost scaling: twice the data costs twice as many servers, not four times the price.
There is a second, subtler reason partitioning wins. A single large server is a single point of failure. No matter how reliable the hardware, it will eventually fail, and when it does, every user is affected simultaneously. Partitioned systems degrade gracefully: losing one node affects only the partition it owns, not the entire dataset.
Partitioning vs. Replication
These two concepts solve different problems and are almost always used together. Replication copies the same data to multiple machines for fault tolerance and read scalability. Partitioning splits different data across machines for write scalability and storage capacity. A production system typically does both: each partition is replicated to 2-3 nodes, so you get fault tolerance per partition and parallel writes across partitions.
Think of it this way: replication makes your data safe (multiple copies survive node failures). Partitioning makes your data fast (each node handles a fraction of the total load). A system with only replication can survive failures but cannot scale writes beyond a single node's capacity. A system with only partitioning can scale writes but loses data when any node fails. Combining both gives you resilience and scalability.
In a typical production setup with 3 partitions and a replication factor of 3, you have 9 total copies of your data spread across nodes. Each partition has a primary replica (handles writes) and two secondary replicas (handle reads and serve as failover). If any single node fails, the partition it owned has two remaining replicas ready to promote. If traffic spikes, read queries can be distributed across all three replicas of each partition.
When Partitioning Becomes Necessary
Not every system needs partitioning. A single PostgreSQL server with 256 GB of RAM and fast NVMe storage can handle millions of rows and thousands of queries per second. Partitioning adds complexity: cross-partition queries, distributed transactions, rebalancing, and routing. Add that complexity only when you need it.
The signals that you need partitioning are concrete and measurable. Your single database server's disk is more than 80% full and growing. Query latency is increasing because indexes no longer fit in memory. Write throughput is hitting the single-node ceiling despite optimization. Backup and restore times are unacceptably long because the dataset is too large. When you see two or more of these signals simultaneously, partitioning moves from theoretical to necessary.
The complexity cost of partitioning is real and should not be underestimated. Cross-partition joins become expensive or impossible. Transactions spanning multiple partitions require distributed coordination protocols (two-phase commit) that add latency and failure modes. Aggregation queries that need data from all partitions trigger scatter-gather. Schema changes must be coordinated across all partitions. Monitoring and alerting must track per-partition health rather than a single database instance. Every one of these costs is worth paying at scale, but premature partitioning adds all this complexity while delivering marginal benefit.
A useful mental model: partitioning is like moving from a single-family house to an apartment building. You gain capacity (more units) and independence (one unit's plumbing issue does not affect others), but you lose simplicity (shared infrastructure, building management, coordination between residents). The apartment building is the right choice when you have outgrown the house, but moving into one when your family of three fits comfortably in a two-bedroom home adds cost and complexity for no benefit.
Cross-Partition Operations
The hardest operations in a partitioned system are those that span multiple partitions. A join between two tables partitioned by different keys requires fetching data from potentially every partition of both tables. An aggregation (COUNT, SUM, AVG) across the full dataset must collect partial results from every partition and merge them.
Distributed transactions across partitions are particularly costly. Two-phase commit (2PC) requires a coordinator to communicate with every partition involved in the transaction, collecting votes and then issuing commit or abort instructions. If any partition is slow or unreachable, the entire transaction blocks. This is why many partitioned databases either limit transactions to single-partition scope (DynamoDB, Cassandra) or accept significant latency overhead for cross-partition transactions (CockroachDB, Spanner).
The practical implication for schema design is clear: design your partition key so that operations within a single business transaction land on the same partition. If a user places an order, the order record, order items, and inventory deduction should all live on the same partition to avoid cross-partition coordination. This sometimes means denormalizing data or using composite partition keys that group related records together.
DynamoDB encourages this pattern explicitly through its single-table design philosophy: store multiple entity types (users, orders, order items) in one table with a shared partition key, so that all related data co-locates on the same partition. This feels unnatural to developers accustomed to relational databases with separate tables per entity, but it is a direct consequence of the cross-partition coordination cost in partitioned systems.
The Partition Key Decision
Everything hinges on how you choose the partition key. The partition key determines which machine owns each record. Choose well and traffic distributes evenly. Choose poorly and one partition gets 90% of the traffic while the others sit idle. This imbalance is called a hot spot, and it is the single most common failure mode in partitioned systems.
Consider a social media platform partitioning posts by user_id. If one celebrity has 50 million followers and posts frequently, the partition holding that user_id receives orders of magnitude more reads than other partitions. The partition key (user_id) is technically correct but practically disastrous for skewed workloads.
A good partition key has three properties: high cardinality (many distinct values to spread across partitions), even distribution (no single value dominates traffic), and query alignment (the most common queries can be answered from a single partition). Finding a key that satisfies all three is the central design challenge of any partitioned system.
Here is a concrete example that illustrates the tension. An analytics platform stores page view events with fields: page_url, user_id, timestamp, country, and device_type. If you partition by page_url, queries for "all views of page X" are fast (single partition), but popular pages create hot spots. If you partition by user_id, user-level analytics are fast, but page-level analytics scatter to every partition. If you partition by country, geographic reports are fast, but countries have vastly different populations (US vs. Luxembourg). Each choice optimizes one access pattern at the expense of others.
The real-world solution is to identify your primary access pattern (the one that handles 80%+ of queries) and optimize the partition key for that pattern. Secondary access patterns are handled through secondary indexes, materialized views, or separate read-optimized copies of the data. Trying to serve all access patterns equally from a single partition key is a design trap that produces mediocre performance everywhere instead of excellent performance where it matters most.
When choosing a partition key, ask two questions: Does this key distribute writes evenly? Does this key let my most common queries hit a single partition? If the answer to either question is no, you need a different key or a compound key strategy.
Terminology: Partitions, Shards, and Regions
Different systems use different words for the same concept. MongoDB calls them shards. HBase calls them regions. Bigtable calls them tablets. Cassandra calls them vnodes. Kafka calls them partitions. The underlying idea is identical: a subset of the data assigned to a specific node. This lesson uses "partition" throughout, but you will encounter all of these terms in practice.
Mid-level engineers should understand why partitioning exists and how partition keys affect query routing. Senior engineers must be able to evaluate partition key choices for skew, hotspot risk, and query patterns. Staff engineers design partitioning strategies that account for data growth, rebalancing overhead, and cross-partition query costs across the full system lifecycle.