Top-K & leaderboards

Ranking the top N of millions — sorted sets for exact scores, sketches for global heavy hitters over a stream, and how to merge per-shard answers.

Ranking is a write problem

A leaderboard looks like a read feature — "show me the top 10" — but every score update has to keep the ranking correct, and there are far more updates than reads. SELECT … ORDER BY score DESC LIMIT 10 over 50 million players re-sorts on every request; a rank query ("where am I?") is worse, because it counts everyone above you. The design question is therefore: which structure keeps the order maintained on write so reads are O(log N) or better?

Mental model: three sizes of top-K

Match the tool to the size of the set and the honesty you need. Exact and it fits in memory: a sorted set (Redis ZSET) — ZINCRBY updates a score in O(log N), ZREVRANGE 0 9 returns the top 10, ZREVRANK gives a player's position, all in one process. Fifty million entries at ~100 bytes each is 5 GB: one node, no sharding. Exact but too big for one node: partition by user or by region, keep a sorted set per shard, and merge the top-K of each shard on read — the global top 10 is always inside the union of the per-shard top 10s. Approximate over a firehose (trending hashtags, hottest products this hour): you cannot store a counter per key, so use a count-min sketch for frequencies plus a small heap of candidate heavy hitters; error is bounded, memory is fixed, and nobody notices that the 9th trending topic is really the 11th.

  • Time windows — "top this week" is a sorted set per bucket (lb:2026-w33); a sliding window is a ZUNIONSTORE of the last N buckets, recomputed on a schedule, or a decayed score where each write adds points × e^(t/τ) so old activity fades without a rebuild.
  • Ties and identity — a score of 1520 for two players needs a deterministic tiebreak; pack score × 2^20 + (MAX − timestamp) into one number so the earlier achiever ranks first and rank queries stay single calls. Count timestamp in seconds from the window's start and make MAX the window's length, so the low 20 bits cover ~12 days — widen them for a longer window, and keep the packed number under 2^53, where a Redis score (a double) stops being exact.
  • Durability — the sorted set is the serving copy; the source of truth is the events ("player 42 scored 30") in a log or database, so a lost cache is rebuilt by replay, not shrugged at.
Score events in, ranked reads out

Diagram components: Client App, Application Service, Redis, PostgreSQL.

The API writes the durable event and the ZINCRBY in one go; reads never touch the database.

Global top-K over a stream

For Trending Topics-style problems the keys are unbounded (every hashtag ever typed) and the answer is needed continuously. The stream is partitioned by key hash across processors; each keeps a count-min sketch plus its local heavy-hitter heap over a tumbling window (say 5 minutes), and a merger unions the local heaps every window into the published top 10. The sketch decides which keys deserve a real counter; the heap holds only those. Memory is a few megabytes per processor whether the stream is a thousand or a million events per second.

Sketch per partition, merge for the answer

Diagram components: Kafka, Stream Processor, Stream Processor, Application Service, Redis.

Interactive diagram — open this article in RodGrid to explore it.

Common mistakes

  • Ranking in the relational database. ORDER BY score LIMIT 10 on the players table is a full sort per read and COUNT(*) WHERE score > mine is a scan per rank query — the sorted set exists to make both O(log N).
  • One global sorted set as the source of truth. A single Redis key holding every score is a single point of loss with no replay; keep the events durable and treat the ZSET as a rebuildable index.
  • Exact counters for unbounded keys. A hash map of every hashtag ever seen grows without limit and gets slower — approximate structures are the correct tool when the key space is open.
  • Rebuilding sliding windows on every read. ZUNIONSTORE over 24 hourly buckets per request is a full merge under load; recompute on a timer and read the materialised result.

Try it on a Grid

Open the Gaming Leaderboard challenge and focus on one thing: the score-update path — one durable event write and one sorted-set update, with the read path labelled as ZREVRANGE, never a database sort. A good run shows a rebuild edge from the event store back into the sorted set. Trending Topics is the approximate variant: draw the per-partition sketch and the merge step, and write the window size on the edge.

All Learn articles