Course Content
System Design Interview
31 sections · 71 lessons
Gaming Leaderboard: sorted sets, scaling and follow-ups
The first lesson showed that a relational database cannot compute an arbitrary player's rank fast enough, and that no amount of hardware changes that. This lesson replaces the counting query with a data structure that answers it in about 25 steps, then covers what to say when the interviewer multiplies the scale by 100 and starts asking follow-ups.
One data structure answers every query in this problem in logarithmic time. That is the whole design.
What a sorted set is
A sorted set stores members, each with a numeric score, kept in score order. Redis is the common implementation and provides it as a first-class type; the same structure appears in other in-memory stores. The operations that matter:
| Operation | What it does | Cost |
|---|---|---|
ZADD key score member | Insert or update a member's score | O(log N) |
ZINCRBY key delta member | Add to a member's score atomically | O(log N) |
ZREVRANK key member | The member's rank, highest score first | O(log N) |
ZREVRANGE key 0 9 WITHSCORES | The top ten with their scores | O(log N + 10) |
ZREVRANGE key 4095 4105 | The ten either side of rank 4,100 | O(log N + 11) |
ZCARD key | Total members | O(1) |
With 25 million members, log₂(25,000,000) ≈ 24.6. Every one of those queries is around 25 steps, against 12.5 million for the counting version — a factor of about 500,000.
Why rank is cheap, which is the non-obvious part
A plain sorted structure gives you ordering. It does not give you rank, because knowing an element's position still requires counting the elements before it.
The trick is that each forward pointer in the skip list stores a span: how many members it jumps over. Walking down from the top level to a member, summing the spans of every pointer you traverse, yields the count of members passed — which is the rank. No counting loop: one addition at each of about 25 hops.
That single design detail is why a sorted set solves this problem and a balanced tree without span bookkeeping would not. It is worth stating precisely in the interview, because it demonstrates you know why the structure works rather than which command to call.
Why a skip list rather than a balanced tree
A skip list is a stack of linked lists. The bottom level holds every member in order; each level above holds a random subset — with a fixed promotion probability, commonly one in four — so higher levels act as express lanes. Search starts at the top, moves right while the next member is smaller than the target, and drops a level when it is not.
Expected O(log N), same as a balanced tree, without rotations. What it buys in practice:
- Range queries are trivial. Find the start, then walk the bottom-level list. A tree needs an in-order traversal with a stack.
- Simpler concurrent modification. No rebalancing means no subtree rewrites.
- Ranks by span, as above, which is what makes this problem tractable.
The honest caveat: skip list bounds are probabilistic, not guaranteed. A pathological random sequence gives worse depth. With 25 million members and a promotion probability of one in four, the expected maximum level is around 12 and the variance is small enough that it has no practical effect — but "expected" is the correct word, and using it shows precision.
Memory, and the second structure
A sorted set is two structures kept in step: a hash map from member to score, giving O(1) score lookup, and the skip list, giving order. Both must be updated on every write, which is why ZADD is O(log N) rather than O(1) and why memory per member is around 100 bytes rather than the 24 bytes the raw data would suggest.
25 million members at 100 bytes is 2.5 GB. Daily and weekly leaderboards add their own sets but only for active players — 5 million daily actives is 500 MB per daily board. A hundred country segments partition the same players rather than duplicating all of them, so the total across every scope and segment lands in the low tens of gigabytes.
Scaling beyond one node
At the stated scale, one node is sufficient and the honest answer is to say so. What follows is what you say when the interviewer raises the scale by two orders of magnitude — which they will.
When one node stops being enough
Three limits, each with a number:
Memory. 100 bytes per member means a 100 GB node holds about 1 billion members across all leaderboards. Past that, the data must split.
Throughput. A single-threaded in-memory store handles on the order of 100,000 simple operations per second per core. At 2,500 writes and 3,750 reads per second we use about 6% of one node. At 100× the players — 50,000 updates per second at peak — we are at roughly 60%, which is where planning starts.
Failure domain. Long before either limit, one node holding every leaderboard is one restart away from an outage. That, not capacity, is usually the real reason to split.
Option 1: shard by score range
Node A holds scores above 10,000, node B holds 5,000 to 10,000, node C holds below 5,000.
Rank is easy. A player's global rank is their rank within their own shard plus the total member count of every higher shard — one ZREVRANK plus a few ZCARD calls, all cheap.
Two real problems. Score distributions are not uniform and shift over a season, so a shard that held a third of players in week one holds two-thirds by week six, and rebalancing means moving members between nodes while ranks are being served. And a score update can move a player across a boundary, which is a delete on one node and an insert on another — two operations that are not atomic together, so a concurrent read can see the player twice or not at all.
Option 2: shard by player, then merge
Hash the player identifier across N nodes. Writes distribute perfectly and never move.
Top-N still works. Ask every shard for its top N, merge the N × shards results, take the top N. With 10 shards and N = 10 that is 100 entries to merge — trivial.
Exact rank does not. A player's global rank requires knowing how many members on every shard score higher. That is a ZCOUNT on each shard for scores above the player's, which is O(log N) per shard, so 10 shards is 10 parallel cheap calls. Workable, at the cost of a scatter-gather on every rank read and a result that is only as consistent as the slowest shard.
Recommendation: shard by player and accept the scatter-gather. Even distribution and stable placement are worth more than the simpler rank arithmetic, because uneven distribution and cross-shard moves are ongoing operational pain while a 10-way parallel ZCOUNT is a fixed, bounded cost.
The hot key
One global all-time leaderboard is a single key on a single node. Sharding by player does not help, because the key itself is the unit of placement — every write to that leaderboard goes to the node holding it.
This is the same problem as the celebrity fan-out in Section 13 (Design a News Feed System), the hot partition in Section 21 (Design a Distributed Message Queue), and the dominant advertiser in Section 23 (Design an Ad Click Event Aggregation). The repairs here:
- Split the key into N sub-leaderboards by player hash, exactly as in Option 2, and merge on read. This is the standard fix.
- Read replicas for the read half, since rank reads outnumber writes 1.5 to 1 and tolerate a few hundred milliseconds of staleness.
- Buffer writes. Instead of one
ZINCRBYper game, accumulate a player's points in a short window — say two seconds — and write once. At 2,500 updates per second this cuts write load substantially, and it costs two seconds of freshness.
Approximate rank for very large boards
At 100 million players, exact rank for someone in the middle is expensive on any topology and almost meaningless to the user — being 43,220,118th and 43,220,140th is the same experience.
The compromise: keep exact ranks for the top few thousand, where precision matters and is cheap, and approximate below. Approximation by score bucketing works well — maintain a count of players in each score band, so rank is the summed count of higher bands plus an interpolated position inside the player's own band. That is a handful of counter reads, independent of player count, and the error is bounded by the band width.
Present it as a product decision, not a technical fudge: exact where it matters, approximate where nobody can tell, and say which is which in the interface.
Time-windowed leaderboards and expiry
Now the follow-ups. Daily and weekly boards are separate sorted sets with the period in the key: lb:daily:2027-03-14, lb:weekly:2027-W11. A score update writes to all applicable boards — three ZINCRBY calls instead of one, which triples write load and is still only 7,500 operations per second at peak.
Expiry is a time-to-live on the key: a daily board lives eight days, long enough for late viewing and end-of-day reporting, then vanishes on its own with no cleanup job.
The subtlety worth raising: what counts as "today" for a player in a different timezone? Two answers, both defensible. A single global reset at a fixed hour is simple and unfair to some regions. Per-timezone boards are fair and multiply the number of boards by the number of timezones you support. Recommend the global reset with the reset time shown clearly in the interface, because the alternative's complexity rarely earns its keep — and say that is a product call.
Regional and friend leaderboards
Regional boards are one sorted set per country, populated by the same write. A hundred countries means a hundred boards, and the total membership is the same 25 million players partitioned rather than duplicated, so the memory cost is a modest overhead rather than 100×.
Friend leaderboards are different, and interviewers like them because the obvious approach fails. You cannot keep a sorted set per friend group — with 5 million players and 200 friends each, that is 5 million sets, and every score change updates 200 of them.
Instead, compute friend leaderboards on read: fetch the player's friend list, fetch each friend's score with one bulk ZMSCORE call against the global board, and sort the result in the application. For 200 friends that is one round trip and a 200-element sort — microseconds. This is fan-out on read, and it is the correct choice for exactly the reason Fan-out on write versus fan-out on read gives: the result set is small and the write amplification of precomputing it is enormous.
Tie-breaking with a composite score
Two players on 9,200 points. The earlier achiever should rank higher.
The elegant implementation packs both values into one number. A sorted set score is a double-precision float with 53 bits of integer precision. Use the high bits for the points and the low bits for an inverted timestamp:
composite = points * 2^20 + (2^20 - 1 - minutes_since_season_start)With 20 bits for the tiebreaker, a season of up to about a million minutes — roughly two years — is representable, and points up to 2³³ still fit within 53 bits. Sorting by the composite gives points-descending, then earliest-first, in one comparison, with no secondary lookup.
The caveat: verify the bit budget for your actual ranges, and be aware that the displayed score must be composite >> 20 rather than the stored value. Getting the display wrong here is an easy and embarrassing bug.
Anti-cheat
A leaderboard is a target. If the client posts its own score, the leaderboard is fiction.
Layered defences: validate server-side that the score is achievable for the game session — duration, actions, and physics all constrain the maximum; rate-limit submissions per player; flag statistical outliers for review rather than auto-rejecting, since real prodigies exist; and keep the raw game session events so a disputed score can be re-examined. Note that this is the same raw-events-for-recomputation argument as ad click reconciliation.
Durability, given an in-memory store
The uncomfortable question: a season's standings live in memory, and memory is lost on restart.
The wrong answer is to rely solely on the in-memory store's own persistence. Point-in-time snapshots lose everything since the last snapshot; append-only logging with a per-second flush loses up to a second. For a leaderboard, a second of lost score updates is survivable but untracked, which is worse than it sounds — you do not know what you lost.
The right answer is that the sorted set is a derived index, not the source of truth. Every score change is written first to a durable store — the game's own event log or transaction database — and then applied to the sorted set. If a leaderboard node is lost, rebuild it by replaying scores from the durable store: 25 million ZADD operations at 100,000 per second is about four minutes.
That reframing appears repeatedly in this course: keep the authoritative record in something durable and treat the fast structure as a rebuildable cache. It is the same relationship as raw events to aggregates in Section 23 (Design an Ad Click Event Aggregation), and metadata to search index in Section 25 (Design a Distributed Email Service).