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.
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
Loading…