Introduction
Imagine a service that stores user profiles, preferences, and account settings. One database has reached its practical capacity, so the service distributes users across several database shards. Given a user ID, every application instance must agree on which shard owns that user’s data.
Routing is straightforward while the shard count stays fixed. Adding a shard is where the design becomes interesting: changing the destination of a request does not move the rows it needs.
Consistent hashing helps limit how many users need to move. A migration protocol makes those moves safe.
Why modulo becomes expensive
A simple routing rule is:
shard = hash(userId) % shardCountConsider these illustrative hash values, with shard IDs numbered from zero:
| User hash | Three shards | Four shards |
|---|---|---|
| 12 | 0 | 0 |
| 13 | 1 | 1 |
| 14 | 2 | 2 |
| 15 | 0 | 3 |
| 16 | 1 | 0 |
| 17 | 2 | 1 |
| 18 | 0 | 2 |
| 19 | 1 | 3 |
Most of these users change destination when the fourth shard arrives. Over uniformly distributed hash values, changing directly from three to four buckets moves roughly three quarters of the keys. This follows from comparing the remainders across a complete cycle of twelve values: only three keep the same bucket number.
If the application switches to the new rule before migrating data, requests go to shards where those users do not exist yet. If it migrates everything first, it still has to handle writes arriving during the copy.
Modulo is not inherently wrong. The difficulty is tying its divisor directly to the number of physical databases, which is precisely the number we expect to change.
Put shard positions on a ring
Consistent hashing places shard positions and user hashes in the same circular space. A user belongs to the first shard position at or after its hash, wrapping to the beginning when necessary. This ring-based partitioning is described in the original Dynamo paper, section 4.2 .
For a small example, use positions from 0 through 99, with shards A at 20, B at 50, and C at 80. Follow increasing positions clockwise; after 99, wrap back to 0.
The positions below are chosen for explanation, not calculated from real user IDs.
| User | Hash position | Owner |
|---|---|---|
| Alice | 10 | A |
| Ben | 35 | B |
| Carmen | 65 | C |
| Dev | 90 | A |
Now add shard D at position 40. The ring becomes:
block-beta
columns 3
space A(("A · 20")) space
C(("C · 80")) space D(("D · 40 · new"))
space B(("B · 50")) space
A --> D
D --> B
B --> C
C -- "99 → 0: wrap" --> ARead the arrows clockwise, starting at A at the top. Positions are schematic: the gaps are not proportional to the hash ranges. Without D, the arrow from A went directly to B; D now intercepts hashes greater than 20 and at most 40.
D takes ownership of the interval (20, 40]. Ben moves from B to D. The other three users keep their owners. A user exactly at position 40 also belongs to D; a user exactly at 20 stays on A.
Only the interval captured by the new position changes ownership. With balanced placement and uniformly distributed keys, adding an equal-capacity shard to N shards moves an expected fraction of about 1 / (N + 1) of the keys. That is an expectation, not a promise about a particular ring’s distribution or the number of bytes transferred.
A small routing function
The lookup itself can use a binary search over sorted positions:
type Token = { position: bigint; shardId: string;};
function ownerForHash(hash: bigint, ring: readonly Token[]): string { if (ring.length === 0) throw new Error("The shard ring is empty");
let low = 0; let high = ring.length;
while (low < high) { const middle = low + Math.floor((high - low) / 2); const token = ring[middle]!;
if (token.position < hash) low = middle + 1; else high = middle; }
return ring[low % ring.length]!.shardId;}
const ring: Token[] = [ { position: 20n, shardId: "A" }, { position: 40n, shardId: "D" }, { position: 50n, shardId: "B" }, { position: 80n, shardId: "C" },];
ownerForHash(35n, ring); // DownerForHash(90n, ring); // A: wrap aroundThis is a routing example, not a complete database client. It assumes an immutable ring sorted by unique position, with hashes and tokens in the same nonnegative range. Ring construction must reject or resolve token collisions deterministically.
In a real service, specify the hash algorithm, seed, input encoding, and user-ID normalization as part of the routing protocol. Do not use a language’s default object hash: another process or language may produce a different result. Use stable logical shard IDs rather than hostnames, so replacing a database host does not implicitly change ownership.
Virtual nodes improve distribution
One position per shard can leave some shards owning much larger intervals than others. Virtual nodes give each physical shard several positions spread around the ring. Each position still points to the same underlying shard; it is not another database or another copy of the data. Dynamo uses this technique to improve distribution and accommodate differences in node capacity. Dynamo paper
A membership configuration could assign tokens for identities such as A:0, A:1, and A:2. Keep those identities stable when adding D. Regenerating every shard’s tokens would discard the placement stability we wanted.
More tokens can smooth out ownership across shards, but more uniform key counts do not necessarily mean more uniform load. Measure bytes, request rates, and latency as well as user counts.
- Adding a shard changes ownership for selected ranges rather than remapping most users.
- Virtual nodes spread each shard’s ownership across several ranges.
- Every router can calculate an owner from the same versioned ring.
Choose what moves with a user
For this design, userId is the placement key for a group of related rows:
user_profiles(user_id, ...)user_preferences(user_id, ...)user_addresses(user_id, address_id, ...)A request first resolves the user ID to a shard, then runs its queries there. Related updates can remain within a local database transaction when all the relevant rows live together.
An email address is usually a poor substitute for the placement key because it can change. A login request that starts with an email address needs another lookup to obtain the stable user ID. That lookup also needs a design for enforcing email uniqueness across the service; a unique index on each individual shard cannot enforce global uniqueness by itself.
The choice changes for tenant-heavy workloads. If most operations span every user in one organization, placing by tenantId may preserve more local transactions. The cost is that a large tenant can dominate one shard. Pick the grouping that matches the operations you need to perform together.
A new ring needs a migration protocol
Suppose version 7 of the ring assigns Ben to B, while version 8 assigns him to D. Updating the routers alone makes Ben’s profile appear missing. Copying his profile once is also insufficient: an update on B during the copy could be lost after cutover.
One possible migration design uses a snapshot, a change stream, and a brief write fence for the moving range:
- Prepare the destination. Create the required tables and indexes on D. Record the range, source, destination, and intended ownership version in durable migration metadata. B remains authoritative.
- Copy a consistent snapshot. Transfer all user-owned rows in the moving range. Establish change capture at a position coordinated with that snapshot, so writes and deletes during the copy can be replayed without gaps.
- Catch up. Apply changes to D in order, with retry-safe processing. Monitor lag and validate the copied data while B still handles reads and writes.
- Fence and cut over. Briefly block or queue writes for that range, drain the remaining changes, then commit the ownership transition. Old owners must reject or forward stale-version requests; a router refresh alone does not fence old writers.
- Retain and clean up. Keep B’s old copy for a defined recovery window, then remove it after confirming traffic and replication state. Once D accepts new writes, rolling back requires reconciling those writes; pointing routers back to B is not enough.
The exact mechanism depends on the database. A durable ownership service must serialize competing migrations, and the write path must enforce its decisions. Naive dual writes introduce another failure case: one database can accept the update while the other rejects it.
Membership versions belong in operational logs and migration state. They make it possible to distinguish a stale router from missing data or a failed copy.
Failure is different from redistribution
If B becomes unavailable, removing it from the ring immediately sends its users elsewhere. That does not make their records available on the new destination.
Treat replica failover separately from changing logical shard ownership. A logical shard can have a primary and replicas across failure domains. Promoting a suitable replica preserves the shard’s identity and placement; changing the ring initiates a different operation.
A busy user can still create a hot shard
Even a beautifully balanced ring cannot split one user’s workload if that user’s data must stay together. A celebrity account, automated client, or unusually large tenant may overwhelm its owner while other shards remain comfortable.
For append-heavy data such as activity events, a separate partitioning scheme can divide that workload into time buckets or subkeys. Reads then need to know which partitions to query. AWS’s write-sharding guide illustrates the same tradeoff: distributing writes across suffixes requires additional read coordination.
The profile and settings might remain colocated by user ID while events use that separate scheme. There is no requirement that every table in the service share one placement strategy.
Consider fixed logical buckets too
For an application-managed SQL deployment, another option is to hash users into a fixed set of logical buckets:
bucket = hash(userId) % 4096shard = placementTable[bucket]Here, 4096 is an illustrative bucket count that stays constant when physical shards are added. Scaling changes selected entries in the placement table, rather than the hash calculation for every user. Changing the bucket count later would reintroduce a repartitioning problem.
This design makes the units of migration explicit. A ring offers flexible token placement; a bucket map offers a bounded set of partitions to track. Both still require coordinated ownership changes and safe data transfer.
Before choosing either, check whether the database already manages partitioning and rebalancing. If it does, an additional application-level ring may add another source of routing state without solving a problem you actually have.
For a user database, the design is ready when it can answer more than “Which shard owns this ID?” It should also answer who owns the next write during a migration, how stale routers are handled, and how the system recovers when the destination fails halfway through the move.