Skip to content

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

0 of 1 checks done
  1. 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.
  2. Check

    All writes for document 42 go to one owner process. What does that buy?

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

Write at least 60 characters (0 so far). Write it as you would say it in a design review. You will compare it against the points a strong answer makes.

Where you practise it

Further reading

Engineers describing it in systems they run.

  • 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.