System Design
X2

Partitioning and sharding

Split one dataset across many machines, so that each machine holds and serves only part of it. The shard key decides which machine owns a row.

Not startedSaved in this browser only.
  1. 1Shard last. Replicas, a bigger machine and partitioned tables come first.
  2. 2The shard key decides everything: the common query must hit one shard, and writes must spread evenly.
  3. 3Hash mod N moves almost every key when N changes. A ring or fixed slots move only the new node's share.
  4. 4Hot keys, cross-shard queries and cross-shard transactions are the cost. Design them out with the key.
X2
    A

    Add a node, move most keys

    hash mod N
    hash value of the key, then its nodehashmod 3mod 40AA1BB2CC3AD4BA5CB6AC7BD8CA9AB10BC11CD✕ 9 of 12 keys change nodeA key stays only if hash mod 3 = hash mod 4: keys 0, 1, 2.With N to N+1 nodes, N/(N+1) of the keys move.
    • Shard when one primary cannot take the writes, or cannot hold the data.
    • Each shard owns a subset of keys. A router maps a key to its shard.
    • The placement rule decides how much data moves when the cluster grows.

    I never place keys with hash mod N on a cluster that grows. Adding one node would move most of the data.

    B

    Strategies

    which for which access pattern
    strategyrange scanadd a nodefits
    Range of the keyApprovedSplit a rangeOrdered keys, time within a source
    Hash mod NNot approvedNot approvedA shard count that never changes
    Hash into fixed slots, slot mapNot approvedMove slotsMost key-value and OLTP shards
    Consistent hash ring, virtual nodesNot approvedMoves 1/(N+1)Peer clusters with no central map
    Directory: key to shard tablePer entryMove one keyBig tenants, uneven key sizes
    Geographic: region in the keyIn a regionPer regionData residency, local latency

    Fixed slots and the ring both move only the new node's share. Slots keep an explicit map; a ring derives the map from token positions.

    I use ranges when queries scan ordered keys, and hashing into slots or a ring when writes must spread evenly.

    C

    Pick a shard key

    four tests, then examples
    writes spread evenly over shards →hoteventhe common query needs one shard →allone✓tenant_idSaaS✓user_idfeeds, profiles✓hotel_idbooking✓(source, day)logs!order_idorders by customer✕created_atlogs✕countryfew valuesshard key candidatesShaded: pick from here. Positions are typical, not measured.
    cardinality
    Many more values than shards. A country column has about 200 values; most traffic sits in a few.
    even spread
    No single value takes a large share of rows or writes.
    alignment
    The common query names the key, so it touches one shard.
    growth
    New data does not all go to one value. A timestamp fails this test.
    workloadshard keythe common querystatus
    SaaS apptenant_idEverything for one tenantBig tenants move out
    Hotel bookinghotel_idHold and confirm nights at one hotelApproved
    Feed, profileuser_idOne user's posts and settingsApproved
    Chatconversation_idMessages of one conversation, in orderApproved
    Logs, metricscreated_atRecent linesOne hot shard
    Logs, metrics(source_id, day)One source over a time rangeApproved
    Ordersorder_idA customer's ordersScatters by customer
    Orderscustomer_idA customer's ordersApproved

    I pick the key that the most common query names. Then I check that no value takes a large share of writes, and that new data spreads too.

    D

    The request path

    click a step; its path lights up
    Clientapp or browserService + routerstateless, N copiesShard mapkey range → shardShard 1Postgres primaryShard 2Postgres primaryShard 3Postgres primary

    Step 1: Find the key

    • Take the shard key from the request: tenant_id = 42.
    • Hash it, or keep it as is for range sharding.

    If it fails

    The request has no shard key. Then it must go to every shard (step 4).

    The router finds the shard from a cached map. A query with the key touches one shard; a query without it touches all of them.

    E

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives this designalso used for
    PostgresDeclarative partitioning: PARTITION BY RANGE, LIST or HASHOne logical table, many physical tables on one server. The step before sharding.Time series, multi-tenant tables
    PostgresPartition pruningA query with the partition key reads one partition. Without it, the plan reads all of them.Date-bounded reports
    PostgresDETACH PARTITION, DROP TABLERemove a whole month at once, with no row-by-row DELETE.Retention for logs and events
    PostgresUnique indexes on a partitioned tableLimit They must contain the partition key. Global uniqueness needs another plan.
    PostgresLogical replication with row filters (version 15 and later)Copy only one tenant's rows to a new server while it stays live.Major-version upgrades, change data capture
    Postgrespostgres_fdw: a partition can be a foreign tableA partition can live on another server and still answer through the parent table.Federated reporting
    CitusDistributed tables by a distribution column, co-located tables, reference tablesSharded Postgres: joins on the same key run on one worker.Multi-tenant SaaS, real-time analytics
    RedisCluster: 16,384 hash slots, CRC16 of the keyFixed slots and a slot map. Resharding moves slots while the cluster serves.Any Redis data set larger than one node
    RedisHash tags {user:42} and MOVED, ASK repliesKeys with one tag share a slot. A client with a stale map is told where the slot went.Multi-key scripts in a cluster
    CassandraToken ring with virtual nodes (num_tokens, 16 by default)A consistent hash ring. Partition key hashed with Murmur3; replicas follow the ring.Write-heavy wide tables
    DynamoDBPartition key hashed to partitions; partitions split as they growManaged hash sharding. A hot partition key still hits one partition.Serverless key-value
    ServiceA router library with a cached shard mapFinds the shard in memory. Reloads the map when a shard refuses a key.Any shard-aware service

    On one Postgres I partition tables and let pruning skip partitions. Across machines I use slots or a ring, and I move data with logical replication or double writes.

    F

    Consistent hashing

    pseudo code
    ring with virtual nodespseudo code
    build(nodes, vnodes):                      // once per membership change
      FOR EACH node IN nodes:
        FOR i IN 0 .. vnodes-1:
          put token hash(node + "#" + i) on the ring1
      sort tokens by position
    
    owner(key):
      h = hash(key)
      t = first token at or after h2             // binary search
      IF there is none: t = the first token     // wrap around3
      RETURN t.node
    
    add(node):
      put the vnodes tokens of node on the ring
      keys just before each new token move to node4
    1. 1A virtual node is one token. Its position depends only on the node name and the index, so other nodes never move.
    2. 2Binary search over the sorted tokens: O(log of the token count).
    3. 3The ring is a circle. A hash past the last token belongs to the first one.
    4. 4Only these keys move, and all of them move to the new node.
    Tested source Go: the ring
    Go: the ringgo
    // Ring is a consistent hash ring. Each node owns vnodes points (tokens) on the ring; a key
    // belongs to the first token at or after its hash, wrapping at the end.
    type Ring struct {
      tokens []Token // sorted by Pos
    }
    
    // Token is one virtual node: a point on the ring and the node that owns it.
    type Token struct {
      Pos  uint64
      Node string
    }
    
    // NewRing places vnodes tokens for each node. A token's position depends only on its node and
    // index, so adding a node never moves another node's tokens.
    func NewRing(nodes []string, vnodes int) *Ring {
      r := &Ring{}
      for _, n := range nodes {
        r.tokens = append(r.tokens, NodeTokens(n, vnodes)...)
      }
      sort.Slice(r.tokens, func(i, j int) bool { return r.tokens[i].Pos < r.tokens[j].Pos })
      return r
    }
    
    // NodeTokens returns the tokens of one node.
    func NodeTokens(node string, vnodes int) []Token {
      out := make([]Token, vnodes)
      for v := range out {
        out[v] = Token{Pos: Hash(fmt.Sprintf("%s#%d", node, v)), Node: node}
      }
      return out
    }
    
    // Owner implements Placer: binary search for the first token at or after the key's hash.
    func (r *Ring) Owner(key string) string {
      h := Hash(key)
      i := sort.Search(len(r.tokens), func(i int) bool { return r.tokens[i].Pos >= h })
      if i == len(r.tokens) {
        i = 0 // past the last token: wrap to the first
      }
      return r.tokens[i].Node
    }
    

    A key belongs to the first token clockwise from its hash. A new node only takes the arcs just before its own tokens.

    G

    What a new node moves

    recorded from the lab simulation
    keys moved, 4 to 5 nodes0%20%40%60%80%100%ideal 1/5hash mod N: 79.6%hash mod N79.6%1,024 fixed slots: 19.9%1,024 fixed slots19.9%ring, 1 token: 14.8%, from 0.2% to 45.7%ring, 1 token14.8%ring, 16 tokens: 19.6%, from 13.9% to 26.9%ring, 16 tokens19.6%ring, 256 tokens: 20.0%, from 16.3% to 23.1%ring, 256 tokens20.0%share of 20,000 keys; whisker: 20 ring layouts
    fullest of 8 nodes, against its fair share1×1.5×2×2.5×3×3.5×hash mod N: 1.02×hash mod N1.02×1,024 slots: 1.01×1,024 slots1.01×ring, 1 token: 2.42×, from 1.66× to 3.32×ring, 1 token2.42×ring, 4 tokens: 1.76×, from 1.35× to 2.42×ring, 4 tokens1.76×ring, 16 tokens: 1.35×, from 1.18× to 1.60×ring, 16 tokens1.35×ring, 64 tokens: 1.18×, from 1.11× to 1.29×ring, 64 tokens1.18×ring, 256 tokens: 1.09×, from 1.02× to 1.17×ring, 256 tokens1.09×ring, 1024 tokens: 1.05×, from 1.03× to 1.07×ring, 1024 tokens1.05×100,000 keys; whisker: 20 ring layouts

    Going from 4 to 5 nodes, hash mod N moves about 80% of the keys. Slots and a ring with enough tokens move about 20%, the new node's fair share.

    H

    Try it: the hash ring

    recorded from the lab simulation
    DBAC4 nodes1 token each
    • A's ranges
    • key on A
    • key moved by the last change
    nodes · click to add or remove
    virtual nodes (tokens) per node
    Add node E to see which keys move.Then raise the virtual nodes and watch the load bars even out.
    share of 20,000 keys · click a node to highlight itA34.3%B24.7%C30.9%D10.1%

    Fair share 25.0% (dashed). Fullest node: 1.37× its fair share.

    Six nodes, any subset, five token counts: 315 rings, each simulated over 20,000 keys. Dots are 48 of those keys.

    With one token per node the arcs are uneven and a new node takes keys from one neighbour. With a few hundred tokens each node owns many small arcs, so load and moves both even out.

    I

    Hot partitions

    recorded from Postgres
    10,000 new log lines, by partitionwrites0Oct 10Oct 210,000Oct 3, today0Oct 4✕ One partition takes 100% of writes. The others are idle.
    fixstatuscost
    Hash a high-cardinality column first, time secondApprovedA time-range query over all sources scatters.
    Split the hot rangeRange shardingThe newest range is hot again after the split.
    Cache the hot readsRead-hot keysDoes nothing for writes.
    Salt the hot keyWrite-hot keysEvery read of that key reads k shards.
    Own shard for the hot tenantWith a directoryOne more entry to keep in the map.

    A timestamp as the key sends every new write to the newest range. I hash the source first and keep time inside the shard.

    J

    Split a hot key

    recorded from the lab simulation
    hottest of 8 shards; one key takes 25.0% of writes0%10%20%30%40%fair share 1/8no salt: 34.4%no salt34.4%2 sub-keys: 23.1%, from 21.9% to 34.4%2 sub-keys23.1%4 sub-keys: 19.8%, from 15.6% to 28.2%4 sub-keys19.8%8 sub-keys: 17.3%, from 15.6% to 25.0%8 sub-keys17.3%16 sub-keys: 16.0%, from 14.0% to 20.3%16 sub-keys16.0%32 sub-keys: 14.8%, from 13.3% to 17.2%32 sub-keys14.8%share of 1,000,000 writes; whisker: 100 key names
    salt a hot keypseudo code
    write(key, value):
      IF key is hot:
        key = key + "#" + (count MOD k1)         // k sub-keys
      shard(key).write(value)
    
    read(key):
      IF key is hot:
        RETURN merge(2shard(key + "#" + i).read FOR i IN 0 .. k-1)
      RETURN shard(key).read
    1. 1Round robin over k sub-keys. Each sub-key hashes to its own shard, so the writes spread.
    2. 2The cost: one read becomes k reads. Salt only keys that are hot for writes.
    Tested source Go: salt
    Go: saltgo
    // Salt spreads the writes of one hot key over k sub-keys ("user:0#0" to "user:0#k-1"), round
    // robin. Each sub-key hashes to its own shard. A read of the hot key must read all k sub-keys.
    func Salt(writes []string, hot string, k int) []string {
      out := make([]string, len(writes))
      next := 0
      for i, w := range writes {
        if w != hot {
          out[i] = w
          continue
        }
        out[i] = fmt.Sprintf("%s#%d", hot, next%k)
        next++
      }
      return out
    }
    

    Sub-keys hash like any key, so two can land on one shard. The whisker shows that spread.

    One key with a quarter of the writes made its shard take 34.4% of all writes. Split into 8 sub-keys, the hottest shard took about 17.3%.

    K

    Cross-shard queries

    schema, then recorded query plans
    one table, 8 hash partitionssql
    -- One logical table, split into 8 partitions by a hash of tenant_id.
    CREATE TABLE events (
      tenant_id  bigint      NOT NULL,
      event_id   bigint      NOT NULL,
      created_at timestamptz NOT NULL,
      body       text        NOT NULL,
      PRIMARY KEY (tenant_id, event_id)2
    ) PARTITION BY HASH (tenant_id)1;
    
    CREATE TABLE events_p0 PARTITION OF events
      FOR VALUES WITH (MODULUS 8, REMAINDER 03);
    CREATE TABLE events_p1 PARTITION OF events
      FOR VALUES WITH (MODULUS 8, REMAINDER 1);
    -- ... events_p2 to events_p7
    1. 1Postgres hashes tenant_id to pick the partition. A sharded cluster does the same across servers.
    2. 2The key must contain tenant_id. Each partition checks uniqueness only for its own rows.
    3. 3Hash mod 8. To grow, detach one partition, replace it with two of modulus 16, and copy its rows back.
    WHERE tenant_id = 42planner or routerp0p1p2p3p4p5p6p7✓ reads 1 of 8 partitionsWHERE event_id = 7planner or routerp0p1p2p3p4p5p6p7✕ reads 8 of 8, then merges
    methodstatuscost
    Scatter-gather to every shardRare queriesLatency of the slowest shard; load on all.
    Global secondary index: value to shard keyApprovedA second write, usually asynchronous.
    Denormalised copy sharded by the other keyHot lookupsTwo copies to keep in step.
    Shard number inside the IDApprovedThe ID format is fixed for good.

    A query without the shard key goes to every shard. For a common lookup I keep a global index or put the shard number in the ID.

    L

    Cross-shard transactions

    avoid them; then use a saga
    move 100 from A to B, time →shard 1account Ashard 2account B1 · debit Acommit2 · credit B ✕account closed3 · refund AcompensationEach step is one local transaction on one shard.A failed step runs the compensations of the steps before it.
    methodstatuswhy
    Co-locate: both rows share the keyApprovedA normal local transaction.
    Saga with compensations, from an outboxApprovedEach step commits locally. Others can see the middle state.
    Two-phase commitFew shards, low rateA prepared transaction holds its locks until the coordinator decides.
    Two independent commits, no planNot approvedA crash between them leaves money debited and not credited.

    I choose the key so that one transaction touches one shard. When it cannot, I run a saga with compensations, not two-phase commit.

    M

    Unique across shards

    IDs and unique values
    methodstatusnote
    Unique index without the shard keyNot approvedPostgres refuses it on a partitioned table.
    UUID, random or time-orderedApprovedNo coordination. 16 bytes; time order keeps inserts local in the index.
    64-bit ID: time, worker, sequenceApprovedEach process needs a unique worker number.
    Shard number inside the IDApprovedA lookup by ID finds its shard with no index.
    Claim table sharded by the value: email to user_idUnique emailsInsert the claim first. A duplicate fails there.
    One central sequenceLow rateEvery insert waits on one server.

    The lab error, verbatim: "unique constraint on partitioned table must include all partitioning columns".

    A unique index works only inside one shard. I generate IDs that are unique by design, and I claim unique values like emails in a table sharded by that value.

    N

    Resharding without downtime

    pseudo code, then the race
    move a rangepseudo code
    move(range, old, new):
      1  turn on double writes:
           write old shard                      // version = version + 1
           write new shard IF its version is older1
      2  backfill: FOR EACH batch of rows IN old:
           copy to new, SKIP rows new already has2
      3  verify: rows that differ == 03
      4  switch reads of the range to new4
      5  stop writes to old, then delete the range from old
    1. 1Two writers can reach the new shard in either order. The version guard keeps the newer row.
    2. 2ON CONFLICT DO NOTHING. A double write made after the backfill read is newer than the copy.
    3. 3A full outer join of the two shards. Switch nothing until it is zero.
    4. 4One entry in the shard map changes. Routers that still use the old map are refused and reload.
    Tested source Go: double write · Go: backfill · SQL: writes and copy
    Go: double writego
    // Write changes one account during a move: old shard first, then the new shard. The version
    // guard on the new shard makes the two writes safe in any order across writers.
    func Write(ctx context.Context, db *pgxpool.Pool, id, delta int64) error {
      var balance, version int64
      if err := db.QueryRow(ctx, reshard["write_old"], id, delta).Scan(&balance, &version); err != nil {
        return fmt.Errorf("write old shard, account %d: %w", id, err)
      }
      if _, err := db.Exec(ctx, reshard["write_new"], id, balance, version); err != nil {
        return fmt.Errorf("write new shard, account %d: %w", id, err)
      }
      return nil
    }
    
    Go: backfillgo
    // Backfill copies the old shard to the new one in batches of size, in id order, while the
    // service keeps double writing.
    func Backfill(ctx context.Context, db *pgxpool.Pool, size int) error {
      var after int64
      for {
        batch, err := ReadBatch(ctx, db, after, size)
        if err != nil {
          return err
        }
        if len(batch) == 0 {
          return nil
        }
        if err := CopyBatch(ctx, db, batch, false); err != nil {
          return err
        }
        after = batch[len(batch)-1].ID
      }
    }
    
    SQL: writes and copysql
    -- The service writes the old shard first. The version counts every change to the row.
    UPDATE accounts_old
    SET balance = balance + $2, version = version + 1
    WHERE id = $1
    RETURNING balance, version;
    
    -- Then it writes the same row to the new shard. An older version never replaces a newer one.
    INSERT INTO accounts_new (id, balance, version) VALUES ($1, $2, $3)
    ON CONFLICT (id) DO UPDATE
    SET balance = excluded.balance, version = excluded.version
    WHERE accounts_new.version < excluded.version;
    
    -- The backfill skips a row the new shard already has. A double write put it there, and it is newer.
    INSERT INTO accounts_new (id, balance, version)
    SELECT * FROM unnest($1::bigint[], $2::bigint[], $3::bigint[])
    ON CONFLICT (id) DO NOTHING;
    
    -- Rows that differ between the shards, or exist on one side only.
    SELECT count(*)
    FROM accounts_old o
    FULL JOIN accounts_new n USING (id)
    WHERE (o.balance, o.version) IS DISTINCT FROM (n.balance, n.version);
    backfill jobold shardnew shardserviceread batch1id 1: 100, v1+50 → 150, v22150, v23copy 100, v14✕ overwrite existing rowsold 150 · new 1001 row differs: the +50 is lost✓ skip existing rowsold 150 · new 1500 rows differ
    • Tested with 8 writers and 3,200 double writes during a backfill of 2,000 rows: 0 rows differ.
    • A delete during the move needs a tombstone. Otherwise the backfill copies the deleted row back.
    • With a ring or slots, only 1/(N+1) of the data takes this path.

    I double write with a version guard, backfill rows the new shard does not have, verify, and only then switch the map. The backfill never overwrites a row.

    O

    Failure cases

    what breaks, and why it stays correct
    eventresultwhy it is safesaved by
    A shard primary failsIts keys fail; other shards serve.Its replica is promoted. Only 1/N of users notice.Replica
    A router has a stale map after a moveIt sends the key to the old shard.The old shard refuses keys it no longer owns. The router reloads and retries.Map version
    The backfill copies a row older than a double writeThe new shard already has the row.The copy skips rows that exist.ON CONFLICT
    Two writers reach the new shard out of orderThe older version arrives second.The upsert changes the row only when its version is older.Version guard
    A row is deleted during the moveThe backfill can copy it back.Delete with a tombstone until the move ends. Verify finds any leftover.Tombstone
    One tenant grows to 30% of the loadIts shard runs hot.A directory entry moves that tenant to its own shard.Directory
    One shard is slow during a scatter queryThe whole query waits.A timeout returns partial results, marked as partial.Timeout
    A transfer fails between two shardsDebited, not credited.The saga runs the refund from its outbox record.Saga
    P

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1Postgres2+ replicas3+ partitions4+ shards5+ reshardmore load →
    Write capacity against demand1k10k100k1MDemand, average: 11,574 single-row inserts per secondDemand, average11,574Demand, 5× peak: 57,870 single-row inserts per secondDemand, 5× peak57,870One Postgres: 50,000 single-row inserts per secondOne Postgres50,0004 shards: 200,000 single-row inserts per second4 shards200,0008 shards: 400,000 single-row inserts per second8 shards400,000single-row inserts per second, log scale
    stepaddit handlesmove up when you see
    1One Postgres with indexes. Scale it up first: more cores, memory, faster disks.About 50,000 single-row inserts a second in the lab, and every query on one machine.Reads slow the primary.
    2Read replicas, as on the replication sheet.Reads grow with each replica. Writes do not.Tables grow so large that indexes and retention slow down.
    3Partitioned tables on the same server.Smaller indexes, fast retention by DROP, pruning by key.One primary cannot take the write rate, even on the largest machine.
    4Shards by key, with a router and a shard map.Writes and storage grow with the shard count. The common query stays on one shard.One shard runs hot, or the cluster needs more nodes.
    5Reshard: move ranges or slots to new shards online.Growth with no downtime. Only 1/(N+1) of the data moves per new node.Top of the ladder.

    Demand example: 1 billion events a day is 1,000,000,000 / 86,400 ≈ 11,574 a second. Shard capacity is the lab measurement times the shard count; real shards scale a little below that. Postgres 16 on an 8-core laptop, 32 clients.

    I add replicas for reads, a bigger machine and partitions for size. I shard only when one primary cannot take the writes, and I pick the key from the main query.

    Q

    Drill

    predict, then reveal

    0 of 10 known

    1. You go from 4 to 5 shards with hash mod N. About how many keys move?

    2. You add node E to a consistent hash ring. Which keys move, and to where?

    3. Why give each node many virtual nodes?

    4. You shard logs by timestamp. What goes wrong?

    5. A SaaS app shards by tenant_id. One tenant sends 30% of the traffic. What do you do?

    6. Orders are sharded by customer_id. How do you find an order by order_id?

    7. Why does Postgres refuse UNIQUE (event_id) on a table partitioned by tenant_id?

    8. During a backfill, why skip rows that already exist on the new shard?

    9. A transfer moves money between users on two shards. How do you make it safe?

    10. Why does Redis Cluster use 16,384 fixed slots instead of a ring?

    R

    Numbers to say

    measured, derived or cited
    mod N
    N to N+1 nodes moves N/(N+1) of keys: 4 to 5 moved 79.6%.
    ring, slots
    Move about 1/(N+1): 19.9% with slots, 20.0% with a 256-token ring.
    tokens
    Fullest of 8 nodes: 2.42× fair share with 1 token, 1.09× with 256.
    hot key
    A key with 25.0% of writes: its shard takes 34.4% of 8 shards' writes.
    Redis
    16,384 hash slots (cluster specification).
    Cassandra
    16 tokens per node by default (configuration reference).
    one node
    About 50,000 single-row inserts a second, measured.

    Simulations are seeded and repeat exactly. Postgres 16 on an 8-core laptop, 32 clients; a server that flushes each commit to durable storage commits slower.