Command Palette

Search for a command to run...

Hectal
PHASE 7Intermediate ~10 min· topic 4 of 4

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.

0/4 · 0%

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

  1. 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.

  2. 02

    Sharded counter: write to likes:post:9:{0..N-1} picked randomly (or by user hash), read with MGET and 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}).

  3. 03

    Buffered counter: each app instance accumulates increments in memory and flushes INCRBY every 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.

  4. 04

    HINCRBY counters in a hash (stats:post:9 with fields likes, shares, views) keep related counters together and compact.

  5. 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.

  6. 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

sharded-counter.pypython
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 total
counter-options.mermaiddiagram
Rendering diagram…

Interview 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

01

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

01

Compare single, sharded and buffered counters.

Practice

01

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.