Scaling reads is the easy half
Replicas, caches and precomputation, and the replication lag that makes a user think their own edit did not save.
Reads scale by duplication: more replicas, more cache layers, or precomputed results. All three are variations on keeping copies, and every copy is a decision about how out of date it may be. This is why scaling reads is mostly a product conversation wearing engineering clothes: someone has to say how stale is acceptable, and if nobody says, an engineer decides by accident.
Replication lag produces the failure users actually report. Write to the primary, read from a replica that is two hundred milliseconds behind, and the user submits a change and sees the old value. They do not conclude that your cluster is eventually consistent; they conclude it did not save, and they submit again. Route reads back to the primary for a short window after a write by that user, or read from the local copy you already have.
Fan-out on write versus on read is the same tradeoff at a different layer. Precompute each follower's feed when a post is made: fast reads, expensive writes, and a celebrity with ten million followers becomes an incident. Or assemble at read time: cheap writes, expensive reads. Most large systems do both, splitting on account size, which is a good reminder that the answer to a design question can be "both, on this boundary".
You should now be able to
- List the ways to serve more reads and their staleness cost
- Explain read-after-write from the user's point of view
- Choose between fan-out on read and on write
Loading…