Most writing about sharding stops at the database: split a table by key and route queries. But the database is rarely the only thing that grows. Caches, queues, background workers, blob storage and search indexes grow with it, and if they stay shared, one noisy customer or one bad deploy still takes everyone down. A sharded architecture partitions the whole stack: each shard is a complete, independent copy of the service's stateful machinery, serving a subset of tenants or users, and only a thin edge layer knows which shard holds whom.
This article is about that system-level design. It assumes you know the key-mapping strategies; sharding strategies compared explains range, hash and directory schemes, and the sharding overview covers the database basics. Here we decide what lives inside a shard, route requests at the edge, place tenants by capacity, keep every component shard-aware, handle the work that crosses shards, and move a tenant's entire footprint while it keeps working.
What makes an architecture sharded
An architecture is sharded when the unit of scale is a whole stack rather than a single component. Shard 7 has its own app servers, workers, cache, queue, database and storage prefix; it serves a known set of tenants and never calls shard 8. Capacity grows by adding stacks. Failures are bounded: a corrupted cache, a runaway job or a bad migration affects one shard's tenants. Deploys can stop after one shard. The same idea appears under other names, such as cells, pods or stamps, with the shared point that the blast radius is a deliberate, sized boundary.
The price is that every feature must answer a question it never faced before: does it work inside one shard, or does it need data from many? The more features answer "one", the better the architecture works. That makes the choice of shard unit the most important decision in the design.
The architecture at a glance
Four kinds of component carry the design. The edge router is stateless: it extracts the tenant from the token or host name, looks it up in a cached directory and forwards. The tenant directory is a small, strongly consistent table mapping tenant to shard, plus a state such as active or moving, changed only by compare-and-set. Shards are identical stacks built from the same templates. Global services hold the few things that cannot be partitioned by tenant: sign-up and identity, billing aggregation, the directory itself and an analytics copy for cross-tenant reporting.
Choosing the shard unit
The shard unit should be the entity that almost every request names and that owns data which changes together. For business software it is nearly always the tenant (the customer organisation): its users, documents, invoices and jobs all live in one shard, and a transaction across them stays local. For consumer software it might be the user or a household; for an IoT platform, a fleet.
Test a candidate against your top twenty request types: each should resolve to one shard, or be rare enough to tolerate fan-out. Then test it against skew. Tenant sizes follow a long-tailed distribution, and the largest tenant may need a shard to itself or more than one shard can provide. Decide now what happens then: a dedicated shard is easy; splitting one tenant across shards reintroduces everything this design avoids. The hot key mitigation article covers the per-key techniques for the cases in between.
Placement and routing
Whole-stack shards favour a directory over a hash. Tenants differ in size by orders of magnitude, shards differ in age and hardware, and you want to choose where each tenant goes, move it later and pin large tenants. A directory lookup costs one cached read at the edge, which is cheap. New tenants are placed by capacity: put each on the open shard with the most headroom after adding it, and close shards to new tenants as they fill so that organic growth of existing tenants has room.
import time
from dataclasses import dataclass
@dataclass
class Shard:
name: str
capacity: float # normalised units, e.g. from load tests of one stamped shard
load: float # current p95 demand in the same units
open: bool = True # closed shards take no new tenants
def place_tenant(tenant_id, expected_load, shards, directory, headroom=0.30):
"""Put a new tenant on the open shard with the most room left after adding it."""
candidates = [s for s in shards
if s.open and s.load + expected_load <= s.capacity * (1 - headroom)]
if not candidates:
raise RuntimeError("no shard has headroom: stamp a new shard first")
best = max(candidates, key=lambda s: s.capacity - s.load - expected_load)
directory.create(tenant_id, best.name) # fails if the tenant already exists
best.load += expected_load
return best.name
class Edge:
"""Route by tenant; refresh one entry when a shard says the tenant has moved."""
def __init__(self, directory, shard_clients, ttl=60):
self.directory, self.clients, self.ttl, self.cache = directory, shard_clients, ttl, {}
def lookup(self, tenant_id, force=False):
hit = self.cache.get(tenant_id)
if force or hit is None or hit[1] < time.time():
entry = self.directory.get(tenant_id) # {"shard": ..., "state": ...}
self.cache[tenant_id] = (entry, time.time() + self.ttl)
return entry
return hit[0]
def handle(self, tenant_id, request):
for attempt in range(4):
entry = self.lookup(tenant_id, force=attempt > 0)
if entry["state"] == "moving": # brief write freeze during cutover
if request.is_read:
return self.clients[entry["shard"]].send(request)
time.sleep(0.2 * (attempt + 1))
continue
response = self.clients[entry["shard"]].send(request)
if response.status != 421: # 421: tenant not on this shard
return response
return request.reply(503, retry_after=2)Two details in the edge code matter. Shards reply with a dedicated status (here HTTP 421, Misdirected Request) when asked about a tenant they do not hold, so a stale cache entry heals in one round trip instead of reading or writing the wrong shard. And while a tenant is in the brief moving state the edge serves reads from the source and holds writes with backoff, so a move appears as a short pause, not an error.
Everything stateful must be shard-aware
Sharding the database while leaving other state shared is the most common way this design fails. Walk every component and put it on one side of the line.
| Component | Where it lives | Why |
|---|---|---|
| Relational data | Inside the shard | Local transactions per tenant |
| Cache | Inside the shard | A cache stampede or eviction storm stays local |
| Queues and workers | Inside the shard | One tenant's backlog cannot starve others |
| Scheduled jobs | Per shard, staggered | Avoids every shard hitting shared services at midnight |
| Blob storage | Per-shard bucket or tenant prefix | Makes moves a copy of one prefix |
| Search index | Inside the shard | Rebuilds and mapping changes are shard-local |
| Identity, sign-up, billing totals | Global | Must see all tenants; keep them small and highly available |
Store every tenant's blobs under a tenant prefix even if you use one bucket, and include the tenant id in every cache key and job payload. Those habits make moves and audits mechanical rather than archaeological.
Work that spans shards
| Need | Pattern | Cost |
|---|---|---|
| Report across all tenants | Replicate shard data into an analytics store; query that | Minutes of lag; another pipeline to run |
| Operator search across tenants | Scatter-gather with a deadline and partial results | Load on every shard; tail latency |
| Globally unique value (an email) | A global uniqueness service keyed by that value | Extra round trip at sign-up |
| Action involving two tenants | Asynchronous messages between shards via a global bus, or a saga | Eventual consistency; compensations |
Scatter-gather is the pattern to keep rare. If each of 50 shards has a 1 percent chance of exceeding its 99th-percentile latency on a request, a 50-way fan-out hits at least one slow shard about 40 percent of the time (one minus 0.99 to the power 50). Serve dashboards from the analytics copy instead, and give fan-out endpoints their own rate limits. When two tenants must interact, such as a shared document between organisations, send a message rather than a distributed transaction; the saga article shows how to handle the failure paths.
Moving a tenant's whole footprint
Moves are how you rebalance, isolate a growing tenant or drain a shard for maintenance. Unlike moving rows between database shards, a whole-stack move must carry rows, blobs, queued jobs and in-flight work together, and the tenant must keep working. The procedure below does that with a write pause of seconds.
def move_tenant(t, src, dst, directory):
"""Each step is idempotent so a crashed move can be re-run from the top."""
dst.db.prepare_tenant(t)
pos = dst.db.copy_snapshot(src.db, t) # rows plus change-log position
dst.blobs.sync_prefix(src.blobs, f"tenants/{t}/") # bulk copy while still live
while src.db.pending_changes(t, pos) > 500:
pos = dst.db.apply_changes(src.db, t, pos)
directory.set_state(t, "moving") # edge pauses writes, reads continue
src.workers.pause_tenant(t) # stop picking this tenant's jobs
src.workers.wait_idle(t, timeout_s=60)
pos = dst.db.apply_changes(src.db, t, pos)
assert src.db.pending_changes(t, pos) == 0
dst.blobs.sync_prefix(src.blobs, f"tenants/{t}/") # final delta
dst.queue.import_jobs(src.queue.drain_tenant(t)) # queued work follows the tenant
verify(src, dst, t) # counts, checksums, blob manifests
directory.compare_and_set(t, expect=src.name, new=dst.name, state="active")
src.mark_moved(t) # src answers 421 for t from now on
dst.cache.warm(t) # optional: precompute hot keys
schedule_later(src.purge_tenant, t, delay_hours=72)Walk through one move. Tenant acme holds 60 GB of rows and 400 GB of blobs on shard 7, which is at 80 percent of capacity. The snapshot and blob copy run for several hours with throttling, while acme works normally. Change replay brings the database within 500 changes. The directory flips acme to moving: the edge holds writes, reads continue. Shard 7's workers stop taking acme jobs and finish the two in flight. The last changes and blob delta are copied, 1,200 queued jobs move to shard 12's queue, verification passes, and the directory's compare-and-set points acme at shard 12. Writes resume after roughly ten seconds. An edge node that still caches shard 7 gets a 421, refreshes and retries. Shard 7 keeps acme's data for 72 hours as a way back.
Rehearse moves on test tenants every week. A move procedure that runs only in emergencies will fail in one.
Stamping and operating shards
- Stamp, do not build. A new shard is created from the same infrastructure-as-code template as every other, load tested, then opened for placement. If creating a shard takes a ticket and a week, you will overfill the existing ones.
- Deploy in waves. One canary shard of internal tenants, then a few percent, then the rest, with automatic halt on error-rate regression in the shards already done.
- Version skew is normal. Shards run different code versions during a rollout, so schema and message changes must be backward compatible (expand, migrate, contract).
- Watch the worst shard. Fleet averages hide one shard at 95 percent disk. Alert on the maximum across shards and keep per-shard SLO dashboards.
- Keep shard sizes bounded. Fix a target capacity per shard and add shards rather than growing them; uniform shards make capacity planning and restore times predictable.
Failure modes
- A shared component nobody noticed. A global Redis or job queue quietly couples all shards, and its outage is total. Audit for it with the table above.
- Writes to the wrong shard. Without the 421 check, a stale directory cache writes a moved tenant's data to the old shard, where it is lost at purge.
- A tenant outgrows a shard. Detect it from per-tenant load trends long before, and move it to a dedicated shard while it still fits the move window.
- Directory outage. Edge nodes must keep serving from cache; only sign-ups and moves should depend on the directory being up.
- Synchronized jobs. Every shard runs the nightly job at 00:00 and floods a global service. Stagger schedules per shard.
- Half-finished moves. A crash leaves a tenant in the moving state. Alert on any tenant in that state for more than a few minutes, and make the mover resumable.
When not to shard the stack
| Situation | Better first step |
|---|---|
| One database is busy but has room on a larger instance | Scale up, add replicas, fix queries |
| No natural owner for most requests | A distributed database with global transactions |
| Most features compare data across all customers | Keep a shared stack; isolate only heavy jobs |
| A few tenants, strict isolation needs | Dedicated deployments per tenant |
| Many tenants, a clear owner per request, growing fast | Whole-stack sharding |
Whole-stack sharding buys bounded blast radius and near-linear capacity and charges for them in fan-out complexity and operational discipline forever. Adopt it when the shard unit is natural, and build the mover and the stamping template before you need them.
What to do next
- Pick the shard unit and check your top twenty request types against it.
- List every stateful component and decide whether it lives in the shard or is global; minimise the global list.
- Build the tenant directory with compare-and-set, the edge router with a short cache, and a 421 response on every shard.
- Put tenant ids in every cache key, job payload and blob path.
- Write and rehearse the whole-tenant move, including queued jobs, blobs and the way back.
- Stamp shards from one template, deploy in waves, and alert on the worst shard rather than the fleet average.