Every sharded system has to answer one question on every request: which shard holds this key? Hash sharding answers it with arithmetic, range sharding with a sorted list of boundaries. Directory-based sharding answers it by looking the key up in a table that records, entry by entry, where each piece of data lives.

That one change buys a lot of freedom. You can move a single customer without touching anyone else, put a noisy tenant on its own hardware, keep a regulated customer in one region, and rebalance by editing rows instead of rehashing data. It also creates a new component that every request depends on, that must be faster than the database behind it, and that must never be wrong about ownership. This page explains how to design that directory: how fine its entries should be, what each entry holds, how routers cache it safely, how to keep it available, and how to move one tenant end to end. The generic router and shard-map machinery is covered in database sharding, in depth; here the directory itself is the subject.

What a directory buys you

Compare the three ways of mapping a key to a shard by what they cost and what they let you do:

SchemeLookupState to storeMove one key or tenantTypical weakness
Hashshard = hash(key) mod N, or a hash ringAlmost noneImpossible without moving everything that hashes with itCannot place data deliberately; resharding moves many keys
RangeBinary search over boundariesOne row per rangeSplit or move a whole rangeHot ranges from sequential keys
DirectoryTable lookup per entryOne row per entryEdit one row after copying dataThe directory is on every request path

A directory is the only scheme that lets placement be a policy decision rather than a consequence of the key. That matters most in multi-tenant SaaS, where tenants differ in size by four or five orders of magnitude, where contracts can demand a region or dedicated hardware, and where the biggest customer will eventually outgrow any shard it shares.

Choosing what one entry maps

The most important design choice is what one directory entry maps. There are three common answers, and they behave very differently:

Entry mapsDirectory sizePlacement freedomWhen to use
One key (row, user, object)As large as the data set itselfTotalRarely; only when keys are few and valuable, such as large files or mailboxes
One tenant or accountTenant count, usually thousands to millionsPer tenantMulti-tenant products where a tenant is the unit of isolation
One virtual bucketFixed, for example 4,096 or 16,384Per bucketMany small entities; bounded directory size

Per-key directories look attractive and rarely survive: the directory grows as fast as the data, every insert needs a directory write, and the cache hit rate falls because there is no locality. Per-tenant directories are the sweet spot for SaaS because the tenant is already how you bill, isolate and delete. Virtual buckets are a hybrid: hash the key to one of a fixed number of buckets, then look the bucket up in a small directory. The directory stays a few thousand rows that every router can hold entirely in memory, and you rebalance by moving buckets.

Many systems combine the last two: most tenants live in hashed virtual buckets, and a short override table pins the largest or most regulated tenants to named shards. The router checks the override table first, then falls back to the bucket.

The directory data model

A directory entry is more than a shard name. It needs enough state to make moves safe and to let routers detect stale caches:

CREATE TABLE shard_directory (
    entry_id      VARCHAR(64) PRIMARY KEY,   -- tenant id, or 'bucket:1234'
    shard_id      VARCHAR(32) NOT NULL,      -- logical shard, resolved to hosts elsewhere
    state         VARCHAR(16) NOT NULL,      -- ACTIVE, MIGRATING, FROZEN, DELETED
    target_shard  VARCHAR(32),               -- set only while MIGRATING
    epoch         BIGINT      NOT NULL,      -- bumped on every change to this row
    updated_at    TIMESTAMP   NOT NULL
);

CREATE TABLE directory_version (
    id            INT PRIMARY KEY CHECK (id = 1),
    version       BIGINT NOT NULL            -- bumped on any change to any row
);

Two indirections keep this table stable. First, shard_id is logical; a separate topology record maps it to the current primary and replicas, so a database failover does not touch the directory at all. Second, the per-row epoch lets a shard reject a request routed with an old view of that one entry, and the global version lets routers ask cheaply whether anything changed since their last sync. Store the directory in a strongly consistent replicated database or a coordination service such as etcd or ZooKeeper; a directory that can return two different owners for the same entry is worse than having none.

Resolving a request: caches, epochs and negative caching

Routers must never read the directory synchronously on the hot path; that would make the directory's latency and availability a floor under every request. Instead each router holds the directory (or the hot part of it) in memory and keeps it fresh in the background, either by watching for changes or by polling the global version every second or so and fetching changed rows. The request path is then a hash-map lookup.

Clientrequest with tenant idRouterin-memory directory cacheDirectory servicereplicated, versionedShard Aowns epoch-checkedShard BentriesShard CMovercopy, catch up, flip, clean1 request2 sync or watchversion + rows3 route with epochflip entrycopywrong epoch: shard rejects, router refreshes
Directory-based sharding: routers resolve from a cached copy, shards enforce per-entry epochs, and the mover changes one entry at a time.
class DirectoryRouter:
    def __init__(self, directory, topology):
        self.directory = directory          # client for the directory service
        self.topology = topology            # shard_id -> current primary address
        self.cache = {}                     # entry_id -> (shard_id, epoch, state)
        self.version = 0

    def sync(self):                         # runs in the background every ~1 s
        v = self.directory.current_version()
        if v != self.version:
            for row in self.directory.rows_changed_since(self.version):
                self.cache[row.entry_id] = (row.shard_id, row.epoch, row.state)
            self.version = v

    def route(self, tenant_id, request, retries=2):
        for _ in range(retries + 1):
            entry = self.cache.get(tenant_id)
            if entry is None:               # miss: one read, never cache a guess
                row = self.directory.get(tenant_id)
                if row is None:
                    raise UnknownTenant(tenant_id)
                entry = (row.shard_id, row.epoch, row.state)
                self.cache[tenant_id] = entry
            shard_id, epoch, state = entry
            if state == "FROZEN" and request.is_write:
                raise RetryLater(tenant_id)  # brief write freeze during a flip
            try:
                return self.topology.primary(shard_id).execute(request, tenant_id, epoch)
            except StaleEpoch:
                self.cache.pop(tenant_id, None)   # refresh this entry and retry
        raise RoutingFailed(tenant_id)

Three details in that code carry the design. A cache miss reads the directory once and caches the answer, but an unknown tenant is not cached as a negative entry for long, because a tenant created a second ago would otherwise be invisible to this router; if you add negative caching to absorb abusive lookups, give it a TTL of seconds and invalidate it on the creation event. The shard checks the epoch on every request and returns StaleEpoch when the router's view is older than the entry it holds, so a lagging cache costs one retry, never a misplaced write. And the frozen state turns a move into a short, explicit retry window rather than a race.

The directory as a service to scale and protect

The directory is now a service in its own right, so size and protect it like one.

Size. A per-tenant entry is around a hundred bytes in memory with its key and indexes. A million tenants is therefore on the order of 100 MB per router: affordable, but large enough that a cold router start should load a snapshot rather than issue a million reads. With virtual buckets the whole directory is a few hundred kilobytes and the question disappears.

Read load. With caching, directory reads are proportional to changes and router restarts, not to traffic. Watch for the opposite: a deploy that restarts every router at once turns into a thundering herd of full loads. Stagger restarts and load from snapshots.

Availability. If the directory service is down, routers keep serving from their caches; only new tenants and moves stop. That is the right failure behaviour and you should test it deliberately. What you must never do is let a router with no cache guess an owner. A router that cannot load the directory should fail its readiness check and receive no traffic.

Write path. Writes to the directory happen on tenant creation, moves and deletions. Tenant creation is the busiest; allocate new tenants to a shard chosen by a placement function (least loaded, or the tenant's region) inside the same transaction that writes the entry, so there is no window where a tenant exists without an owner.

Worked example: moving one tenant

Moving one tenant is where a directory pays for itself. Suppose tenant acme lives on shard A, has 40 GB of data and needs to move to shard C because A is at 80% disk. The steps are:

  1. Mark. Set the entry to MIGRATING with target_shard = C and bump its epoch. Traffic still goes to A.
  2. Bulk copy. Copy acme's rows from A to C from a consistent snapshot, recording the change-log position (binlog or WAL) at which the snapshot was taken. This takes most of the time and runs throttled so A's other tenants do not notice.
  3. Catch up. Replay changes for acme from the recorded position until C is seconds behind A. Filter the change stream by tenant id; this is why every table needs the tenant id column.
  4. Freeze and flip. Set the entry to FROZEN and bump the epoch, so A rejects acme's writes with the old epoch and routers return retry-later. Apply the last few changes to C, verify row counts and checksums per table, then write shard_id = C, state ACTIVE and a new epoch in one transaction. The freeze lasts as long as the final catch-up, typically well under a few seconds.
  5. Drain and clean. Routers converge on the new entry within one sync interval; any request that reaches A with the old epoch is rejected and retried against C. After a safety period, delete acme's rows from A.

Two rules keep this safe. Ownership changes only by a single directory write, never by two independent updates. And the source shard must refuse writes from the moment of the freeze, because that refusal, not the router cache, is what prevents two shards accepting writes for the same tenant. Keep the move resumable: record each step in the directory row or a move log, so a crashed mover restarts from the last completed step instead of from scratch.

Isolating a hot tenant

A directory is also the cleanest tool for the tenant that is too big or too busy for a shared shard. Because the entry can point anywhere, you can give that tenant a dedicated shard, a larger instance class or a different region by running the move above. Pair it with per-tenant metrics on each shard (queries, CPU time, bytes written) so the decision is driven by measurement, and keep a short pinned list so the rebalancer never moves those tenants back automatically.

When even a dedicated shard is too small, the directory can map one tenant to several shards with a secondary key inside the tenant, but that turns single-tenant queries into scatter-gather and is a much bigger change. For hot individual keys inside a tenant, caching and key splitting work better than resharding; see hot key mitigation.

Failure modes

FailureSymptomCausePrevention
Split brain on a tenantWrites visible on one shard but not the otherOwnership changed without freezing the sourceFreeze plus epoch check at the shard; single-write flip
Invisible new tenantSign-up succeeds, first request 404s on some nodesLong-lived negative cache on routersNo or short negative TTL; invalidate on creation event
Herd after deployDirectory saturated when routers restart togetherEvery router does a full loadSnapshots, staggered restarts
Orphaned rowsDisk use on the old shard never fallsCleanup step skipped after a mover crashResumable move log; reconcile rows against the directory
Hidden cross-tenant joinsQueries break after a moveCode joined tenants that used to share a shardForbid cross-tenant queries; test with tenants split across shards
Directory driftDirectory says C, data is on AManual edits or partial restoresOnly the mover writes entries; periodic reconciliation job

Trade-offs and when to choose something else

Directory-based sharding trades a stateless calculation for a stateful service. You gain per-entry placement, cheap rebalancing, compliance pinning and the ability to isolate one customer. You pay with a component on every request path (mitigated by caching), a consistency protocol for moves, and a reconciliation job that keeps the directory and the data in agreement.

Choose it when the unit of data is a tenant or another entity you need to place deliberately and when sizes are badly skewed. Choose hashing when keys are uniform and you never need to move one on purpose. Choose ranges when you need ordered scans across keys. The virtual-bucket hybrid is a good default when you are unsure, because it starts as hashing and becomes a directory the day you need one. Whatever you choose, every other stateful piece of the stack must know about it; sharded architecture covers making caches, queues and jobs shard-aware, and caching covers the cache layer in front of the shards.

What to do next

  1. Write down the unit of placement (tenant, bucket or key) and the reason; if you cannot name a reason to place anything deliberately, use hashing.
  2. Create the directory table with logical shard ids, state, target shard and per-entry epoch, plus a global version.
  3. Build the router cache with background sync, a bounded negative cache and retry on stale epoch.
  4. Make every shard check the epoch and refuse writes for FROZEN or foreign entries.
  5. Write the mover as a resumable state machine: mark, copy, catch up, freeze, flip, clean.
  6. Rehearse a move of a test tenant in production, measure the freeze window, and alert if it exceeds your budget.
  7. Add a nightly reconciliation that compares directory entries with where rows actually are.
  8. Kill the directory service in a test environment and confirm traffic keeps flowing from caches.
Key takeaway: Directory-based sharding stores where each entry lives instead of computing it, which lets you place, move and isolate tenants one at a time. Pick a coarse entry (tenant or virtual bucket), give each entry a logical shard, a state and an epoch, cache the directory in every router and refresh it in the background, and make shards reject requests with a stale epoch. Move tenants with a resumable copy, catch-up, freeze and single-write flip, and reconcile the directory against the data regularly.