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.

Advertisement

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

Whole-stack sharding: every stateful component lives inside a shardClientstenant in token or hostEdge routerstateless, caches directoryTenant directorytenant to shard, versionShard 7 (one of N identical stacks)App serversWorkersCacheQueueDatabaseBlob prefixShard 8, shard 9, ...same templates, different tenantsGlobal servicesidentity, billing, directory, analytics copyMovercopies one tenant's whole footprintlookupto new shardflip entryA shard is a vertical slice: losing one affects only its tenants, and a deploy can stop at one shard.The edge needs one lookup per request; nothing below it ever talks to another shard.
An edge router looks up each tenant in a directory and forwards to one shard. Inside the shard, every stateful component is local. A small set of global services sits outside all shards.

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.

Advertisement

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.

ComponentWhere it livesWhy
Relational dataInside the shardLocal transactions per tenant
CacheInside the shardA cache stampede or eviction storm stays local
Queues and workersInside the shardOne tenant's backlog cannot starve others
Scheduled jobsPer shard, staggeredAvoids every shard hitting shared services at midnight
Blob storagePer-shard bucket or tenant prefixMakes moves a copy of one prefix
Search indexInside the shardRebuilds and mapping changes are shard-local
Identity, sign-up, billing totalsGlobalMust 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

NeedPatternCost
Report across all tenantsReplicate shard data into an analytics store; query thatMinutes of lag; another pipeline to run
Operator search across tenantsScatter-gather with a deadline and partial resultsLoad on every shard; tail latency
Globally unique value (an email)A global uniqueness service keyed by that valueExtra round trip at sign-up
Action involving two tenantsAsynchronous messages between shards via a global bus, or a sagaEventual 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&#x27;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

SituationBetter first step
One database is busy but has room on a larger instanceScale up, add replicas, fix queries
No natural owner for most requestsA distributed database with global transactions
Most features compare data across all customersKeep a shared stack; isolate only heavy jobs
A few tenants, strict isolation needsDedicated deployments per tenant
Many tenants, a clear owner per request, growing fastWhole-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

  1. Pick the shard unit and check your top twenty request types against it.
  2. List every stateful component and decide whether it lives in the shard or is global; minimise the global list.
  3. Build the tenant directory with compare-and-set, the edge router with a short cache, and a 421 response on every shard.
  4. Put tenant ids in every cache key, job payload and blob path.
  5. Write and rehearse the whole-tenant move, including queued jobs, blobs and the way back.
  6. Stamp shards from one template, deploy in waves, and alert on the worst shard rather than the fleet average.
Key takeaway: A sharded architecture partitions the whole stack, not only the database: each shard is an identical, independent slice with its own app servers, workers, cache, queue, database and storage, serving a set of tenants that an edge router finds through a small directory. Place tenants by capacity, make every stateful component shard-aware, keep cross-shard work rare and explicit, move whole tenants with a copy, catch-up, brief pause and flip, and run the fleet shard by shard.