Interview Setup
Interview Prompt
Design a leaderboard system for 100M players with 100K score updates/sec peak, 50K top-K queries/sec, and 200K rank lookups/sec across ~1,000 leaderboards.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| Global leaderboard or per-game/per-season boards? | ~1,000 boards at 10M players each = 1.5 GB memory per board in Redis. |
| Exact ranks or approximate top-K acceptable? |
|
| Time-scoped boards (daily or weekly): reset or rolling window? |
|
| Tie-breaking rule when scores are equal? | Composite score = score + (1 - timestamp/max_time) ensures earliest achiever ranks higher. |
Scope
In scope
- Redis sorted sets
- Rank queries
- Sharded leaderboards
- Time-scoped boards
- Tie-breaking
- Capacity estimation with shown math
Out of scope (state explicitly)
- Search engine and product catalog indexing
- Payment gateway checkout flow
- Fraud and abuse ML pipelines
Functional Requirements
Start by asking your interviewer which leaderboard scopes matter: global top-K, friends-only, and "my rank among 100M players" represent three distinct read patterns. Clarify tie-breaking rules and whether boards reset daily, weekly, or seasonally before you commit to a Redis key design.
- Real-time global leaderboard showing top-K users by score
- User can query their own rank among millions of players
- Update scores in real time (increment, set, decay)
- Multiple leaderboards: daily, weekly, seasonal, all-time, per-game-mode
- Friends leaderboard: rank among friends only
- Historical leaderboards: view past week/season final standings
- Tie-breaking: same score → earlier achievement time ranks higher
- Percentile ranking: "You are in the top 5% of all players"
Non-Functional Requirements
Your interviewer will stress-test write throughput and read latency under peak gameplay: 100K score updates/sec is a much harder requirement than top-K display. They will also probe whether eventual consistency of 1 to 2 seconds is acceptable; state that explicitly so you do not over-engineer strong consistency or cross-shard transactions.
- Low Latency: Top-K query < 10ms; rank lookup < 50ms
- High Write Throughput: 100K score updates/sec during peak
- Consistency: Eventual consistency OK (1 to 2 second delay acceptable)
- Scalability: 100M+ players per leaderboard
- Availability: 99.99%: leaderboard is core game feature
Capacity Estimations
Redis sorted-set memory scales with player count x boards; run the math before assuming one global ZSET fits. Sharding by score range or user_id prefix is the staff-level follow-up when ZREVRANK latency degrades.
| Metric | Calculation | Value |
|---|---|---|
| Total players | Given (assumption documented in value) | 100M |
| Active players per leaderboard | Given (assumption documented in value) | 10M |
| Score updates / sec (peak) | From Score updates / day ÷ 86400 (+ peak factor in value) | 100K |
| Top-K queries / sec | From Top-K queries / day ÷ 86400 (+ peak factor in value) | 50K |
| Rank lookup queries / sec | From Rank lookup queries / day ÷ 86400 (+ peak factor in value) | 200K |
| Leaderboards (total) | Given | ~1,000 |
| Memory per leaderboard | Given | ~1.5 GB |
Architecture Diagram
Walk your interviewer through the write path first: game server posts a score delta, the leaderboard service runs ZINCRBY on a Redis sorted set, and Kafka logs the event for rebuild and analytics. This path must stay under 5 ms at 100K updates/sec.
The read path is separate: ZREVRANGE returns top-K in O(log N + K), ZREVRANK answers "where am I?" in O(log N). At 100M players, exact rank lookup in Redis is feasible; sorting on every read in PostgreSQL is not.
Shard by board_id across a Redis cluster when ~1,000 boards x 1.5 GB each exceeds single-node memory. Hot boards get dedicated shards; historical snapshots land in PostgreSQL after weekly rotation.
In the room
Ask whether exact global rank is required or approximate top-5% is sufficient, because that single answer decides between ZREVRANK and bucketed percentile histograms at 100M scale.
Component Deep Dives
Redis sorted sets are the entire design in one data structure: lead with ZADD and ZINCRBY semantics, then tie-breaking (encode timestamp in the score's fractional bits), then sharding when a single board exceeds memory or hot-key limits.
Redis Sorted Set: Core Data Structure
Redis sorted sets are the canonical leaderboard primitive, and leading with ZADD and ZINCRBY semantics anchors every subsequent design choice to this data structure.
ZADD lb:weekly {score} {user_id} → O(log N) insert/update
ZREVRANK lb:weekly {user_id} → O(log N) get user's rank
ZREVRANGE lb:weekly 0 99 WITHSCORES → O(log N + K) top K
ZINCRBY lb:weekly {delta} {user_id} → O(log N) atomic increment
ZREVRANGE lb:weekly rank-5 rank+5 → O(log N + K) neighbors
Internal: Skip List (balanced probabilistic data structure)
10M members: ~20 operations per query (log₂(10M) ≈ 23)
Each operation: ~1µs → total: ~20µs per query
Memory: ~1.5 GB for 10M entries (key + score + overhead)Tie-Breaking Mechanism
Once scores tie, the interviewer will ask who ranks higher. Encoding the achievement timestamp into the score makes ordering deterministic without requiring a secondary sort pass.
Problem: Two users have score=1000 → who ranks higher?
Rule: Earlier achievement = higher rank
Implementation: Encode timestamp into the score
score_key = score x 10^10 + (MAX_TIMESTAMP - actual_timestamp)
Example (MAX_TIMESTAMP = 9999999999, must be < 10^10 so it never
disturbs the score ordering):
User A: score=1000 at t=100 → key = 1000_0000000000 + 9999999899
User B: score=1000 at t=200 → key = 1000_0000000000 + 9999999799
User A's key > User B's key → User A ranks higher ✓Scaling Beyond 100M Members
When a single board outgrows one Redis node or hot-key limits, evaluate sharding strategies such as score-range shards, global heap merge, or Fenwick trees for percentile-only queries.
Approach 1: Single Redis Sorted Set (recommended up to ~100M) 100M x 150 bytes = 15 GB → fits in single large Redis instance All operations still O(log N) ≈ O(27) Approach 2: Sharded by Score Range (for > 100M) Shard 0: scores 0-999 → Redis instance 0 Shard N: scores 9000-9999 → Redis instance N Top-K: query highest shard → get top from that shard Rank: sum counts in all higher shards + ZREVRANK within shard Approach 3: Fenwick Tree / Binary Indexed Tree (for percentile rank) Score range [0, MAX_SCORE] → array of counts rank(score) = prefix_sum(MAX_SCORE) - prefix_sum(score) O(log S) where S = max score range, regardless of user count
Global board merge at 100M players: heap-merge top-K from 16 shards costs O(K x 16) ≈ 1.6K ZREVRANGE ops every 60s, requiring ~27ms compute compared to 15 GB RAM if held in one ZSET. Exact global ZREVRANK at 100M is O(log 100M) ≈ 27 hops x 200K lookups/sec, which is Redis CPU bound; use Count-Min Sketch for approximate rank on the global board only. Async merge every 60s trades 60s staleness for a sub-second read path, which is acceptable for display boards though not for payout-critical ranks.
Friends Leaderboard
Friends leaderboards are a filtered read rather than a separate data structure. The service fetches friend IDs from the social graph, pipelines ZSCORE calls for each friend, and sorts the results in application memory. The query cost scales as O(F) where F represents friend count.
Percentile Rank at Scale
Percentile rank ("top 5%") is a common staff-level follow-up at massive scale. While exact ZREVRANK works up to roughly 100M members, maintaining 1,000 score count buckets yields O(1) lookup with negligible error, calculated as the sum of buckets above the user's score divided by total users and refreshed every minute.
Event Bus Design (Kafka)
Kafka acts as the durability layer behind Redis. If a Redis shard fails, score events are replayed to rebuild the leaderboard without losing player progress.
Topic: score-events
Partitions: 64 (partition by leaderboard_id)
Partition key: leaderboard_id (per-board ordering for replay)
Retention: 7 days (Redis rebuild buffer after failover)
Producers: Score Writer after ZINCRBY/ZADD succeeds
Payload: {lb_id, user_id, delta, new_score, tiebreak_key, timestamp}
Consumers: analytics pipeline, cross-region sync, PostgreSQL audit log
Topic: leaderboard-snapshots
Low volume; emitted on weekly rotation
Consumers: History Service → PostgreSQL archive of top 10K
Score update path: validate → Redis ZINCRBY (< 1ms) → publish score-events → 200
Rank reads never touch Kafka; Kafka enables rebuild if Redis shard is lostAPI Design
Separate write APIs (game server posts score deltas) from read APIs (client fetches top-K and self-rank). Batch score updates where possible because while each individual ZINCRBY is cheap, 100K/sec creates significant contention on a single hot board.
# Score updates
POST /api/leaderboard/{lb_id}/score
{ "user_id": "u_123", "score_delta": 50 }
# Read operations
GET /api/leaderboard/{lb_id}/top?k=100 → Top K players
GET /api/leaderboard/{lb_id}/rank/{user_id} → User's rank + score
GET /api/leaderboard/{lb_id}/around/{user_id}?n=10 → 10 above + 10 below
GET /api/leaderboard/{lb_id}/friends/{user_id} → Rank among friends
GET /api/leaderboard/{lb_id}/percentile/{user_id} → "Top 5%" info
# Historical
GET /api/leaderboard/{lb_id}/history?date=2026-03-01 → Archived standingsCommon Error Responses
400 Bad Request: invalid input, missing required fields, or malformed JSON payload 401 Unauthorized: missing or invalid authentication token or API key 403 Forbidden: authenticated caller lacks required permissions for this resource 404 Not Found: requested resource ID does not exist 409 Conflict: duplicate write or version conflict, retry with a unique idempotency key 422 Unprocessable Entity: syntactically valid request failed semantic business validation 429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header 500 Internal Error: unexpected server failure, retry safely with an idempotency key 503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff 504 Gateway Timeout: search index shard responded slowly, narrow query parameters or retry
Data Model
-- PostgreSQL (Historical Snapshots)
CREATE TABLE leaderboard_snapshots (
leaderboard_id TEXT,
snapshot_date DATE,
rank INT,
user_id UUID,
score BIGINT,
PRIMARY KEY (leaderboard_id, snapshot_date, rank)
);
-- Redis (Active Leaderboard)
-- ZADD lb:weekly:game1 {score_with_tiebreak} {user_id}
-- HSET user:display:{user_id} name "Alice" avatar_url "..." level 42Kafka (Score Events)
Topic: score-events (partition by leaderboard_id)
{ "lb_id": "weekly:game1", "user_id": "u_123", "delta": 50,
"new_score": 1250, "tiebreak_key": 9999999500, "timestamp": "..." }
Consumers: Redis rebuild on failover, analytics, cross-region syncFault Tolerance
| Technique | Application |
|---|---|
| Redis persistence | RDB snapshots + AOF → recover on restart |
| Redis Cluster |
|
| Rebuild from Kafka | If Redis lost → replay score events → rebuild leaderboard |
| Read replicas | Separate read replicas for top-K (heavy read load) |
| Leaderboard rotation | Weekly reset → archive to PostgreSQL → DEL key |
Leaderboard Rotation: Race Condition
Problem: At midnight, rotate weekly leaderboard
Naive approach (WRONG):
At 00:00:
1. ZRANGEBYSCORE → snapshot top 10K to PostgreSQL
2. DEL lb:weekly
3. New week begins
Race: score update arrives between step 1 and 2 → lost forever
Correct approach: Key rotation with grace period ⭐
Leaderboard key includes time bucket: lb:weekly:2026-W11
At 00:00:
1. Update "current week" pointer: SET current_lb_week "2026-W12"
2. All NEW scores go to lb:weekly:2026-W12
3. Background job: snapshot old to PostgreSQL
4. After snapshot confirmed: DEL old key
✓ No race: "current week" atomically switches with SETCache Stampede
When many users request top-K simultaneously, read from Redis replicas with eventual consistency. If a Redis master fails, a replica promotes within seconds.
Score Tampering
The game server validates scores rather than trusting client submissions, and anti-cheat anomaly detection flags and reviews score increases that occur too rapidly.
Additional Considerations
Score Decay
Inactive players sitting at the top indefinitely make leaderboards feel stale. Effective remedies include time-windowed leaderboards with weekly resets, a daily decay background job using ZINCRBY, or recency-weighted scoring formulas.
Interview Walkthrough
- 25-minute cut
Skip arch50/arch75 depth unless staff.
- Redis ZSET per board (8 min)
- ZADD O(log N) (9 min)
- ZREVRANK for rank lookup (8 min)
- Explain Redis sorted set (ZADD and ZREVRANK) as the canonical leaderboard structure for O(log N) updates.
- Cover sharded leaderboards for global vs friends vs weekly scopes.
- Discuss approximate rank at billion-player scale via HyperLogLog or bucketed tiers.
- Mention async score updates from game events via a queue so gameplay is never blocked on leaderboard writes.
- Cover tie-breaking policy (timestamp and user_id) and state it explicitly.
- Common pitfall: SQL ORDER BY on every read because relational databases cannot serve real-time global rankings at scale.
Engineering Trade-offs
Why Redis Sorted Set Wins
| Solution | Top-K | User Rank | Update | Memory | Complexity | |---|---|---|---|---|---|---| | MySQL ORDER BY score | O(N log N) | O(N) | O(log N) | Disk | High | | Redis Sorted Set ⭐ | O(log N + K) | O(log N) | O(log N) | RAM | Low | | Fenwick Tree | O(S) | O(log S) | O(log S) | RAM | Medium | | Pre-computed ranks | O(1) | O(1) | O(N) recompute | RAM | High | Redis Sorted Set wins because: 1. O(log N) for ALL operations without compromise 2. Built-in ZREVRANK for instant rank lookup 3. ZREVRANGE for top-K in one command 4. ZINCRBY for atomic score updates 5. Fits in memory for realistic user counts (100M = 15 GB)
Multi-Game, Multi-Region at Scale
With 2,500 leaderboards averaging 1.5 GB each, total memory reaches 3.75 TB, which exceeds a single Redis cluster. Resolve this by sharding by leaderboard_id using consistent hashing: hot leaderboards with over 1M active players run on dedicated Redis shards, while smaller boards share pooled shards.
Percentile at 1B+ Players
A score histogram of 1,000 buckets requires only 8 KB of RAM while retaining ±0.1% accuracy. Each score update decrements the old bucket and increments the new bucket. Precomputing prefix sums provides O(1) rank lookups, allowing games like Candy Crush to display percentile rankings without requiring massive 15 GB sorted sets.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.