skip to content

A leaderboard in Redis holds tens of millions of players and product now wants every player's exact global position plus deep pagination through the standings. Which parts of that stay cheap, which do not, and how would you architect around the limits?

level: principalimportance: should knowfreq 34%

answer

  1. ZREVRANK O(log N) — exact rank stays cheap
  2. index paging cost includes the offset walked
  3. cursor by score, not by offset
  4. one key = one node: memory + hot-key ceiling
  5. score-bucket histogram = approximate global rank

basics

~20 s

Exact rank and top-N stay cheap — both are O(log N). Deep pagination by index is not: cost grows with the offset. Memory and single-key hotness are the real ceilings. Paginate by score cursor, serve reads from replicas, shard into cohort boards, and approximate global rank with score-bucket counters.

solid answer

~1 min

**Stays cheap.** `ZREVRANK` is O(log N) — exact global rank over 50 million players is still microseconds, because the ordered structure keeps subtree sizes. Top-N and an "around me" window are O(log N + M) with M tiny. **Does not.** `ZRANGE key 200000 200049` is O(log N + offset + count): the engine walks the skipped elements. Deep pages get linearly worse and block the server while they run. Replace offsets with a **score cursor**: `ZRANGE key (lastScore -inf BYSCORE REV LIMIT 0 50`, carrying the last score (plus member for tie-breaking) in the page token. **Real ceilings.** - **Memory** — tens of millions of members is gigabytes on one node; it also lengthens snapshot forks and replica full-sync time. - **Hot key** — one Sorted Set cannot be split across cluster nodes, so a global board is one node's CPU and RAM. - **Write amplification** — every scoring event is an O(log N) write that must replicate. **Architecture.** Reads to replicas (a slightly stale rank is fine). Shard into per-region/per-cohort boards and merge top-N client-side. For a global rank across shards, keep a **score-bucket histogram** and compute rank as "players in higher buckets + local rank within the bucket" — approximate at the boundaries, O(buckets) to evaluate.

code

text · 10 lines
text
# expensive: cost grows with the offset walked
> ZRANGE lb 200000 200049 REV WITHSCORES

# constant cost per page: resume from the last score seen
> ZRANGE lb "(1450" "-inf" BYSCORE REV LIMIT 0 50 WITHSCORES

# cheap regardless of size: exact rank + a window around one player
> ZREVRANK lb player:42
(integer) 411998
> ZRANGE lb 411996 412000 REV WITHSCORES

go deeper

for a junior

Know that rank and top-N lookups are logarithmic and cheap, and that reading large ranges is proportional to how much you read.

for a middle

Explain why offset-based paging degrades with depth and how a score cursor fixes it, and note that one key lives on one node.

for a senior

Cover the operational ceilings — memory driving fork, replica sync and failover; write amplification per scoring event; replica-served reads — and propose time partitioning with TTLs.

for a principal

Make the exact-versus-approximate trade explicit: exact ranks for the visible top slice, bucket-histogram approximation below it, cohort sharding with client-side merges, and a written statement of which guarantees the product is buying at what cost.

## First, separate what actually hurts The instinct is that ranking fifty million players must be expensive. It is not — at least not the part people expect. **Rank lookup is logarithmic.** The ordered structure behind a Sorted Set stores, at each level, how many elements a link spans. Counting the elements before a given member is therefore a descent, not a walk: `ZREVRANK` is O(log N). Exact global position for one player is cheap at any realistic size. So is `ZSCORE` (O(1)), `ZCARD` (O(1)), `ZCOUNT` (O(log N)) and top-N (O(log N + N)). **Range reads by index are the expensive shape**, because their cost includes everything skipped: `ZRANGE key <start> <stop>` is O(log N + M) where M spans the *whole window from the start index*. Page 4 000 of a 50-element pager means walking 200 000 entries inside one command, during which the server serves no one else. Deep pagination is thus not a data-size problem but a **latency blast-radius** problem: one user's request degrades everyone. ### The fix: cursor pagination by score Carry the last row's score and member in the page token and query `ZRANGE key (lastScore -inf BYSCORE REV LIMIT 0 50`. Each page costs O(log N + 50) regardless of depth. Because equal scores are common on leaderboards, the cursor must include the member to disambiguate the boundary — the standard trick is to filter out members already emitted at the boundary score, or to make scores unique by packing a tiebreaker in (see below). Cursor paging is also more stable under concurrent writes: offsets shift as scores change, which makes offset paging skip and duplicate rows even when it is fast enough. ## The ceilings that actually bind **Memory.** Each member costs the member string plus per-entry overhead in two internal structures. Tens of millions of members is measured in gigabytes, and memory is the master resource: it caps how much else the instance can hold, it determines snapshot fork behaviour and copy-on-write cost, it sets replica full-sync duration, and it raises the cost of every failover. Shortening member identifiers — numeric ids rather than UUID strings — is one of the highest-leverage optimisations available and should be decided before launch. **Single-key hotness.** A Sorted Set is one key, and one key lives in one hash slot on one node. You cannot shard a single leaderboard across a cluster. A globally popular board concentrates its write traffic and memory on one node, and adding cluster nodes does not help that node at all. This is the structural limit, and every scaling strategy below is a way to avoid hitting it. **Write amplification.** Every scoring event is an O(log N) write, plus replication to every replica, plus AOF appends if enabled. A game emitting events per action rather than per match can generate orders of magnitude more writes than the product needs. Debounce in the application — accumulate in memory and flush a player's best or delta on an interval — and the leaderboard's write load becomes a design parameter rather than a consequence. ## Scaling strategies **1. Reads to replicas.** Ranks, top-N and neighbour windows are read-only and tolerate staleness of a second or two — nobody notices their rank is one refresh behind. Routing reads to replicas removes the majority of the load from the primary and is usually the cheapest large win available. **2. Time-partitioned boards.** `lb:daily:*`, `lb:weekly:*` with TTLs keep each key small and let old data expire instead of accumulating. Rollups via `ZUNIONSTORE` run on a schedule, off-peak, ideally on a dedicated instance, because a union over huge sets is one long blocking command. **3. Cohort / regional sharding.** Split the population into K boards by region, tier or a hash of the player id. Top-N per shard is cheap and a global top-N is a K-way merge of K small results client-side. Each shard is a separate key that can live on a separate node, so both memory and write load spread. **4. Approximate global rank across shards.** Merging exact ranks across K boards is not possible from per-shard ranks alone. The standard solution is a **score-bucket histogram**: maintain a second structure (a Hash or a small Sorted Set) counting how many players fall in each score bucket, updated on every score change. A player's approximate global rank is then `sum(counts of all higher buckets) + local rank within their bucket`. With buckets sized to the score distribution this is accurate to within a bucket, costs O(number of buckets) to evaluate, and is a genuinely good answer for the "you are roughly #412 000" display that no user verifies to the unit. Reserve exact ranks for the top few thousand, where users *do* check. **5. Unique scores.** Packing a tiebreaker (an inverted timestamp, or the player id's low bits) into the score makes every score distinct, which makes score-cursor pagination exact and removes the boundary-duplicate problem. The constraint is double precision: exact integers only to 2^53, so budget the bits deliberately. ## What to say out loud The complete answer names three things: the cost model (rank cheap, offset paging expensive, both bounded by one node), the ceilings (memory, single-key hotness, write amplification), and the trade you are proposing (exact ranks for the top slice, approximate ranks below it, cursor paging everywhere, replicas for reads, cohort shards when one node stops being enough). Presenting approximation as a deliberate product decision rather than a technical concession is what distinguishes a principal-level answer.

  • Why is looking up one player's exact rank cheap while paging to the 4000th page is not?
    Rank lookup is a descent through the ordered structure, which stores how many elements each link spans, so counting predecessors costs O(log N) without visiting them. An index-based range must actually traverse from the start index, so its cost includes every element skipped before the window you want. That traversal happens inside a single command, so a deep page stalls the whole server rather than just the requesting client.
  • How can you produce a global rank when the leaderboard has been split into several cohort shards?
    Per-shard ranks cannot be combined directly, so you maintain a shared histogram counting how many players fall into each score bucket, updated whenever a score moves between buckets. A player's approximate global rank is the total count in all higher buckets plus their rank within their own bucket, which is accurate to within one bucket width. Exact ranks are then reserved for the top slice where the number of players is small enough to keep in a single unsharded board.
  • What makes memory the master constraint for a very large leaderboard, beyond simply fitting in RAM?
    Memory size drives the cost of the operations around the data as much as the data itself: snapshotting forks the process and the copy-on-write cost scales with how much of that memory is being written, a replica joining must transfer the whole dataset before it is useful, and failover time grows with both. It also caps what else can share the instance and how much headroom eviction has. That is why shortening member identifiers and partitioning by time period pay back far beyond the raw bytes saved.

A stadium where the scoreboard can instantly tell any one fan their seat number, but sending an usher to row 4 000 means walking past every row on the way — so you hand out tickets that say where the last usher stopped.

saying these in an interview costs you the question

  • Assuming exact rank lookup becomes slow at large sizes when it is logarithmic
  • Paginating deep standings with large index offsets and treating slowness as a client-side problem
  • Believing a single Sorted Set can be sharded across cluster nodes
  • Ignoring that every scoring event is a replicated write and letting the event rate be set by gameplay rather than by design
  • Rejecting approximate ranks on principle when no user can verify a rank of 412,000 to the unit

context