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
Two partitioning styles
- Vertical — split wide tables (
usersvsuser_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.
# 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.
| Strategy | Mechanism | Risk |
|---|---|---|
| Hash | hash(key) % N | Even spread; range queries fan out |
| Range | user_id 1–1M on shard A | Hot ranges if IDs cluster |
| Directory | Lookup service maps key → shard | Lookup service is critical path |
| Geographic | EU users on EU shard | Compliance win; cross-region queries hard |
Hash
Mechanismhash(key) % NRiskEven spread; range queries fan outRange
Mechanismuser_id 1–1M on shard ARiskHot ranges if IDs clusterDirectory
MechanismLookup service maps key → shardRiskLookup service is critical pathGeographic
MechanismEU users on EU shardRiskCompliance 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_INCREMENTper shard. - Transactions — ACID is per shard; cross-shard checkout uses sagas.
- Rebalancing — adding shards moves data; consistent hashing reduces remapping pain.
-- 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:
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.
- Vertical splits columns/tables; horizontal splits rows across nodes.
- Pick a shard key that matches hot queries and spreads load evenly.
- VaultCommerce shards orders by
user_idfor local "my orders" reads. - Cross-shard joins and global transactions do not work — design around them.
- 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.