Choosing a shard key is the decision everyone talks about. Living with a sharded database is mostly about everything around that decision: how a query finds the one shard that owns its row, how every router agrees on which shard that is, how you answer the queries that do not carry the shard key, and above all how you move data between shards while the system keeps serving traffic. Get the key right and this machinery wrong, and you still have outages: writes landing on the wrong shard after a split, or a resharding job run at 2 a.m. in maintenance mode.

This article covers that machinery. Key selection, data skew and the cost of cross-shard queries are covered in the companion sharding architecture article and are assumed here. The running example is an orders database for a marketplace, sharded by customer_id across four shards, which has to grow to eight without downtime. The ideas map directly onto real systems such as Vitess (a MySQL sharding layer), Citus (PostgreSQL), and MongoDB's sharded clusters, which all solve the same three problems with different names.

Advertisement

The three layers: shards, the shard map and the router

A sharded database has three distinct parts, and keeping them distinct is most of the design. The shards are ordinary databases, each holding a subset of rows, each usually a primary with replicas. The shard map is a small piece of metadata saying which subset lives where. The router reads the shard map and sends every query to the right shard or shards.

The trick that makes resharding tractable is an indirection between the key and the shard. You do not map customer_id directly to a shard number. You map it to a keyspace id, typically a fixed-width hash such as the first 8 bytes of a hash of the key, and the shard map assigns contiguous ranges of keyspace ids to shards. Hashing spreads keys evenly; ranges make splits cheap, because splitting shard C means cutting its range in two, and only rows whose keyspace id falls in the moved half have to travel. Vitess calls the key-to-keyspace-id function a vindex; MongoDB calls the ranges chunks; Citus assigns hash ranges to shard tables. The concept is the same.

Three layers of a sharded database: router, shard map, shardsApplicationquery with shard keyRouterlibrary, proxy or coordinatorShard mapversioned, strongly consistentSQLwatchmap v42key -> hash -> keyspace id -> range -> shardShard Arange 00-3fShard Brange 40-7fShard C (splitting)range 80-bfShard Drange c0-ffsingle-shardC180-9fC2a0-bfcopy + catch upLookup tableemail -> keyspace idnon-key queryEvery shard enforces its own range: a write for a key it does not own is rejected,so a router holding a stale map fails loudly and refreshes instead of corrupting data.Each shard is a primary plus replicas; the router also picks primary or replica per query.
The router consults a versioned shard map, hashes the key to a keyspace id and picks the shard whose range contains it. Shard C is mid-split into C1 and C2; a lookup table serves queries that lack the shard key.

Compare the naive shard = hash(key) % N: changing N from 4 to 8 moves about half of all rows, and you cannot split just the one shard that is full.

Mapping keys: a router in forty lines

The core of any router is small. The version below hashes the key, finds the owning range with a binary search over sorted range starts, and remembers which map version it used so that a shard can reject it if the map has moved on.

import bisect, hashlib
from dataclasses import dataclass

@dataclass(frozen=True)
class ShardMap:
    version: int
    starts: list      # sorted range starts, e.g. [0x00.., 0x40.., 0x80.., 0xc0..]
    shards: list      # shards[i] owns [starts[i], starts[i+1])

def keyspace_id(customer_id: int) -> int:
    digest = hashlib.blake2b(customer_id.to_bytes(8, "big"), digest_size=8).digest()
    return int.from_bytes(digest, "big")      # 64-bit, uniformly spread

class Router:
    def __init__(self, map_store):
        self.map_store = map_store
        self.map = map_store.current()

    def shard_for(self, customer_id: int):
        kid = keyspace_id(customer_id)
        i = bisect.bisect_right(self.map.starts, kid) - 1
        return self.map.shards[i], kid

    def execute(self, customer_id, sql, params):
        for attempt in range(3):
            shard, kid = self.shard_for(customer_id)
            try:
                return shard.run(sql, params, keyspace_id=kid, map_version=self.map.version)
            except WrongShardError:              # shard no longer owns kid
                self.map = self.map_store.current()   # refresh and retry
        raise RoutingError("map kept changing; giving up")

Two details matter more than they look. The keyspace id is passed to the shard with every write, so the shard can check ownership itself. And the router treats WrongShardError as a routine signal to refresh, not as a failure. That pair is what makes a stale cached map safe.

Advertisement

Where the router lives

The same logic can run in three places, and the choice shapes operations for years.

PlacementExamplesStrengthsCosts
Client library in each serviceCustom libraries, many in-house ORMsNo extra hop; no proxy to scaleEvery language needs a library; upgrading routing logic means redeploying every service
Proxy tierVitess vtgate, ProxySQL-style proxies, mongosApplications speak plain SQL or the native protocol; one place to upgradeExtra network hop; proxy fleet must be scaled and made highly available
Coordinator inside the databaseCitus coordinatorDistributed planning, cross-shard joins and transactions handled by the engineCoordinator can become a hotspot; tied to one database engine

For most teams with more than one service language, a proxy tier wins: routing rules change without touching services. Run proxies stateless, behind a load balancer, with their own headroom, because every query passes through them. A client library suits a single latency-critical service that owns the data.

The shard map: versioned, watched and enforced

The shard map must be stored somewhere strongly consistent, typically etcd, ZooKeeper or a small replicated metadata database, because two routers disagreeing about ownership is how rows get written to the wrong place. Routers cache it and subscribe to changes; they do not read it per query.

Caching means some router will always be briefly stale, so correctness cannot depend on routers being up to date. Three mechanisms make staleness harmless:

  • Monotonic versions. Every change to the map bumps a version number. Routers attach the version they used; shards and the metadata store can reason about which writer saw which map.
  • Ownership enforced at the shard. Each shard knows the ranges it currently owns and rejects reads and writes outside them. This is the real safety net: a stale router gets an error, never a silent misplaced write.
  • Write freezes during moves. While a range changes owner, the source shard refuses writes for that range for a short window, so there is no moment when two shards both accept writes for the same key.

The rule: routers are allowed to be wrong, shards are not.

Queries without the shard key

Sharding by customer_id makes the order history page a single-shard query. But support staff search orders by order_id, and login looks customers up by email. There are three ways to serve those.

  • Scatter-gather. Send the query to every shard and merge. Fine for rare admin queries; terrible for anything hot, because latency becomes the slowest shard and load multiplies by the shard count.
  • Embed the shard key in the id. Make order_id carry the keyspace id or customer id in its high bits. The router can then route an order lookup without any extra read. This is the cheapest option when you control id generation.
  • Lookup tables. Keep a separate table, email -> keyspace_id, itself sharded by email. A login query does one lookup, then one routed query. Vitess offers this as a lookup vindex. The lookup table must be updated in the same logical write as the row, or reconciled asynchronously, which is a consistency decision you must make explicitly.
-- lookup table, sharded by email (its own key)
CREATE TABLE customer_by_email (
  email        VARCHAR(255) PRIMARY KEY,
  keyspace_id  BINARY(8)    NOT NULL
);

-- login: two single-shard hops instead of an N-way scatter
-- 1) SELECT keyspace_id FROM customer_by_email WHERE email = ?;
-- 2) SELECT * FROM customers WHERE customer_id = ?;   -- routed by keyspace_id

Global uniqueness of email is also enforced by this table: two shards cannot both insert the same email, because the lookup row lives on exactly one shard. Without it, a unique constraint on a non-key column is not enforceable at all in a sharded database.

Generating ids that route themselves

Auto-increment columns are per-shard, so two shards will hand out the same number. Common replacements are a central sequence service that leases blocks of ids, or a time-ordered 64-bit id composed of a timestamp, a generator id and a counter. For routing, the best variant reserves bits for the keyspace id prefix of the owning entity:

EPOCH_MS = 1_767_225_600_000   # 2026-01-01; keeps ids positive in a signed BIGINT

def new_order_id(customer_kid: int, now_ms: int, seq: int) -> int:
    # 1 sign bit (0) | 41 bits ms since EPOCH_MS | 10 bits keyspace prefix | 12 bits sequence
    prefix = customer_kid >> (64 - 10)          # top 10 bits of the 64-bit keyspace id
    return ((now_ms - EPOCH_MS) << 22) | (prefix << 12) | (seq & 0xFFF)

Only the top bits are embedded, so keep shard boundaries aligned to that prefix granularity; otherwise fall back to a lookup.

Worked example: splitting a shard while serving traffic

Shard C holds range 80-bf and is at 75 percent disk and 70 percent CPU at peak. We split it into C1 (80-9f) and C2 (a0-bf). Every step below is reversible until step 6.

  1. Provision C1 and C2 as empty primaries with replicas, schema applied, not yet in the shard map.
  2. Copy. Take a consistent snapshot of C, recording its change-stream position (binlog GTID in MySQL, LSN in PostgreSQL). Stream rows into C1 or C2 by keyspace id. Throttle the copy against C's replica lag so production does not suffer.
  3. Catch up. From the recorded position, apply C's ongoing changes to C1 and C2, filtered by range, until the lag is seconds. Vitess implements steps 2 and 3 as a VReplication-based Reshard workflow; elsewhere it is a change-data-capture pipeline you build or buy.
  4. Verify. Compare row counts and per-range checksums between C and the union of C1 and C2, ideally against a quiesced snapshot. Mismatches block the cutover.
  5. Switch reads. Publish a map version that routes reads for 80-bf to C1 and C2 replicas while writes still go to C. Watch error rates and latency. Rolling back is just republishing the old map.
  6. Freeze and switch writes. C stops accepting writes for 80-bf; wait until C1 and C2 have applied every change up to C's final position; publish the new map; C1 and C2 start accepting writes. The write pause is typically seconds. This is the point of no return unless you also run reverse replication from C1 and C2 back to C, which Vitess sets up by default and you should copy if building your own.
  7. Retire. After a soak period, stop reverse replication, drop C from the map, and archive it.
def switch_writes(src, dsts, ranges, map_store, timeout_s=10):
    src.deny_writes(ranges)                        # 1. freeze
    final_pos = src.current_position()
    for d in dsts:
        if not d.wait_applied(final_pos, timeout_s):   # 2. drain
            src.allow_writes(ranges)               # abort: nothing changed yet
            raise CutoverAborted(f"{d} lagging")
    new_map = map_store.current().reassign(ranges, dsts)
    map_store.compare_and_set(new_map)             # 3. publish, bumping version
    for d in dsts:
        d.allow_writes(ranges)                     # 4. open
    start_reverse_replication(dsts, src)           # 5. keep rollback possible

The ordering is the whole design: freeze before draining, publish only after draining, open only after publishing. Any crash before compare_and_set leaves the old map in force and the old shard authoritative.

Failure modes

FailureSymptomPrevention
Stale router writes to old ownerRows missing after a split; duplicates on two shardsShard-side range enforcement and a write freeze during moves
Copy overloads the sourceProduction latency spikes during reshardingThrottle the copy on replica lag and CPU; copy from a replica
Catch-up never convergesLag grows while write rate exceeds apply rateSplit earlier, at 60 to 70 percent of capacity, not at 95
Lookup table driftsLogin cannot find an existing customerWrite the lookup in the same transaction where possible; run reconcilers
Proxy tier saturatesAll shards healthy, all queries slowScale proxies separately and alert on their CPU and connection counts
Schema change half appliedSome shards reject new columnsRun migrations per shard with a tracker; make application code tolerate both schemas
Hot key on one shardOne shard pegged regardless of splitsSplitting cannot fix a single key; see hot-key mitigation

Operating a sharded fleet

Track size, QPS, p99 latency and replication lag per shard, and alert on the spread between the busiest and the median shard, because imbalance is the early warning for skew and for an upcoming split.

Backups are per shard, so a cluster-wide restore point needs either a coordinated snapshot or point-in-time recovery to a common timestamp on every shard. Practise that restore. Cross-shard business operations, such as moving credit between two customers on different shards, should be designed as a saga with compensations or as events published through the outbox pattern, rather than relying on two-phase commit across shards. Read-heavy traffic can be pushed to shard replicas with the same rules described in read replica routing.

Trade-offs

A proxy tier buys language independence at the cost of a hop and a fleet to run. Range-over-hash buys cheap splits at the cost of a strongly consistent metadata service. Lookup tables make non-key queries fast but add a second write. Reverse replication makes cutovers reversible but doubles replication load while it runs. How much to automate depends on how often you reshard: yearly splits tolerate manual steps; weekly ones do not.

What to do next

  1. Write down your current key-to-shard mapping and check whether you can split a single shard without moving data on the others; if not, introduce a keyspace-id indirection.
  2. Make every shard enforce the key ranges it owns and reject writes outside them.
  3. Store the shard map in a strongly consistent store with a version number, and make routers refresh on a wrong-shard error.
  4. List every hot query that does not carry the shard key and choose scatter, embedded id or lookup table for each.
  5. Set a split threshold, for example 65 percent of disk or peak CPU, and alert on it per shard.
  6. Rehearse a full split on staging, including the write freeze and a rollback using reverse replication, and time each step.
Key takeaway: A sharded database is shards plus a shard map plus a router. Hash keys to keyspace ids and assign ranges of them to shards so one shard can split without disturbing the rest; keep the map versioned in a consistent store; let routers be stale but make shards enforce ownership; serve non-key queries with embedded ids or lookup tables; and move data with copy, catch-up, verify, freeze, publish and reverse replication so every cutover can be undone.