A Redis sorted set on a single node handles tens of millions of members with O(log N) operations. That’s enough for most games. Then you have 500 million players and a score update on every match completion. Suddenly one node isn’t enough.

The Single-Node Limit#

A Redis sorted set for 500 million players needs roughly 5-8 GB of memory just for the skip list nodes and hash table entries. More importantly, every score update and rank query hits one CPU core, because Redis is single-threaded per shard. At peak times (match completion bursts after major tournaments), the sorted set becomes a bottleneck.

Sharding by Score Range#

Partition the leaderboard by score ranges. Players with scores 0-9999 on Shard 1, 10000-19999 on Shard 2, and so on. A score update goes to exactly one shard. A rank query calculates: “sum of all players with scores above mine” by querying all shards for their member count above the player’s score, then summing.

The trade-off: global rank requires a scatter-gather across all shards. For a top-10 leaderboard query, you ask every shard for its top 10, then merge. That’s N round trips, where N is the number of shards.

Player score = 15400, on Shard 2.
Global rank = (players above 15400 on Shard 3) + (players above 15400 on Shard 2).
Query Shard 3 and Shard 2 in parallel, sum results.

Approximate Ranking#

For games where players don’t need an exact global rank, approximate ranking dramatically simplifies the architecture. Use HyperLogLog to estimate player count in score buckets. Show “approximately rank 1.2 million” instead of exact rank 1,183,447. Most players can’t distinguish the two, and the storage and compute cost drops dramatically.

graph TD A[Player Score Update] --> B{Score Range?} B -->|0-9999| C[Shard 1] B -->|10000-19999| D[Shard 2] B -->|20000+| E[Shard 3] F[Global Rank Query] --> C F --> D F --> E C --> G[Aggregate: sum counts above player score] D --> G E --> G style A fill:#000000,stroke:#00ff00,stroke-width:2px,color:#fff style B fill:#000000,stroke:#00ff00,stroke-width:2px,color:#fff style C fill:#000000,stroke:#00ff00,stroke-width:2px,color:#fff style D fill:#000000,stroke:#00ff00,stroke-width:2px,color:#fff style E fill:#000000,stroke:#00ff00,stroke-width:2px,color:#fff style F fill:#000000,stroke:#00ff00,stroke-width:2px,color:#fff style G fill:#000000,stroke:#00ff00,stroke-width:2px,color:#fff

Top-N Is Cheaper Than Global Rank#

Serving a top-100 leaderboard is much simpler than computing a specific player’s global rank. A top-100 query asks each shard for its top-100, then merges. You can cache the top-100 result aggressively because it changes slowly. Cache it for 30 seconds and serve it to millions of users from Redis.

A player’s global rank changes with every score update from every other player. Caching it is hard. Pre-computing it is expensive. This is why many games show “your percentile” (top 5%) rather than exact rank.

At Salesforce#

We had report dashboards that ranked tenants by usage metrics. Global ranking across 150K tenants by query count was expensive to compute on demand. We ran a nightly job that computed percentile buckets and stored them in a summary table. During the day, rank queries hit the summary table, not the full dataset. Stale by up to 24 hours, but close enough for what the dashboard showed.

What I’m Learning#

Top-N lists are easy to scale. Global rank is hard. Approximate rank is the practical middle ground. Understanding which of the three your product actually needs often changes the whole architecture.

How do your leaderboard or ranking systems balance accuracy against performance?