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.

LIVE ARTIFACTYOU CONTROL THE SHARDINGOPEN FULL SCREEN ↗
THE IDEANormally the storage system decides the sharding logic (which machine holds which data), and your application just asks it. Application-layer sharding moves that decision into your code: you keep a mapping from each customer to the machine that holds their data, and every read and write goes through it. Now the storage is just a pool of interchangeable machines you control. You can take a machine offline for maintenance, isolate a heavy customer, or add capacity by editing the mapping - none of which you could do on your own schedule when the storage owned the sharding. The price is that the mapping is now yours to keep correct, and nothing routes if the source of truth behind it is down.
WHAT TO TRYTry to take a cluster offline while the storage system owns the sharding - you cannot, you just wait. Then switch to application-layer routing and take it offline instantly by rerouting its customers in the mapping.

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.

WHERE ROUTING LIVES
The application routes through a mapping layer to a pool of independent storage units it controls.
The application asks a mapping layer where each customer's data lives - a fast cache backed by a source of truth - then talks to that storage unit directly; the storage never decides the sharding.

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

01Your customers differ a lot in size and behavior. When one customer is huge and another tiny, your application can place them far more sensibly than a generic hash that treats them all the same.
02You will need to change how data is split later. If the storage system's built-in sharding cannot be reworked without big downtime or a painful migration, owning the split yourself keeps that change in your hands.
03Keeping certain customers together or apart matters more than even load. Routing in your own code lets you group customers so one failure hits only that group, or separate a heavy customer so it cannot slow the others down.
04You want to run maintenance on your own schedule. With routing in your code, you can take a cluster offline for an upgrade or replacement by editing the mapping, instead of waiting for the storage system to rebalance on its own timeline.
05Routing needs facts the storage system does not have. Things like a customer's priority tier, a legal rule about where their data must sit, or per-customer settings live in your application, not in the storage engine.

Tradeoffs

You now own the mapping, and its bugs. Every read and write has to consult it, and every change has to be rolled out carefully. That brings a new class of bugs - wrong or stale routing - that did not exist when the storage handled it.
Simple operations get more complicated. What was one line - 'write this row' - now means look up the destination, route to the correct cluster, handle a failed lookup, and fall back to the source of truth. And all of that has to be maintained as the system grows.
You give up some built-in storage features. Cross-shard queries, automatic rebalancing, and global indexes usually assume the storage controls the split, so you either do without them or rebuild them yourself, which is real extra work.
The mapping becomes something you cannot route without. If the source of truth behind it goes down, your application cannot decide where anything goes. The cache covers popular lookups, but a new customer or a cache miss still needs the source of truth.

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.

Discord
Discord Engineering
2025
Discord's foundational decision in 2017, kept through the 2025 redesign. The routing logic lives in application code, not inside Elasticsearch. This gives Discord control over which clusters and indices each message goes to, control that proved essential when the time came to evolve the architecture. The value is clearest in hindsight: the redesign was possible because sharding was already an application-level concern the team could rework, without waiting on Elasticsearch's internal coordination. Read the breakdown →
Figma
Figma Blog
2024
Figma shards its primary transactional store through DBProxy's query engine (shard keys scrambled with a hash, related tables grouped so common joins and transactions keep working), where Discord shards a search-index fleet. Same idea, different storage underneath: the application layer owns where data goes and how queries are routed, so the database itself can stay stock. Read the breakdown →
Notion
Notion Blog
2021
Notion names the approach outright: build your own partitioning scheme and route queries from application code, chosen over packaged tools like Citus and Vitess because their clustering logic is opaque and Notion wanted control over how its data was distributed. Three companies in the library have now made this same call over different starting points: Discord moving off Elasticsearch, Figma adding a query-routing service in front of Postgres, and Notion routing directly in application code. The shared reason is the same each time: owning the routing is worth the ownership cost when the packaged option hides how your data is placed. Read the breakdown →
Pinterest
Pinterest Engineering Blog
2015
This is the earliest instance of the pattern, designed in 2012, years before the other systems here hit the same wall. Discord routes messages to Elasticsearch shards in application code (2017); Figma and Notion partitioned an existing single Postgres database under duress; Pinterest built the virtual-shards-over-machines architecture from scratch mid-hypergrowth. The shared skeleton is exact: many small logical shards (4,096 here), a config mapping ranges to machines, and capacity added by remapping. So are the shared renunciations: no cross-shard joins, foreign keys, or global indexes. Read the breakdown →

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.