A viral post's like counter in a sharded cache takes 40,000 increments per second on one key; how do you split that write load across nodes?
answer
- replicas do not take writes
- one key becomes several
- random suffix on increment
- sum on read, cache the sum
- growing N versus shrinking N
basics
~20 sReplace the single counter with N sub-counters whose suffixed keys hash to different nodes. Each increment goes to one randomly chosen sub-counter, and a read sums all N. Write load per sub-counter drops to about 1/N, but each read now fetches N keys.
solid answer
~40 sA write-hot key can't be fixed with read replicas, because every increment still goes to the single primary that owns the key. Instead, shard the counter itself: `likes:post:42#0` through `likes:post:42#7`. Each increment picks a random suffix, so 40,000 increments per second become roughly 5,000 per sub-counter. A read fetches all eight and adds them up, which costs an 8-key multi-get. You can cache that sum briefly so reads don't multiply. Check where the suffixes actually land, because with few nodes some sub-counters will share a node. If readers always sum up to the largest N ever used, growing N is safe. A cheaper alternative is to buffer increments in each app process and flush one accumulated delta every few hundred milliseconds. That trades a small loss window on a crash for far fewer writes.
go deeper
Recall that a key written very often can overload its node, and that splitting it into several smaller counters spreads the writes.
Walk through the write path with a random suffix and the read path that sums N keys, and show the per-node arithmetic for a chosen N.
Cover the operational details: placement collisions, safely growing or shrinking N, flushing to durable storage, and caching the summed read.
Decide how much counting accuracy the product needs, and pick buffered deltas, sub-counters or a durable event log based on the acceptable loss window and read cost.
## Why a counter is a different kind of hot key Most hot-key advice assumes a **read-hot** key: many reads, few writes. You can copy or cache such a value because it rarely changes. A counter for likes, views or a live score is **write-hot**: - every event is an increment on the same key - the owner node must apply each increment in order, so this is the busiest possible write path - in a primary-replica setup only the primary accepts writes, so **read replicas add nothing** for the increment load - a local cached copy doesn't help either, because the problem is the writes So the value itself must be spread out. ## Sub-counters A **sharded counter** replaces one key with N sub-keys: 1. **Write:** pick `r = random(0, N)` and increment `counter#r`. 2. **Read:** fetch `counter#0 … counter#N-1` in one multi-get and sum them. 3. **Place:** each suffixed key hashes on its own, so the sub-counters spread across nodes. Illustrative sizing, assuming one node can absorb about 10,000 increments per second for this key with headroom: - 40,000 / 10,000 means at least 4 sub-counters - choosing 8 gives about 5,000 increments per second each, which leaves room for growth The totals are the same because each increment lands in exactly one sub-counter. ## The costs | Aspect | Single counter | 8 sub-counters | |---|---|---| | Increments per key | 40,000/s | ~5,000/s each | | Keys fetched per read | 1 | 8 | | Read consistency | One value | Sum of 8 values read at slightly different moments | | Operational state | None | N must be known to every reader | Keep a few points in mind: - **Read fan-in.** A page that shows the count now does an N-key fetch. Cache the summed value for a second or so, because a like count doesn't need to be exact to the millisecond. - **Placement collisions.** Suffixes hash independently, so two sub-counters can land on the same node. With 8 sub-keys on 20 nodes and uniform hashing, the chance that all 8 get distinct nodes is only about 20%, so most of the time at least two share a node. Check placement, or choose suffixes that map to distinct nodes. - **Approximate reads.** The sum is taken across keys at slightly different moments. It can be a few increments behind, which is acceptable for a like count and not for money. ## Changing N - **Growing** from 4 to 8 is safe if readers sum all 8. The new sub-counters start at zero and the old ones keep their counts, so the total is unchanged. - **Shrinking** needs a fold step: move each retired sub-counter's value into a remaining one before readers stop summing it. - Store N somewhere every reader and writer can see, or keep it fixed per key class. ## Durability A cache counter is volatile. The usual pattern is to **flush** the summed total to a durable store on a schedule, or to append each like to a durable log and treat the cache counter as a fast view that can be rebuilt. ## The buffered-delta alternative Instead of one cache write per event, each application process can keep a local in-memory delta and flush it periodically: 1. On a like, add 1 to the process-local delta for that post. 2. Every 500 ms, send one add of the accumulated delta to the shared counter and reset the local value. With 100 processes flushing twice a second, the shared key sees at most 200 writes per second instead of 40,000. The price is that a process crash loses up to 500 ms of its unflushed increments, and the displayed count lags by up to the flush interval. The two techniques combine well: buffered deltas cut the write rate, and sub-counters spread what remains.
- How would buffering increments in each process change the load?Each process adds likes to a local delta and flushes one combined add every few hundred milliseconds. With 100 processes flushing twice per second, the shared counter sees at most 200 writes per second. The cost is a loss window equal to the flush interval if a process crashes, plus a small display lag.
- Why is a sharded counter a poor fit for an account balance?The read sums sub-counters fetched at slightly different moments, and increments are only as durable as the cache. That is fine for approximate like counts. A balance needs exact, durable, atomic updates with invariants such as never going below zero, which belong in a transactional store.
saying these in an interview costs you the question
- Add read replicas to absorb the increment load
- Cache the counter locally on each app server
- Every increment must be written to all sub-counters
- Suffixed sub-keys are guaranteed to land on distinct nodes
- Shrinking the sub-counter count needs no data migration