Problem and requirements
A mobile game awards a point each time a player wins a match. The leaderboard resets every month as a tournament. Players want to see the top 10, their own rank, and ideally the players four places above and below them. Updates should appear in near real time. Scale targets are 5 million daily active users and 25 million monthly active users, and the board must be highly available, because a broken leaderboard is visible to everyone.
The most important non-obvious requirement is trust. Scores must be set by the server, never by the client, because any value a client sends can be forged with a proxy or modified binary. The game service determines that a match was won and then tells the leaderboard service to increment the score.
Back-of-the-envelope estimation
Spreading 5M DAU evenly over a day gives 5,000,000 / 86,400 ≈ 58, or about 50 users per second using the 1e5 shortcut. Traffic is not uniform, so assume a peak of 5x, about 250 users per second. If each user plays around 10 games per day, score updates are about 500 per second on average and roughly 2,500 at peak. Fetching the top 10 is lighter: if each user loads the board once per session start, that is about 50 QPS.
These rates are modest; a single Redis node handles tens of thousands of operations per second. Memory is also modest: storing a 24-character user ID and a small integer score is about 26 bytes per entry, so 25M x 26 B ≈ 650 MB. Even doubling for skip-list and hash-table overhead keeps us near 1.3 GB, which fits one machine. So the design priority is not raw capacity but low-latency ranking and resilience.
APIs and high-level design
Three endpoints cover the needs: POST /v1/scores (internal only, called by the game service with user ID and points), GET /v1/scores for the top 10, and GET /v1/scores/{user_id} returning the user's rank and score. The flow: a player wins, the game service validates the result and calls the leaderboard service, which updates the leaderboard store; clients read from the leaderboard service.
If other systems care about wins, for example analytics or push notifications, the game service can instead publish win events to a message queue like Kafka, with the leaderboard service as one of several consumers. This decouples producers from the number and speed of downstream consumers.
Why relational ranking does not scale
A naive design keeps a leaderboard table of (user_id, score) and computes rank with ORDER BY score DESC. Finding the top 10 is acceptable with an index, but finding one user's rank requires counting how many rows have a higher score, which scans a large fraction of the table each time. With millions of rows and continuous updates, that query becomes seconds rather than milliseconds, and caching does not help because scores change constantly. Batch-computing ranks periodically breaks the real-time requirement.
The underlying problem is that B-tree indexes are not order-statistic trees: they can find a key quickly but cannot tell you its position in sorted order without walking entries.
Deep dive: Redis sorted sets
A sorted set stores unique members, each with a score, ordered by score. Internally Redis combines a hash table mapping member to score, for O(1) lookup, with a skip list ordered by score. A skip list is a linked list with extra express lanes at multiple levels, giving O(log n) search, insert and delete; Redis also tracks span widths, so computing a member's rank is O(log n) as well.
The commands map directly to requirements. ZINCRBY leaderboard_feb_2021 1 user42 adds a point on a win, inserting the user if new. ZREVRANGE leaderboard_feb_2021 0 9 WITHSCORES returns the top 10 in O(log n + 10). ZREVRANK leaderboard_feb_2021 user42 returns the user's zero-based rank, and a second ZREVRANGE over rank minus 4 to rank plus 4 returns the neighbours. ZADD sets a score directly. Each month uses a fresh key such as leaderboard_feb_2021, so the reset is just switching keys, and old boards can be archived or expired.
- ZADD / ZINCRBY: O(log n) update.
- ZREVRANGE: O(log n + m) for m returned entries.
- ZREVRANK: O(log n) rank lookup.
ZINCRBY updates the score, ZREVRANK finds the position, and a small ZREVRANGE window around that rank returns the neighbours.Persistence, caching and serverless variants
Redis keeps data in memory, so a node failure must not lose the tournament. Run a replica for failover, enable persistence, and also record every win in MySQL: a users table for profiles and a point history table of (user_id, score, timestamp). If Redis is lost, the leaderboard can be rebuilt by replaying the history, iterating over rows and calling ZINCRBY for each. The top-10 user profiles (names, avatars) can be cached because that list is read by everyone and changes slowly in composition.
A cloud-native variant uses API Gateway in front of AWS Lambda functions that call ElastiCache Redis and a managed database. Serverless scales automatically and removes server management, which suits spiky game traffic, though cold starts and connection management to Redis need attention.
Scaling to 500M DAU and alternatives
At 100 times the users, memory reaches tens of gigabytes and peak writes approach a quarter-million per second, so we shard. Fixed partitioning by score range puts, say, scores 1 to 100 on shard one and so on; the top 10 lives on the highest shard, and a user's rank is their local rank plus the counts of all higher shards, which are cheap to get. The catch is that the application must track each user's current shard and move users as scores cross boundaries, and ranges must be tuned to balance load.
Hash partitioning via Redis Cluster spreads users evenly by key slot, but ranking breaks: the top 10 requires a scatter-gather that fetches the top 10 from every shard and merges them, and an exact rank for one user requires asking every shard how many players outrank them, which gets slower as shards grow. An alternative is NoSQL, for example DynamoDB with a global secondary index on month and score. Putting all of a month's items under one partition key creates a hot partition, so write sharding appends a suffix to spread items across N partitions, at the cost of querying N partitions and merging. For exact user ranks at that scale, many teams settle for a percentile instead.
ZCARD of every higher shard. The cost is the user-to-shard map, which must move players as their scores cross a boundary.Tie-breaking, failure handling and wrap-up
Ties need a deterministic rule; commonly the player who reached the score first wins. Store the timestamp of the last increment and encode it into the score, or use a secondary sort in Redis by combining score and inverted timestamp into a single numeric value. For failures, the Redis replica handles node loss quickly, and the MySQL history allows a full rebuild. In summary: keep scoring server-authoritative, use a sorted set for O(log n) ranking, use monthly keys for resets, persist history for recovery, and shard by score range when a single node is no longer enough.
Key numbers
Key terms
- Sorted set
- A Redis structure of unique members ordered by a numeric score, backed by a skip list and a hash table.
- Skip list
- A layered linked list that achieves O(log n) search and insert by keeping express pointers at higher levels.
- ZREVRANK
- A Redis command returning a member's zero-based position when ordered from highest to lowest score.
- Scatter-gather
- Sending a query to every shard and merging the partial results.
- Fixed (range) partitioning
- Assigning entries to shards by score ranges so global order can be derived from shard order.
- Write sharding
- Appending a suffix to a hot partition key to spread writes across several partitions.
- Server-authoritative scoring
- Only trusted backend services can change scores, never the client.
Common mistakes
- Letting the client submit its own score, inviting trivial cheating.
- Computing rank with a COUNT over a relational table on every request.
- Hash-sharding the leaderboard and then discovering exact rank requires querying every shard.
- Relying on Redis alone with no durable history to rebuild from.
- Leaving ties undefined so two players swap places unpredictably.
Further study
- Skip Lists: A Probabilistic Alternative to Balanced Trees (William Pugh, 1990)
- Redis sorted sets documentation
- Redis Cluster specification
- Amazon DynamoDB global secondary indexes and write sharding guidance