Command Palette

Search for a command to run...

Hectal
PHASE 9Advanced ~9 min· topic 3 of 5

Topic 9.3

Sharding: Shard Keys and Strategies

In one line

Sharding spreads data across multiple independent database servers so writes and storage scale horizontally. The shard key decides everything: it should spread load evenly, keep each request's data on one shard, and never need to change. Strategies are hash, range, directory (lookup) and geographic, usually with many logical shards mapped onto fewer physical nodes.

0/5 · 0%

Think of it like this

A national bank splitting customers across regional offices. Assign by customer ID hash and every office gets a fair share; assign by city and Mumbai's office is overwhelmed; either way, a customer's accounts and transactions should all be at the same office.

Key ideas

  1. 01

    Shard only when necessary: a well-tuned PostgreSQL primary with replicas handles tens of thousands of writes/sec and many TB. Sharding adds cross-shard queries, distributed transactions, rebalancing, and operational load.

  2. 02

    Good shard keys: high cardinality, even distribution of data and traffic, present in almost every query, and grouping data that's transactionally related (tenant_id for B2B SaaS, user_id for consumer apps). Co-locate child tables by the same key (orders and order_items by customer_id).

  3. 03

    Hash sharding: shard = hash(key) mod N spreads evenly but range queries hit all shards. Range sharding keeps ranges together but concentrates sequential keys (time, sequential IDs) on one hot shard. Directory sharding stores key → shard in a lookup service: flexible (move a tenant individually) but the directory must be fast and highly available. Geographic sharding places users' data in their region (latency, data residency).

  4. 04

    Logical shards: create many more logical shards (e.g. 4,096) than physical nodes and map them to nodes. Rebalancing moves whole logical shards instead of rehashing every key. Consistent hashing is another way to minimise movement.

  5. 05

    Tools: Citus (PostgreSQL extension), Vitess (MySQL), application-level sharding with a routing library, or natively distributed SQL (CockroachDB, YugabyteDB, Spanner) that shards by key range automatically.

Code & diagrams

shard-router.javajava
public final class ShardRouter {
    private static final int LOGICAL_SHARDS = 4096;
    private final int[] logicalToPhysical;          // loaded from config / directory
    private final List<DataSource> physical;

    public DataSource forTenant(long tenantId) {
        int logical = Math.floorMod(Long.hashCode(mix(tenantId)), LOGICAL_SHARDS);
        return physical.get(logicalToPhysical[logical]);
    }

    private static long mix(long x) {               // spread sequential ids (splitmix64 finaliser)
        x = (x ^ (x >>> 30)) * 0xbf58476d1ce4e5b9L;
        x = (x ^ (x >>> 27)) * 0x94d049bb133111ebL;
        return x ^ (x >>> 31);
    }
}
citus.sqlsql
-- Citus: distribute tables by tenant and co-locate them
SELECT create_distributed_table('orders', 'tenant_id');
SELECT create_distributed_table('order_items', 'tenant_id', colocate_with => 'orders');
SELECT create_reference_table('currency');     -- small table copied to every node

-- single-shard query (routed): includes tenant_id
SELECT * FROM orders o JOIN order_items i USING (tenant_id, order_id)
WHERE o.tenant_id = 42 AND o.order_id = 1001;

Interview problem

The problem

Choose a shard key for a multi-tenant SaaS

A project-management SaaS has 200K tenants; the largest has 5% of all data, most are tiny. Queries are almost always within a tenant; there is an internal cross-tenant analytics dashboard. The single database is at 30 TB and write capacity is exhausted. Choose a key and strategy.

When it breaks

Sharding by a monotonically increasing key with range strategy

What you see

All new writes go to the last shard; one node is at 100% while others idle.

Fix & prevent

Hash the key or use a key with natural spread (tenant/user); for time-series, combine time with a hashed bucket.

Explain it without notes

01

Why create many logical shards per physical node?

Practice

01

Compute how many keys move when growing from 4 to 5 nodes with hash mod N vs with 4,096 logical shards.

Trade-offs

  • ↔

    Sharding scales writes and storage horizontally, at the price of cross-shard complexity, rebalancing and harder operations.

Done when you can

  • I can choose a shard key and strategy and justify it against access patterns and hotspots.