PrepZone Logo
PrepZone

Partitioning and Sharding

Vertical vs horizontal splits, hash vs range keys, and rebalancing without downtime.

Read these first

Why this matters

  • VaultCommerce crossed 8K sustained order writes/sec during a flash sale — vertical scale and read replicas were no longer enough.
  • Choosing the wrong shard key creates hot shards that negate distribution benefits.
  • Interviews expect you to defend a partition key and explain what breaks when you add shards.

Vertical vs horizontal partitioning

Vertical
usersid, name, email
user_profilesbio, avatar
Horizontal
Shard Auser_id 0–999
Shard Buser_id 1000–1999
Vertical splits tables by column group. Horizontal splits rows across nodes by key.

Two partitioning styles

  • Vertical — split wide tables (users vs user_profiles) or move cold columns to archive storage.
  • Horizontal — split rows across nodes by hash or range of a key.
  • When vertical helps — reduce row width, isolate LOB columns, separate read profiles.
  • When horizontal helps — write QPS or storage exceeds one node's capacity.

VaultCommerce vertically split product descriptions into MongoDB while keeping SKU inventory in Postgres. Horizontal sharding arrived later when the orders table exceeded comfortable single-node write throughput.

Shard key selection

The shard key must match your hottest queries and spread load evenly.

Java
# VaultCommerce order shard router
def shard_for_user(user_id: str) -> str:
    bucket = hash(user_id) % NUM_SHARDS
    return f"orders-shard-{bucket}"

Sharding by user_id co-locates a customer's order history on one shard — "my orders" never crosses shards. Global admin reports aggregate asynchronously from all shards.

StrategyMechanismRisk
Hashhash(key) % NEven spread; range queries fan out
Rangeuser_id 1–1M on shard AHot ranges if IDs cluster
DirectoryLookup service maps key → shardLookup service is critical path
GeographicEU users on EU shardCompliance win; cross-region queries hard
  • Hash

    Mechanismhash(key) % N
    RiskEven spread; range queries fan out
  • Range

    Mechanismuser_id 1–1M on shard A
    RiskHot ranges if IDs cluster
  • Directory

    MechanismLookup service maps key → shard
    RiskLookup service is critical path
  • Geographic

    MechanismEU users on EU shard
    RiskCompliance win; cross-region queries hard

Cross-shard challenges

What sharding breaks

  • Joins — no cross-shard SQL joins; denormalise or fetch-and-merge in application code.
  • Global uniqueness — use UUIDs or a central id generator, not AUTO_INCREMENT per shard.
  • Transactions — ACID is per shard; cross-shard checkout uses sagas.
  • Rebalancing — adding shards moves data; consistent hashing reduces remapping pain.
Java
-- Anti-pattern: cross-shard join at query time
-- SELECT o.*, u.email FROM orders_shard_2 o JOIN users_shard_0 u ...
-- Instead: store denormalised email on order row or batch-fetch users by ID

VaultCommerce denormalises customer_email on the order row at write time — a deliberate trade for shard-local reads.

Migration ladder

Do not shard on day one:

Java
vaultcommerce_scaling_ladder:
  1: single_postgres_primary
  2: read_replicas_for_browse
  3: redis_cache_for_hot_products
  4: horizontal_shard_orders_by_user_id
  5: per_shard_replicas_for_read_scale

Each step has operational cost. VaultCommerce reached step 4 at ~3M daily active users when write metrics, not storage alone, forced the decision.

Quick recall

Everything you need if you only revisit this box.

  1. Vertical splits columns/tables; horizontal splits rows across nodes.
  2. Pick a shard key that matches hot queries and spreads load evenly.
  3. VaultCommerce shards orders by user_id for local "my orders" reads.
  4. Cross-shard joins and global transactions do not work — design around them.
  5. Shard only after replicas and caching are exhausted; each shard gets its own replicas.

Test yourself

Answer these before moving on — recall is what makes it stick.