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:
| Scheme | Lookup | State to store | Move one key or tenant | Typical weakness |
|---|---|---|---|---|
| Hash | shard = hash(key) mod N, or a hash ring | Almost none | Impossible without moving everything that hashes with it | Cannot place data deliberately; resharding moves many keys |
| Range | Binary search over boundaries | One row per range | Split or move a whole range | Hot ranges from sequential keys |
| Directory | Table lookup per entry | One row per entry | Edit one row after copying data | The 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 maps | Directory size | Placement freedom | When to use |
|---|---|---|---|
| One key (row, user, object) | As large as the data set itself | Total | Rarely; only when keys are few and valuable, such as large files or mailboxes |
| One tenant or account | Tenant count, usually thousands to millions | Per tenant | Multi-tenant products where a tenant is the unit of isolation |
| One virtual bucket | Fixed, for example 4,096 or 16,384 | Per bucket | Many 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.
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:
- Mark. Set the entry to MIGRATING with
target_shard = Cand bump its epoch. Traffic still goes to A. - 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.
- 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.
- 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. - 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
| Failure | Symptom | Cause | Prevention |
|---|---|---|---|
| Split brain on a tenant | Writes visible on one shard but not the other | Ownership changed without freezing the source | Freeze plus epoch check at the shard; single-write flip |
| Invisible new tenant | Sign-up succeeds, first request 404s on some nodes | Long-lived negative cache on routers | No or short negative TTL; invalidate on creation event |
| Herd after deploy | Directory saturated when routers restart together | Every router does a full load | Snapshots, staggered restarts |
| Orphaned rows | Disk use on the old shard never falls | Cleanup step skipped after a mover crash | Resumable move log; reconcile rows against the directory |
| Hidden cross-tenant joins | Queries break after a move | Code joined tenants that used to share a shard | Forbid cross-tenant queries; test with tenants split across shards |
| Directory drift | Directory says C, data is on A | Manual edits or partial restores | Only 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
- 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.
- Create the directory table with logical shard ids, state, target shard and per-entry epoch, plus a global version.
- Build the router cache with background sync, a bounded negative cache and retry on stale epoch.
- Make every shard check the epoch and refuse writes for FROZEN or foreign entries.
- Write the mover as a resumable state machine: mark, copy, catch up, freeze, flip, clean.
- Rehearse a move of a test tenant in production, measure the freeze window, and alert if it exceeds your budget.
- Add a nightly reconciliation that compares directory entries with where rows actually are.
- Kill the directory service in a test environment and confirm traffic keeps flowing from caches.