Topic 7.4
Distributed Counters: Single, Sharded, Approximate
In one line
A single INCR key is exact and simple until one key gets too hot. Sharded counters spread writes across N keys and sum on read; buffered counters aggregate locally and flush; approximate counters (HyperLogLog, Count-Min) trade exactness for tiny memory. Persist counts that matter to a durable store.
Think of it like this
Counting votes at a big election. One counting table for the whole country is exact but hopelessly slow. Many local tables count in parallel and report subtotals, and the national total is the sum. Exit polls (approximate) give a good answer much sooner.
Key ideas
- 01
Single counter:
INCR likes:post:9. Exact, atomic, O(1). Fine up to tens of thousands of increments/sec on one key, but a viral post pins one shard's CPU. - 02
Sharded counter: write to
likes:post:9:{0..N-1}picked randomly (or by user hash), read withMGETand sum. Writes scale with N; reads cost N lookups, so cache the sum for a second or two. To spread across cluster nodes, make sure the shards hash to different slots (put the shard number inside the hash tag,likes:{post:9:3}). - 03
Buffered counter: each app instance accumulates increments in memory and flushes
INCRBYevery second. 1,000 instances → at most 1,000 writes/sec per counter no matter the traffic. Costs up to 1 second of lag and loss of buffered counts on crash. - 04
HINCRBYcounters in a hash (stats:post:9with fieldslikes,shares,views) keep related counters together and compact. - 05
Uniqueness: "likes" usually means one per user, so you need a set or a database row per like to prevent double-liking; the counter is derived from that. "Views" are usually just increments.
- 06
Durability: counters that matter (billing, payouts, visible like counts) should be persisted: periodic snapshots to the database, or derive them from an event log (Kafka or a Redis Stream) that can rebuild the count.
Code & diagrams
import random, redis
r = redis.Redis()
N = 16
def incr_like(post_id: int):
shard = random.randrange(N)
r.incr(f"likes:{{post:{post_id}:{shard}}}") # different tags -> different slots
def like_count(post_id: int) -> int:
cache_key = f"likes:total:{post_id}"
cached = r.get(cache_key)
if cached is not None:
return int(cached)
keys = [f"likes:{{post:{post_id}:{s}}}" for s in range(N)]
pipe = r.pipeline(transaction=False) # cluster client splits by node
for k in keys:
pipe.get(k)
total = sum(int(v or 0) for v in pipe.execute())
r.set(cache_key, total, ex=2) # cache the sum briefly
return totalInterview problem
The problem
Social media likes at 100M scale
Design like counts: 100M likes per day overall, some posts go viral with 200K likes per second, each user can like a post once, and counts are shown on every post view. Compare a single counter, a sharded counter and database aggregation.
You're given
- 100M likes/day
- Viral peaks 200K likes/sec on one post
- One like per user per post
- Counts visible on 1M views/sec
The interviewer follows up
Why not store each like as a set member and use SCARD as the count?
When it breaks
Sharded counter shards all share one hash tag
What you see
All shards land in the same slot and node, so sharding spreads keys but not load; the hot shard stays hot.
Fix & prevent
Put the shard number inside the hash tag (or don't use a tag), verify with CLUSTER KEYSLOT that shards map to different nodes.
Explain it without notes
Compare single, sharded and buffered counters.
Practice
Benchmark INCR on one key vs 16 sharded keys from 50 concurrent clients on a small cluster; compare per-node CPU.
Trade-offs
- ↔
More shards mean more write capacity and more expensive reads.
- ↔
Buffering reduces load but adds lag and needs a durable path for correctness.
Done when you can
I can design counters that survive viral traffic.
I separate uniqueness (sets or database) from counting.
I persist or can rebuild counters that matter.