Learning on Web Dev Open is free for all.

System Design & Performance > Scale, consistency and failureWrites, where it gets awkward
Phase 06Scale, consistency and failure312 of 434

Writes, where it gets awkward

One primary until it is not enough, then partitioning, and the hot key that undoes your beautifully even shard scheme.

Concept15 minAI adversary

Replicas do not help writes, because every write must reach every copy. The single-writer ceiling is genuinely high on modern hardware, often much higher than the traffic people are designing for, so the first honest answer to a write-scaling question is usually to check whether you have one. Batching, removing write amplification and trimming unnecessary indexes buy more than people expect.

Beyond that you partition: split data across independent shards by a key so each takes a fraction of the writes. The key decides everything. Partition by tenant and every query naturally scopes to one shard, but your largest tenant now sits on one machine. Partition by hash of id and load spreads evenly, but any query that spans ids fans out to every shard. Choose from the queries you must serve, not from the elegance of the distribution.

Hot keys are the recurring problem: one product, one celebrity, one enormous customer generating a disproportionate share. Mitigations exist: salt the key into sub-partitions, cache the hot item aggressively, or give the outlier dedicated capacity, but they are deliberate work. Assuming even distribution is the most common way a sharding plan fails in its first month.

You should now be able to

  • Explain why writes do not scale by adding replicas
  • Choose a partition key from access patterns
  • Recognise and mitigate a hot partition
Ask the community

Loading…