Partitioning
Splitting data or work by key so each part is handled independently: scaling out, and giving each key a single owner.
Distribution
Learn it
Eventually one machine can't handle all the load. Copies of a stateless tier are easy to add; copies of state aren't, because two copies accepting writes for the same thing must coordinate.
Partitioning (sharding) assigns each key (a user, a document, an account) to exactly one partition, and each partition is handled independently.
- The key is the important decision. Operations that must be atomic or ordered together should share a partition key; operations across partitions lose those properties or need expensive coordination.
- Mapping: hashing spreads load evenly but scatters ranges; range partitioning keeps neighbours together but invites hotspots. Consistent hashing keeps most keys in place when partitions change.
- Routing: something must find the partition that owns a key, including during moves.
Check
All writes for document 42 go to one owner process. What does that buy?Partitioning doesn't fix hot keys. If one key (a celebrity account, an all-hands document) gets more load than one partition can handle, adding partitions doesn't help: it still lives on one. You have to split the work for that key: separate the single-writer part from reads and fan-out, cache it, or batch it.
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Hot partition
- One key or range receives disproportionate load.
- Cross-partition operations
- Transactions or queries spanning partitions become slow, complex or non-atomic.
- Split ownership during rebalancing
- Two nodes believe they own a partition while it moves; fence ownership.
Instead, consider
- Vertical scaling
- A bigger machine still fits the load. It is simpler, and often enough for longer than expected.
- Read replicas
- Reads dominate and can tolerate slight staleness, while writes still fit one primary.
- Caching
- Load is mostly repeated reads of the same data.
In practice
- Application-level sharding
- Shard ID derived from a key and mapped to a database.
- Kafka partitions
- Per-partition order and consumer ownership.
- Consistent-hash routing at the load balancer
- Route all of a document's connections to one server.
- Distributed databases
- Automatic range or hash partitioning (DynamoDB, Spanner, CockroachDB).
It assumes
- Most operations touch a single partition key.
- Load is spread across many keys, so no single key dominates.
- There is a mechanism to move partitions and route correctly during the move.
Explain it in your own words
Where you practise it
- A URL shortener like bit.ly
- A reliable video processing pipeline
- Rate limiting a public API
- Product analytics over billions of events
Where the events live · Two visitors who were one person · One customer is a third of the traffic
- Storing trillions of chat messages
What the numbers say · Choose the store · Choose the partition key · Everyone opens the same channel · Defend the design
- A real-time collaborative editor
Where does a document's live state live? · The all-hands document
- Sharding Postgres while it is running
What is actually running out? · Choose the shard key · How many shards? · From 32 hosts to 96 · Queries that cross shards · Why not a distributed database?
Further reading
Engineers describing it in systems they run.
- The Great Re-shard: adding Postgres capacity (again) with zero downtime
Notion · Arka Ganguli and others · Post, Jul 2023
The payoff of fixed logical shards: growing to more machines by moving shards whole.
- Herding elephants: lessons learned from sharding Postgres at Notion
Notion · Garrett Fidalgo · Post, Oct 2021
Picking a shard key and a shard count, then moving a live database across with no data lost. Includes an honest list of what they would do differently.
- How Discord Stores Billions of Messages
Discord · Stanislav Vishnevskiy · Post, Jan 2017
Choosing a database and a partition key for chat history, and the surprises that followed: tombstones, and an edit racing a delete.
Related concepts
- Ordering
There is no global 'now' in a distributed system. Order exists only where something assigns it, so decide which order you need and who assigns it.
- Leases and fencing tokens
Ownership that expires unless renewed, plus a token that lets the rest of the system reject an owner that has lost its claim without knowing it.
- Caching
Keeping a copy of data closer to where it is used, trading freshness and complexity for speed and reduced load on the source.
- Backpressure and capacity
When work arrives faster than it can be done, something has to give: the queue grows, the producer slows, or work is shed. Choose which on purpose.
- Persistent connections
Long-lived connections such as WebSockets turn a stateless request tier into one that holds per-client state, with consequences for routing, deploys and failure detection.
- Online data migrations
Moving live data to a new schema or store without downtime: write to both, backfill the past, verify, switch reads, then switch writes, with a way back at every step.
- Log-structured storage (LSM trees)
Storage engines that turn every write into a sequential append and merge files in the background: very fast writes, at the cost of compaction, tombstones and more expensive reads.
- Columnar storage
Storing each column of a table separately, so analytical queries read only the columns they use and compress them well, at the cost of slow single-row lookups and updates.
- Generating unique identifiers
Making ids that are unique across machines and time, and choosing what else they reveal: order, volume, guessability, length.
- Replication
Keeping copies of data on several machines for durability, read capacity and locality, and living with copies that briefly disagree.
- Consistent hashing
Mapping keys to nodes so that adding or removing a node moves only a small share of keys, instead of reshuffling almost all of them.
- Rate limiting
Capping how fast a client may use a resource, to protect capacity, enforce fairness, and stay within the limits of the systems you depend on.
- Fan-out on write and fan-out on read
When one write must reach many readers, do the work when it is written (precompute every reader's view) or when it is read (assemble it on demand). Most real feeds do both.