Counting at scale: top-K and rank
Counting what you can't count exactly: sharded counters, heavy hitters and top-K with sketches, distinct counts, and rank on boards with millions of entries.
Counting at scale is choosing which answers must be exact and which can be approximate. This page gets you ready for hot counters, heavy hitters, distinct counts and global rank when one machine is too small.
Read this if your last attempt…
- Your last leaderboard answer was just "put ORDER BY score DESC LIMIT K on a table"
- You counted every item in a stream even though the key space was unbounded
- You merged shard top-K lists and assumed the result was exact
- You used an approximate sketch for a money, quota or privacy decision
The concept
Counting starts with one question: what decision will use the count? If the count moves money, blocks a user, enforces a privacy threshold, or decides a contest prize, make it exact. If it drives a trending module, a dashboard hint, a cache warm-up list, or a long-tail percentile, an approximate answer can be better because it stays cheap and fresh.
Exact counters first
A single row counter is the simplest model and the first write hotspot. Every increment of a viral key contends on the same row, lock, partition or Redis slot. Buffering helps: collect deltas in memory, a log consumer, or a small queue, then write larger deltas. Striping helps a hotter key: write to salted counter rows, then sum the stripes on read. That is the same hot-key idea as sharding and partitioning: spread writes, then pay a read-side merge cost.
Hot writes first land in a log or buffer. Exact counters protect correctness-sensitive totals. Sketches and bounded candidate sets find heavy hitters. Sorted sets and histograms serve rank queries.
Counting choices: match precision to the decision that consumes the answer.
| Need | Good fit | Watch out for |
|---|---|---|
| Money, quotas, privacy thresholds | Exact counter, ledger, or replayable aggregate | One hot row becomes the bottleneck unless you batch or stripe. |
| Viral likes or views | Striped exact counter with read-side sum | Reads pay a merge cost and stale cache policy must be explicit. |
| Heavy hitters in an unbounded stream | Count-Min Sketch plus Space-Saving style candidates | Sketches estimate; they do not prove a strict top-K without bounds or verification. |
| Unique viewers at scale | HyperLogLog per window, merged across shards | You get cardinality only, not the member list. |
| Exact leaderboard rank | Redis sorted set or materialised rank table | Memory and cross-shard rank counts become the scaling limit. |
| Long-tail percentile rank | Score histogram plus exact top-N | Bucket width controls error, and ties still need a rule. |
How interviewers grade this
- You separate exact counters from approximate discovery paths before choosing a data structure.
- You cool a hot key with batching, stripes or shard-aware ownership instead of hammering one row.
- You describe Count-Min Sketch as over-estimating and pair it with a candidate set for top-K.
- You say that shard-local top-K merge can miss a globally heavy item when counts are split across shards.
- You use HyperLogLog for unique counts that tolerate error, and keep exact sets for thresholds that matter.
- You define rank ties and explain how a global rank works after the board is sharded.
Variants
Striped exact counters
Split writes for one hot logical counter across several physical rows or keys.
Use stripes when the value must be exact but one key is too hot. Writers pick a stripe, increment it, and readers sum stripes. This is a good fit for likes, views, rate counters and per-object engagement where read freshness can tolerate a small merge or cache delay. Keep the number of stripes a tuning knob, because every extra stripe cools writes and adds read work.
Pros
- +Preserves exactness for the logical counter.
- +Spreads write contention across partitions.
- +Works with SQL rows, key-value stores and Redis keys.
Cons
- −Reads become a fan-out and sum.
- −Increasing stripes later needs a migration or versioned read path.
- −It does not solve global top-K by itself.
Choose this variant when
- A small set of keys is hot and the answer must stay exact.
Count-Min Sketch
A mergeable frequency sketch with one-sided over-estimates when stream updates only add.
Use Count-Min Sketch when the key space is too large for exact per-item counts and false positives are acceptable. It is useful for candidate discovery, abuse signals and cache warm-up lists. Because collisions only add counts, treat a high estimate as "worth checking", not as a billable fact. Sketches merge naturally by adding corresponding counters, which is why they fit sharded streams.
Pros
- +Bounded memory independent of distinct-key count.
- +Fast stream updates and mergeable shard summaries.
- +One-sided error is easy to explain.
Cons
- −Can over-rank colliding keys.
- −Needs a companion candidate set to emit top-K.
- −Total stream weight affects the absolute error.
Choose this variant when
- You need cheap frequency estimates for a very large stream.
Space-Saving / Misra-Gries candidates
Keep a bounded set of candidate heavy hitters and evict the smallest counter when a new key arrives.
Use candidate algorithms when the product asks for "what is hot right now?" rather than "what is the exact count for every item?" They hold a configured number of counters. Existing candidates increment; new keys fill empty slots; once full, the smallest counter is replaced with an error bound. This gives you a shortlist to verify or serve as approximate top-K.
Pros
- +Memory is controlled by the counter budget.
- +Outputs candidates directly.
- +Works well inside time windows.
Cons
- −The candidate list is approximate unless you verify counts.
- −Small budgets miss items near the cutoff.
- −Window resets and decay policy change the result.
Choose this variant when
- Trending lists, hot hashtags, top search queries and abuse-heavy hitters.
HyperLogLog distinct counts
Approximate unique users, devices or viewers with fixed-size, mergeable state.
Use HLL when the question is cardinality: how many unique actors saw this item in this window? Each shard can maintain an HLL per key and window, then merge those HLLs for the global estimate. Do not use it when you need to list the users or enforce a hard privacy threshold; keep an exact set or thresholded ledger for that path.
Pros
- +Small fixed memory per sketch.
- +Mergeable across shards and windows.
- +Good for dashboards and trending features.
Cons
- −Approximate by design.
- −Cannot return the members.
- −Wrong primitive for hard thresholds.
Choose this variant when
- Unique viewers, daily active users by segment and approximate reach.
Sorted-set rank
Maintain an ordered set by score and ask for top ranges or a member rank.
A player-sharded board can update cheaply, but a global rank asks every shard how many scores are higher. A histogram gives a cheap percentile estimate for the long tail while the top-N remains exact.
Use a sorted set when the board fits the memory and durability story. It is excellent for exact top ranges, neighbourhood reads and a player's current rank. At very large scale, treat the sorted set as a derived serving view over a durable score source. Shard by score range when rank queries dominate; shard by player when writes and ownership dominate.
Pros
- +Exact rank inside one board.
- +Simple top-N and around-me queries.
- +Easy to rebuild from durable scores.
Cons
- −RAM is the ceiling.
- −Global rank across shards needs per-shard counts.
- −Tie semantics must be explicit.
Choose this variant when
- Game leaderboards, contest boards and live score tables.
Worked example
Numbers in this section are illustrative.
Scenario: design the counting layer for a consumer app with trending posts and a game-style points board. Numbers in this section are illustrative.
Requirements. The product wants top 100 trending posts per region for the last hour and day, unique viewers per post, exact like counts on the post page, and a player rank for 80M players. Trending can be approximate for a few minutes. Likes that affect creator payouts and player prize ranks must be exact.
Write path. Every play, like and score event goes to a durable log. A consumer coalesces repeated increments before touching storage. Post likes use striped exact counters: post:{id}:like:{stripe}. The post page sums stripes and caches the result briefly. Creator payout jobs read the durable log or exact aggregate, not a sketch.
Trending path. For each region and hour, each stream shard keeps a Space-Saving candidate set and a Count-Min Sketch. The shard emits more than 100 candidates because local top-K merge can miss a globally steady item. The merge service dedupes candidates and fetches exact counts. It publishes a strict top 100 only if shard cutoffs prove no omitted item can catch up, or after a second pass over the exact window aggregate. During an incident, it can fall back to approximate order with a visible freshness flag.
Unique viewers. Each shard keeps hll:post:{id}:hour:{h} and hll:post:{id}:day:{d}. The region service merges HLLs for approximate uniques. If a metric crosses a privacy threshold, the threshold check uses an exact suppressed-count table, not HLL.
Rank. The top 1000 players are exact in a global Redis sorted set backed by a durable scores table. The long tail uses a score histogram per shard for approximate percentile and rank band. If the interview requires exact global rank for every player, I shard by player and query every shard for count(score > my_score), then apply the tie policy: competition rank adds 1, deterministic order counts equal-score players ahead, and dense rank counts distinct higher scores. That is slower but exact.
A player-sharded board can update cheaply, but a global rank asks every shard how many scores are higher. A histogram gives a cheap percentile estimate for the long tail while the top-N remains exact.
Tie rule. Equal scores share a competition rank for public display. Prize settlement breaks ties by earliest time reaching the score, then by player id, and writes the final order to an immutable results table.
Good vs bad answer
Numbers in this section are illustrative.
Interviewer probe
“You need the top 100 trending hashtags and exact player rank on a board with millions of players. What do you build?”
Weak answer
"I will count everything in a database table, sort it by count, and cache the result. For rank I will use Redis and shard it if it gets big."
Strong answer
"I split the problem. For trending hashtags, exact counts for every tag do not fit the stream path, so each shard keeps windowed heavy-hitter candidates with Space-Saving and a Count-Min Sketch. The merge service over-fetches candidates, then verifies exact counts for that candidate set before publishing if the product needs a strict top 100.
I would not claim that merging local top 100 lists is exact, because an item can be moderate on every shard and absent locally. For rank, if the board fits in memory I use a Redis sorted set with ZADD or ZINCRBY for scores, ZREVRANK for a player rank, and ZREVRANGE for the top page. If it no longer fits, I keep exact top-N and use either per-shard rank counts or a histogram estimate for the long tail, with a stated tie policy."
Why it wins: It separates approximate discovery from exact rank, names the algorithms, catches the shard-merge trap, and states what happens when one sorted set is too large.
When it comes up
- The prompt says leaderboard, trending, top-K, most popular, heavy hitters, unique viewers or global rank.
- The event stream has too many distinct keys to count every key exactly on the hot path.
- A single object can become viral and make one counter key hot.
- The interviewer asks whether a merged shard result is exact.
Order of reveal
- 11. Classify the count. First I separate exact decisions from approximate discovery.
- 22. Cool hot exact counters. For exact hot keys I batch increments or stripe counters, then sum on read.
- 33. Pick stream summaries. For heavy hitters I use bounded candidates, often with Count-Min Sketch estimates.
- 44. Handle windows and shards. Trending is windowed; shard-local top-K merge is approximate unless I verify candidates.
- 55. Rank deliberately. One sorted set gives exact rank while it fits; sharded boards need per-shard counts or approximate histograms.
Signature phrases
- “Exact for enforcement, approximate for discovery.” — Shows the main safety boundary.
- “Local top-K lists are candidates, not proof.” — Catches the common shard-merge error.
- “A rank is count-above plus a tie policy.” — Makes global rank concrete after sharding.
Likely follow-ups
?“Why does Count-Min Sketch over-estimate?”Reveal
Every update only increments counters. Hash collisions add other items into the same counters, so the minimum row can still be above the true count, but not below it.
?“How do you make trending decay instead of resetting every hour?”Reveal
Use a score such as weighted recent counts, with older buckets multiplied by a decay factor, or maintain several windows and blend them. I would still keep exact enforcement counters separate.
?“What if the global board no longer fits in one Redis sorted set?”Reveal
Keep the durable score source, then shard the serving view. Exact global rank asks shards for counts above a score. For the long tail, a histogram can give a percentile band while the top-N remains exact.
Code examples
ZADD leaderboard 4820 player:7
ZINCRBY leaderboard 5 player:7
ZREVRANK leaderboard player:7
ZREVRANGE leaderboard 0 99 WITHSCORESfunction observe(key: string) {
const hit = counters.get(key);
if (hit) return counters.set(key, hit + 1);
if (counters.size < budget) return counters.set(key, 1);
const [victim, count] = smallestCounter(counters);
counters.delete(victim);
counters.set(key, count + 1); // remember count as this key's error bound
}Common mistakes
If the same item can receive events on many shards, local top-K lists can hide a globally heavy item. Over-fetch candidates and verify, or use a sketch threshold that proves no hidden item can win.
Approximate counts are fine for discovery and dashboards. They are the wrong authority for payouts, bans, inventory, prize ranks or privacy thresholds unless a separate exact path confirms the decision.
A single hot row serialises all increments. Batch deltas, stripe the counter, or route ownership so the write path can scale while reads sum or cache the result.
Lifetime top-K favours old winners. Trending needs a named hour, day, sliding window or decay curve, plus an expiry and late-event policy.
A sorted set is elegant while it fits. Once the board is very large, memory, rebuild time and cross-shard rank counts become first-class design constraints.
Equal scores happen often. Say whether ranks are shared, dense, competition-style, or broken by time and id. Settlement paths should persist the final tie order.
Practice drills
Numbers in this section are illustrative.
Why is merging each shard's top 100 not enough for a global top 100?Reveal
Because the same item can be below the local cutoff on every shard but high after all shard counts are summed. Local lists are candidates. Over-fetch and verify counts, or use a bound that proves no hidden item can win.
When would you choose HLL over an exact set?Reveal
Use HLL for approximate unique counts such as viewers or daily active users when you only need cardinality and can tolerate error. Use an exact set or ledger when you need membership, deletion semantics, or a hard threshold.
How do you rank a player on a player-sharded board?Reveal
Find the player's score on its owning shard, ask every shard how many members have a higher score, then apply the tie policy. Shared competition rank is one plus that higher-count total; a deterministic tie-break also counts equal-score members ahead of the player; dense rank counts distinct higher scores. That is exact but costs a fan-out.
What is the safe answer for creator payouts based on views?Reveal
Use exact, replayable aggregates from the event log or an authoritative counter table. Sketches can suggest suspicious or hot items, but the payout path should verify exact counts.
Deep dives
Why shard top-K merge is only a candidate step
Each shard keeps bounded candidates for the current window. The coordinator over-fetches candidates, fetches exact counts when it needs a strict answer, and publishes the window result.
There are two different shard shapes. If items are partitioned by item id and every event for an item goes to the owning shard, the shard's exact local top-K contains any item that could be global top-K from that shard. The harder and more common stream shape is different: events are partitioned by user, region, producer, Kafka partition or time. Then one item's count is split across shards.
In the split-count shape, a globally hot item can be just below the cutoff everywhere. Asking each shard for exactly K items can drop it from the candidate pool before the coordinator ever sees it. Over-fetching helps because it widens the candidate set. It is still a heuristic unless you have a bound: for example, each shard also returns the local cutoff count, and the coordinator can prove that the sum of possible hidden counts cannot exceed the current global Kth count. If you cannot prove that, scan the exact per-window aggregate in a second pass, or accept and label the result as approximate.
Approximate rank from score histograms
A player-sharded board can update cheaply, but a global rank asks every shard how many scores are higher. A histogram gives a cheap percentile estimate for the long tail while the top-N remains exact.
A histogram keeps counts per score bucket. To estimate rank for score S, add the counts in buckets above S, then add a fraction of S's bucket depending on the product rule. This is cheap to merge across shards because histograms add bucket by bucket.
The trade-off is bucket error. Wide buckets are cheap but make rank bands fuzzy. Narrow buckets cost more memory and update work. Use this for long-tail "you are around the 82nd percentile" feedback, not for awarding prizes. A practical design keeps exact top-N in a sorted set or table, stores final prize ranks immutably, and uses histograms only where an approximate band is acceptable.
Aggregating ratings
Ratings look like a counter, but the read path usually wants both a displayed average and a ranked list. Keep the serving row boring: rating_sum, rating_count, and optionally a small histogram by star value. The displayed mean is rating_sum / rating_count when rating_count > 0. For write-hot items, reuse the write-cooling ideas in striped exact counters, then publish a compact per-item aggregate.
Writes, edits and deletes
Treat a review write and its aggregate delta as one logical change. On create, add the new rating and increment the count. On edit, add new_rating - old_rating and leave the count unchanged. On delete, subtract the old rating and decrement the count. Do this in the same database transaction as the review row when the aggregate lives in the same store. If the aggregate is maintained by a worker, write an outbox event in the review transaction and make the consumer idempotent with a stable review id and version, so replay does not double-apply a delta.
Aggregates drift when bugs, retries or manual fixes bypass the normal path. Schedule periodic recounts from the review table or event log and repair differences. Keep the recount scoped, for example recently changed items first, then a slower full sweep. Fraud and spam filtering should happen before a rating counts. A pending or rejected review can exist for moderation, but it should not update rating_sum until it is eligible for ranking.
Ranking is not the raw mean
Do not sort by raw mean when items have different review counts. A place with one 5-star review should not outrank a place with thousands of strong reviews by default. A simple Bayesian-style score is (C·m + Σ ratings)/(C + n): m is the global prior mean, C is the strength of that prior measured as pseudo-reviews, Σ ratings is the item's rating sum, and n is its rating count. It is the mean after adding C pseudo-reviews at prior mean m.
As an illustrative example, with m = 4.0, C = 20, and one 5-star review, the score is about 4.05 rather than 5.0. Evan Miller's Bayesian Average Ratings, published on 2012-11-06, explains the broader prior-and-update intuition for sparse rating lists.
For binary up/down votes, a lower-bound confidence score is a good alternative. Evan Miller's How Not To Sort By Average Rating, published on 2009-02-06, recommends the lower bound of the Wilson score confidence interval. With p̂ = positive / n, confidence multiplier z, and total n, the lower bound is (p̂ + z²/(2n) - z·sqrt((p̂(1-p̂)+z²/(4n))/n))/(1+z²/n). It balances observed approval with uncertainty from small samples.
Keep board mechanics separate: sorted-set rank covers exact serving when a ranked board fits in memory, and approximate rank from score histograms covers percentile-style feedback. The rating card's job is to produce a fair score for each item, not to repeat the board data structure.
Cheat sheet
- •Start by asking which counts must be exact.
- •A single hot counter row is a write hotspot; batch or stripe it.
- •Count-Min Sketch over-estimates; use it for estimates, not enforcement.
- •Space-Saving / Misra-Gries gives bounded-memory heavy-hitter candidates.
- •Shard-local top-K merge is not exact when item counts are split across shards.
- •HyperLogLog estimates unique counts and merges well.
- •Redis sorted sets give exact rank while the board fits in memory.
- •Global rank after sharding is count-above plus tie policy.
- •Use histograms for long-tail approximate rank and exact top-N for winners.
Practice this skill
These problems exercise Counting at scale: top-K and rank. Try one now to apply what you just learned.
Read this if