Pattern · seen in 4 breakdowns across 4 companies
Application-Layer Sharding
Application-layer sharding puts the decision of sharding (which machine each piece of data lives on) into your application code, instead of the underlying storage system splitting the data across machines.
The mechanism
At its core: instead of letting the storage system decide where each customer's data lives, your application keeps a mapping from customer to machine. Owning that mapping means you can reroute, take clusters offline, and rebalance by editing it - without depending on the underlying storage system.
When you own the mapping, taking a cluster offline is one edit; when the storage owns it, all you can do is wait.
Definition
Big systems split their data across many machines so no single one has to hold it all - that splitting is called sharding. The usual way is to let the storage system decide where each piece goes. Application-layer sharding makes that decision in your own code instead. Your application knows which customer or account each piece of data belongs to, and uses that to send every read and write to the correct machine. The storage system just runs the operations it is handed; it does not decide where anything lives.
The reason to own the data sharding is control. When the storage system owns the sharding, you are stuck with its rules: how it scales, how it coordinates, and what operations it lets you run. When your application owns the sharding, it can route using things the storage system never knows: which customers are most important, how they use the system, where their data must legally sit. It can also change how data is split without waiting on the storage system's own migration tools, which are often the hardest part to work around once you outgrow the first design.
In practice, this is a mapping layer: a lookup from each customer to the machine that holds their data, kept in a fast cache in front of a slower, authoritative store (the source of truth). Reads check the cache; writes go to whatever machine the cache points to. The storage becomes a pool of interchangeable units - clusters, partitions, nodes - that your application controls. That splits two things that used to be tangled. Adding capacity just means putting another unit in the pool, and changing how data is split just means editing the mapping. Because your code now controls the sharding, it can also keep related data together on one unit - see shard-key colocation - so joins and transactions still work.
This is worth doing when you have outgrown one machine's worth of data and the storage's own way of splitting it has become the thing holding you back: too rigid to change, or blind to what your application knows. It is not free: your code now owns the mapping and everything that can go wrong with it. When the storage's built-in sharding still fits, let it do the job.
When it applies
Tradeoffs
The same move, 4 ways
Every row is a production system that bet on this pattern — the note says how, in that system's own terms.
Often used together
Patterns sharing breakdowns with this one — derived from co-occurrence, threshold ≥2 shared.
Problems this pattern answers
The walls where its breakdowns live — each opens the cross-company comparison.