System Design Interview

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:

OperationWhat it doesCost
ZADD key score memberInsert or update a member's scoreO(log N)
ZINCRBY key delta memberAdd to a member's score atomicallyO(log N)
ZREVRANK key memberThe member's rank, highest score firstO(log N)
ZREVRANGE key 0 9 WITHSCORESThe top ten with their scoresO(log N + 10)
ZREVRANGE key 4095 4105The ten either side of rank 4,100O(log N + 11)
ZCARD keyTotal membersO(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.

A skip list: express lanes over a sorted listlevel 2120span 3480span 3940level 1120span 2310span 1480span 3940level 0 —everyplayer120span 1240span 1310span 1480span 1520span 1690span 1940anabocydiedfigurank(di) = 1 + 3 = 4thwalk from the head at the top level, adding each span you traverse — that running total is the rank, computed without counting the nodes below.Insert, delete and rank are all O(log n), which is why a sorted set is the right primitive for a leaderboard and a sorted SQL query is not.
The spans are what turn "where am I on the leaderboard" from a full scan into a walk down three levels.

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.

Two ways to split a sorted setShard by score range• Rank is a sum of shard counts• Ranges drift and need rebalancing• Updates move players between shardsShard by player, then merge• Even distribution, trivially• Top k merges cheaply across shards• Exact rank needs every shard
Exact global rank is what resists sharding, which is why very large boards quote an approximate rank instead.

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 ZINCRBY per 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.

Five extensions worth preparingAfter thesorted setDaily and weekly keysFriend leaderboardsComposite tie-breaksAnti-cheat at ingestRebuild from the log
Writing three sorted sets instead of one triples the write rate and is still only seven thousand operations a second.

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:

Text
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).