SysDesignPrep.com
Study guide 112 of 183

Counting at scale

How to count likes, views, clicks and unique visitors at high volume: exact versus approximate counts, sharded counters, batching through streams, idempotent counting, HyperLogLog for uniques, windowed counts and reconciliation.

Reading is half of it. See this used in a real interview: walk through Design an Ad Click Aggregator →

"Show the number of likes" and "count video views" sound like UPDATE posts SET likes = likes + 1. At a few writes per second that works. At thousands of writes per second on one popular item, that single row becomes the bottleneck of the whole system. Counting is a recurring sub-problem in interviews (likes, views, ad clicks, rate limits, trending), and the right answer depends on how exact and how fresh the number must be.

First, ask what the number is for

CountMust be exact?FreshnessTypical approach
Likes shown on a postno, "1.2M" is finesecondssharded or buffered counter
Video views for displaynominutesstream aggregation
Ad clicks for billingyes, and auditableminutes to hoursidempotent stream plus batch reconciliation
Unique visitorsapproximate is fineminutesHyperLogLog
Inventory remainingyesimmediatetransactional decrement
Rate limit countersapproximatelyimmediatein-memory counters per window

Saying this out loud often lets you choose a much cheaper design.

The single-row bottleneck

Every increment of one row takes a lock or a compare-and-set on one node. Throughput is capped by that row, and at high rates writes queue, time out and retry, which makes it worse. This is a write-hot key. See hot keys and skew.

Sharded counters

Split the counter into N shards (post:123:likes:0 to :15). Each increment picks a random shard; a read sums all shards (or reads a cached total refreshed every few seconds). Throughput scales roughly N times. Redis INCR on sharded keys handles very high rates in memory.

Buffer and batch

Instead of writing every event:

  1. Events go to a log or stream (Kafka), partitioned by item id.
  2. A stream processor aggregates counts per item over a short window (one to ten seconds).
  3. It writes one increment per item per window to the database.

A million likes per second on a few hot posts become a few writes per second per post. The displayed count lags by seconds, which users never notice. See message queues and streams.

Counting exactly once

Retries and redeliveries cause double counting. For counts that matter (billing, payouts):

  • Give each event a unique id at the source.
  • Deduplicate within a window (a keyed state store or a set of recent ids), or use the stream processor's exactly-once mode with transactional sinks.
  • Make the write idempotent: store per-window aggregates keyed by (item, window) and overwrite rather than increment.
  • Reconcile: a daily batch job recomputes counts from the raw event log and corrects the fast path. The batch number is the one you bill on.

See distributed transactions and idempotency and Design an Ad Click Aggregator.

Likes that a user can undo

"Did I like this?" needs per-user state, not just a counter. Store a (user, post) like record as the source of truth, which also makes liking idempotent (liking twice does nothing). The counter is a derived, eventually consistent view of those records, updated asynchronously and periodically recomputed.

Unique counts

Counting distinct users exactly needs memory proportional to the number of users. HyperLogLog estimates distinct counts within about 1 % using around 12 KB per counter, and sketches can be merged: daily uniques combine into weekly uniques. Redis has it built in (PFADD, PFCOUNT). See probabilistic data structures.

  • Fixed windows: counts per minute or hour, stored as separate keys with TTLs, so old data expires automatically.
  • Sliding windows: sum the last few fixed buckets for a cheap approximation. See Design a Rate Limiter.
  • Top-K and trending: count-min sketch plus a heap for the heaviest hitters in a stream; trending compares the recent rate with a baseline rather than raw totals.
  • Time decay for rankings like "hot" on Reddit: the score combines votes with age, recomputed periodically. See Design Reddit.

Serving the count

Reads usually outnumber writes. Cache the displayed count with the post, round it for display ("12K"), and accept that two viewers may see slightly different numbers for a few seconds. For a leaderboard of counts, a sorted set holds the ranking. See Design a Leaderboard.

Checklist

  • Decide exactness and freshness per counter.
  • Sharded counters or stream batching for hot items.
  • Unique event ids, deduplication and idempotent writes for counts that cost money.
  • Batch reconciliation from the raw log as the source of truth.
  • Per-user records for undoable actions; counters derived from them.
  • HyperLogLog for uniques, sketches for top-K, TTL buckets for windows.

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.