End-to-End Designs
End-to-end designs come down to a few recurring moves: make the critical write atomic, give each ordering domain a single owner, choose fan-out on write or on read deliberately, and split hot keys when one entity outgrows a shard.
Key points
- 1
Check-then-act across separate statements is a race. Use a conditional
UPDATE ... WHERE <expected state>and check the affected row count, or lock the row in the same transaction. - 2
Model timeouts as data (
expires_at) so that correctness doesn't depend on a sweeper job running on time. - 3
Order per entity, not globally: route each room, account or order to one owner that assigns increasing sequence numbers, and let clients sync with "everything after N".
- 4
Fan-out on write is cheap to read but expensive for huge audiences; fan-out on read stores once and costs more at read time. Hybrid designs treat celebrities and big rooms specially.
- 5
One hot key caps throughput at one partition and one consumer. Salt the key into sub-keys and add a second stage that merges mergeable partial aggregates.
- 6
Exactly-once effects across services need an idempotency key created by whoever retries, and passed through to every system that performs the side effect.
Common traps
Adding partitions doesn't help a single hot key: hashing still sends it to one partition.
Averages of averages and averages of percentiles are wrong. Merge (sum, count) pairs or quantile sketches instead.
Client or server wall clocks can't give a shared message order. Use sequence numbers from a single owner.
Read the source
Test yourself on End-to-End Designs
Ten questions, with the answer and explanation after each one.