Technologies referenced in this case study: Redis · Kafka · Flink · Cassandra · DynamoDB · PostgreSQL
Related: Scaling Writes · Contention · Data Pipelines · Stream Processing · Real-time Updates · Sharding
How to Use This Case Study#
Organized for interview use first, reference second. Read front-to-back once. Return to individual sections for targeted review.
| Mode | Time | What to Read |
|---|---|---|
| Quick Review | 15 min | Executive Summary → Interview Walkthrough → Fault Lines table → Active Drills 1–3 |
| Targeted Study | 1–2 hrs | Executive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failure Modes) → two Deep Dives |
| Deep Dive | 3+ hrs | Everything, including the Principal Lens (Section 11) and Appendices |
What is a Leaderboard? — Why interviewers pick this topic
A leaderboard ranks entities (players, creators, sellers, posts) by a score that changes continuously, and answers three questions: who's on top (top-K), where am I (my rank), and who's near me (neighbors). Its sibling problem is counting at scale — likes, views, votes — where the hard part is absorbing millions of increments per second on a handful of hot keys.
Before vs After — the tournament-final incident:
Without a designed leaderboard (SQL ORDER BY on the scores table):
t=0: Esports final ends; 4M viewers open the leaderboard
t=+5s: SELECT ... ORDER BY score DESC LIMIT 100 — 40K QPS, each a 2s index scan
t=+20s: Primary CPU 100%; score writes from 12K live matches queue up
t=+60s: "Your rank" query (COUNT(*) WHERE score > mine) times out for everyone
t=+3min: Prize-eligible top-100 page shows stale data; two players dispute 1st place on stream
t=+2h: Ops restores; legal asks which ranking was authoritative at 20:00:00
With a designed leaderboard:
t=0: Same 4M viewers
t=+0s: Top-100 served from a 1s-TTL cache fed by Redis ZREVRANGE (~0.2ms)
t=+0s: "Your rank" = ZREVRANK, O(log N), ~20µs server time for 10M members
t=+1s: Score writes flow through Kafka → single-writer updaters; no lock contention
t=20:00:00 Snapshot of the authoritative board persisted with event offset; prizes paid from snapshot
Why interviewers reach for this question: Every candidate knows Redis sorted sets, so the question filters on what comes after ZADD: what "rank" means when a board has 500M members and can't fit on one shard, whether exact rank is worth its cost, how windows (daily/weekly) reset without a thundering herd, how to count a viral post's 2M likes per minute, and who is accountable when a prize depends on a number. It is a data-modeling and correctness-budget question dressed as a caching question.
Mechanics Refresher: Ranking and Counting Primitives
| Primitive | How It Works | Pros | Cons |
|---|---|---|---|
| Redis sorted set (ZSET) | Skip list + hash; ZADD, ZINCRBY, ZREVRANK, ZREVRANGE all O(log N) | Exact rank; µs ops; ~100K+ ops/s per shard | One key lives on one shard; ~100 B/member → 100M members ≈ 10+ GB |
SQL ORDER BY + index | B-tree on (score DESC, id) | Durable; flexible queries | Rank = COUNT(*) WHERE score > x → O(N) scan per call |
| Score histogram / buckets | Count members per score range; rank ≈ sum of higher buckets + offset | O(buckets); tiny memory; shards trivially | Approximate within a bucket |
| Sharded counter | Split one hot counter into N sub-counters; read = sum | Write throughput × N | Reads cost N; eventual sum |
| Count-Min Sketch | d hash rows × w counters; estimate = min | Fixed memory; heavy-hitter counts | Overestimates by ≤ ε·total with prob 1−δ |
| Space-Saving / heavy hitters | Track top-m candidates with counts; evict min | Top-K with bounded memory | Approximate for tail; needs m ≫ K |
| HyperLogLog | Unique-count sketch in 12 KB | 0.81% std error; mergeable | Distinct counts only, not ranks |
For most production systems: Redis ZSET for exact top-K and rank of boards up to ~50–100M members per shard; score-bucket histograms for approximate rank (percentile) beyond that; sharded counters or stream pre-aggregation for hot counts; sketches only for trending/top-K over unbounded key spaces. The data structure is almost never the interview question — the correctness budget and the hot keys are.
Executive Summary
If you only read one section, read this. Everything in the case study flows from the contrast below.
What This Interview Actually Tests#
A leaderboard is not a sorted-set question. Everyone knows ZADD.
It is a "what is rank allowed to mean, and who pays when it's wrong" question that tests:
- Whether you separate the top of the board (few rows, many readers, often money attached) from the long tail (millions of rows, each read by one person who wants a rough position)
- Whether you know that exact global rank does not shard — and choose what to give up
- Whether you see write hotness (a viral post, a live match) as a different problem from read hotness
- Whether you define the authoritative snapshot when a prize, payout, or public claim depends on the number
The key insight: Different parts of the same board have different correctness bars. Top-100 must be exact and auditable; rank #4,812,339 can be "top 12%" and nobody is harmed. Staff candidates spend exactness where it has an owner who cares, and approximate everywhere else.
The L5 vs L6 vs L7 Contrast — Start Here#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Redis sorted set, ZADD on every score" | Asks "game ranking, engagement counts, or a prize contest?" and what 'rank' must mean for each | Asks who else in the org builds counters/rankings and whether this is a product feature or a platform primitive |
| Rank semantics | Exact rank for everyone | Exact top-K, approximate percentile for the tail; states the error bound | Makes "exact vs approximate" a declared product contract per surface, reviewed by legal where money is involved |
| Scale | "Shard the sorted set" | Knows hash-sharding breaks global rank; uses score-range tiers or bucket histograms; top-K via merge | Prices exactness: "exact rank for 500M users costs ~3× memory and a cross-shard read per request; we pay it only for prize boards" |
| Hot writes | "Redis is fast enough" | Sharded counters / pre-aggregation in the stream; single-writer per key; quantifies per-key ceilings | Standard counter service with tiers (exact, sharded, sketched) that 20 teams adopt instead of building their own |
| Windows | "New key per day" | Key-per-window with TTL; staggered expiry; tumbling vs sliding choice; late events policy | Aligns window boundaries with business calendars and time zones across products; one definition of "this week" |
| Integrity | Trusts score events | Server-authoritative scores, anti-cheat hooks, audit log, snapshot for payouts | Owns the dispute process: who adjudicates, from which snapshot, within what SLA |
Why "rank semantics" separates levels
L5: Promises exact rank for every user because ZREVRANK returns it. It works — until the board has 800M members, which don't fit on one shard, and hash-sharding makes every "my rank" call a scatter-gather ZCOUNT across 64 shards.
L6: "Exact rank matters for the top few thousand — people screenshot it, prizes attach to it. For everyone else, 'you're in the top 12%' is more useful and costs O(1). I'll keep an exact ZSET for the top 10K and a score histogram with 10K buckets for everyone; rank error for the tail is at most the population of one bucket, which I'll size to keep error < 0.1%."
L7: "Which surfaces show exact rank is a product contract. Game boards: exact top-10K. Creator rankings used for payouts: exact, from a nightly batch snapshot, not real-time. Engagement 'trending': approximate by design. I'd publish this so no team promises exact rank where the platform only provides approximate."
Why "hot writes" separates levels
L5: Treats a single INCR or ZINCRBY per event as fine because Redis does ~100K ops/s. Then a celebrity post receives 2M likes in a minute (~33K/s) on one key, on one shard, on one thread — and every other key on that shard stalls.
L6: "A single key has a single-thread ceiling regardless of cluster size. For hot counters I'll pre-aggregate in the stream — Flink tumbles 1s windows per key and emits one INCRBY 33000 instead of 33K INCRs — and split known-hot keys into 16 sub-counters. Reads sum 16 keys, ~0.3ms."
L7: Turns the pattern into a counter service with declared tiers so product teams stop discovering the hot-key ceiling in production, one viral event at a time.
Why "integrity" separates levels
L5: Scores arrive from the game client and are written to the board.
L6: Scores come from the authoritative game server, not the client; events carry match ID and sequence so replays are idempotent; outlier detection quarantines suspicious jumps before they reach the public board; the payout board is a snapshot at a defined offset.
L7: Recognizes that any board with money attached is a financial record: it needs an audit trail, a dispute window, and a named adjudicator — and it's the legal and trust-and-safety teams, not engineering, who sign off on the rules.
The Staff Positions#
| Position | Rationale |
|---|---|
| Exact top-K, approximate tail | Exactness has owners at the top of the board and none at rank 4M |
| Server-authoritative scores via an event log | Replayable, auditable, idempotent; clients never write scores |
| Single writer per board partition | Removes read-modify-write races; ordering by offset makes updates deterministic |
| Pre-aggregate hot counters in the stream | One key has one thread; batch 1s of increments into one write |
| Key-per-window with TTL, not reset jobs | A "reset at midnight" job is a herd and a race; new keys make windows free |
| Payout boards are snapshots, not live reads | Money needs a reproducible number tied to an event offset |
| Cache the top page for 1–5s | 4M viewers reading top-100 should cost one Redis call per second, not 40K |
The Three Intents#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Real-time competitive board (games, contests) | Low-latency rank for millions; exact at the top | ZSET per board/window for top tier; histogram for tail; event-sourced scores | Rank disputes; board lag during finals | Exact top-K within 1–2s; tail within ±0.1% |
| Engagement counting & trending (likes, views, top posts) | 100K–10M increments/s; extreme key skew | Stream pre-aggregation, sharded counters, heavy-hitter sketches | Hot-key shard meltdown; counts that go backwards | Monotonic, eventually accurate (minutes); trending approximate |
| Money-bearing ranking (payouts, creator funds, sales contests) | Auditable, reproducible, disputable | Batch recomputation from the event log; signed snapshots | Paying the wrong person; unreproducible result | Exact, reproducible from a defined cutoff; human dispute process |
🎯 Staff Move: "These three share a sorted order and nothing else. I'll design the real-time competitive board — exact top-K, approximate rank for everyone else, daily and weekly windows — because it forces the sharding and freshness tradeoffs. I'll handle hot counting as a write-path concern inside it, and I'll say explicitly that any payout comes from a batch snapshot, not the live board."
The Five Fault Lines#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Exact vs Approximate Rank | Exact rank for all (expensive, doesn't shard) vs exact top + approximate tail? |
| 2 | Freshness vs Write Amplification | Apply every event immediately vs batch/pre-aggregate and accept seconds of lag? |
| 3 | One Sorted Set vs Sharded Board | Single key (exact, capped by one shard) vs hash / score-range sharding (scales, rank gets hard)? |
| 4 | Window Semantics | Tumbling vs sliding, UTC vs local, late events — who defines "this week"? |
| 5 | Trust vs Throughput in Score Ingestion | Accept scores fast from producers vs validate, dedupe, and anti-cheat before they're public? |
In the Wild: Real Production Systems#
Why this section belongs here: Citing specific production systems demonstrates you've studied operational reality, not textbook designs.
Redis Sorted Sets — The Default Game Leaderboard#
Redis's own documentation and cloud providers' reference architectures (for example, AWS ElastiCache materials) present sorted sets as the canonical leaderboard: ZADD/ZINCRBY to update, ZREVRANGE for top-K, ZREVRANK for "my rank," all O(log N). A single sorted set with tens of millions of members answers rank in microseconds.
Staff insight: The ZSET is the right answer per shard. The interview is about what you do when one board outgrows one shard or one thread — cite the default, then say where it stops.
YouTube's "301 Views" — Counting With a Verification Gate#
For years YouTube public view counts famously appeared to stall around 301 while the system verified that views were legitimate; YouTube has since changed how counts update. The public explanation was that early views were validated before being reflected in the public number.
Staff insight: A public count is a product surface with an integrity bar. Showing a lagging-but-trustworthy number was a deliberate choice of integrity over freshness — the same choice a prize board must make.
Google Cloud Firestore — Distributed Counters#
Firestore's documentation describes a sustained write limit of roughly one write per second per document and recommends distributed counters: split a counter into N shard documents, increment a random shard, and sum on read.
Staff insight: This is the sharded-counter pattern documented by a major vendor as the required workaround for a per-key write ceiling. Every store has such a ceiling — Redis (one thread per shard), DynamoDB (~1,000 WCU/s per partition) — and a Staff answer names it and the split factor.
What Interviewers Probe#
| After You Say... | They Will Ask... | (What They're Evaluating) |
|---|---|---|
| "Redis sorted set" | "500M players. Does it fit? What's rank when it's sharded?" | Whether you know rank doesn't compose across hash shards |
| "ZINCRBY per event" | "A streamer's match gets 50K score events/s. One key." | Single-key/single-thread ceiling |
| "Reset the board at midnight" | "Which midnight? What happens to a score event at 23:59:59.9?" | Window semantics, late events |
| "Exact rank for everyone" | "What does exactness cost at rank 40 million, and who needs it?" | Correctness budget |
| "Clients send their score" | "How do you stop a cheater from posting 999,999,999?" | Trust boundary |
| "The live board decides the winner" | "Two players dispute 1st. Which number is authoritative?" | Snapshots and audit |
System Architecture Overview#
Reading the diagram: Producers are authoritative servers, never clients. Every score is an event in Kafka, partitioned so each board-entity pair has one writer. Flink validates, pre-aggregates 1s of events per key, and assigns events to windows. The serving layer splits by correctness bar: an exact ZSET for the top tier, a histogram for approximate rank of everyone, sharded counters for hot counts, and a 1s cache in front of the top page. Snapshots tie board state to a Kafka offset so any published number can be reproduced — and payout boards come from batch recomputation, not the live ZSET.
Quick-Reference: The 30-Second Cheat Sheet#
| Topic | The L5 Answer | The L6 Answer — Say This |
|---|---|---|
| Data structure | "Redis sorted set" | "ZSET for the exact top tier; score histogram for everyone's approximate rank." |
| Sharding | "Shard the ZSET by user" | "Hash-sharding breaks rank. Top-K = merge of per-shard top-K; rank = histogram, or score-range tiers." |
| Hot writes | "Redis is fast" | "One key = one thread ≈ 100K ops/s ceiling. Pre-aggregate 1s in Flink; shard known-hot counters 16 ways." |
| Windows | "Cron resets at midnight" | "Key per window (lb:{board}:2026-W40) with TTL. Nothing resets; old keys expire." |
| Integrity | "Validate input" | "Server-authoritative events, event_id dedupe, outlier quarantine, snapshot for payouts." |
| Read load | "Redis handles reads" | "Top page cached 1s — 4M readers become 1 ZREVRANGE/s. Rank is per-user, uncached, O(log N)." |
Key Numbers Worth Memorizing#
| Metric | Value | Why It Matters |
|---|---|---|
| ZSET op latency (server) | ~1–20 µs for O(log N) at 10M members | Rank is cheap on one shard |
| Redis single-shard throughput | ~100–200K simple ops/s (one thread) | Hot key ceiling — adding shards doesn't help one key |
| ZSET memory | ~80–120 bytes/member (short member IDs) | 100M members ≈ 8–12 GB; 1B doesn't fit one shard |
| Per-partition write ceilings | DynamoDB ~1,000 WCU/s; Firestore ~1 write/s per doc | Why sharded counters exist |
| Count-Min Sketch error | ε = e/w, δ = e^−d; w=2,000, d=5 → ~0.14% of total, 99.3% confidence | Sizing sketches for trending |
| HyperLogLog | 12 KB, ~0.81% std error | Unique viewers per post, mergeable |
| Histogram buckets | 10K buckets × 8 B = 80 KB per board | Approximate rank for any board size |
| Pre-aggregation reduction | 1s window turns 33K INCR/s into 1 INCRBY/s per key | The hot-counter fix |
| Top-page cache TTL | 1–5s | Converts reader fan-in to O(1) backend reads |
| Double precision | 53-bit mantissa | Limits composite score+tiebreak encoding |
Interview Walkthrough
The most common mistake: Candidates draw
ZADD+ZREVRANGEin minute 4 and then spend 20 minutes on API shapes and caching headers. The interviewer's real question — "it has 500 million members now; what is rank?" — arrives at minute 35 with no time left. Get the basics done by minute 10.
Phase 1: Requirements & Framing (2–3 minutes)#
Functional requirements in 30 seconds:
"Update a player's score, show the top N, show my rank and the players around me, for all-time plus daily and weekly windows."
Then intent and non-functionals — the Staff move:
"Is this a live game board, an engagement counter like likes and trending, or a payout ranking where money depends on the order? They have different correctness bars. I'll design the live competitive board, and I'll treat payouts as a batch snapshot problem layered on top."
Commit to numbers:
"Assume 200M registered players, 20M daily active, 200K score updates/s at peak during events, 50K rank reads/s and a top-page spike to 4M concurrent viewers during finals. Top-1,000 exact within 2 seconds of the event. Everyone else gets rank as a percentile within ±0.1%. Scores come only from game servers."
🎯 Staff Move: "Exact top, approximate tail" is the load-bearing sentence. It decides the sharding strategy, the memory budget, and the failure behavior before you draw anything.
Phase 2: Core Entities & API (1–2 minutes)#
- ScoreEvent:
event_id,board_id,entity_id,deltaornew_score,source_match_id,event_time - Board:
board_id,aggregation(max,sum,latest),windows(all,daily,weekly),tiebreak(earliest),exact_top_k - RankView:
rank(exact ornull),percentile,score,neighbors[]
POST /boards/{id}/events (internal, game servers only) { event_id, entity_id, delta, event_time }
GET /boards/{id}/top?window=weekly&n=100 → [{ rank, entity_id, score }]
GET /boards/{id}/me?window=weekly → { rank | null, percentile, score }
GET /boards/{id}/around/{entity}?window=weekly&k=5 → neighbors (exact only in top tier)
🎯 Staff Move: "The response shape encodes the correctness contract:
rankis present only when it's exact; otherwise the client showspercentile. The UI can't accidentally promise precision the backend doesn't have."
Phase 3: High-Level Architecture (≤5 minutes)#
Walk the flow in 90 seconds:
- Game servers publish
ScoreEvents to Kafka, keyed by(board_id, entity_id)— per-entity ordering, one writer per entity - Flink dedupes by
event_id, assigns events to windows (all-time,2026-09-30,2026-W40), and folds 1s of events per entity into one update - For each window: update the entity's score in the canonical score store, update the histogram (decrement old bucket, increment new), and
ZADDinto the top-tier ZSET if the score is above the tier threshold - Top page:
ZREVRANGE 0 99behind a 1s cache - "My rank": if my score ≥ tier threshold,
ZREVRANK(exact); else histogram lookup (approximate percentile)
Key points to state:
- Scores are events — replayable, auditable, idempotent
- Single writer per entity — no read-modify-write races on
maxaggregation - Two serving structures by correctness bar — exact ZSET for the top, histogram for the tail
- Windows are keys, not resets
- Top page is cached; rank is not — fan-in on the top, fan-out on rank
🎯 Staff Move: Then say: "This works up to maybe 50M members per board on a single ZSET. The interesting part is what breaks at 500M, what happens when one entity gets 50K events a second, and what happens at the window boundary."
Phase 4: Transition to Depth (1 minute)#
"Three places worth going deep: what rank means when the board doesn't fit on one shard, hot writes on a single entity or counter, and window boundaries plus late events. Which is most interesting to you?"
Default if no preference: sharding and rank semantics — it's the question every leaderboard interview eventually reaches.
Phase 5: Deep Dives (25–30 minutes)#
Deep dive A: Rank at 500M members (8–10 min)
"At ~100 bytes a member, 500M is ~50 GB — past what I want in one Redis shard, and one shard's thread would serve every rank read. Options: hash-shard the ZSET 64 ways — top-K becomes a merge of 64 per-shard top-Ks, fine, but rank becomes Σ ZCOUNT(score, +inf) across 64 shards on every call, 64× the read load. Score-range sharding — shard 0 holds scores ≥ 10,000, shard 1 holds 5,000–9,999 — keeps rank at 'offset of higher shards + local rank', but score distributions shift over a season, so ranges need rebalancing, and the top shard is the hot one. Or the one I'd pick: keep only the top tier exact — the top 100K in one ZSET, ~10 MB — and serve everyone else from a 10K-bucket histogram. Rank error below the tier is bounded by bucket population; with adaptive bucket edges I size buckets so none holds more than 0.1% of members."
Deep dive B: Hot writes (6–8 min)
"Two kinds of hot. A hot entity — a streamer's team gets 50K score events/s during a raid — is handled by Flink pre-aggregation: 1s windows per key collapse 50K events into one ZINCRBY. A hot counter with no natural owner — likes on a viral post at 2M/min — gets both pre-aggregation and a 16-way split: likes:{post}:{0..15}, writes pick hash(event_id) % 16, reads MGET 16 keys. Single-key ceiling is about 100K ops/s on one Redis thread, and pre-aggregation drops us to ~1 write/s/key anyway, so the split matters mainly for stores with low per-key limits like DynamoDB."
Deep dive C: Windows and late events (5–7 min)
"Weekly board = key lb:{board}:2026-W40. At the boundary, nothing resets; new events go to W41, and W40 gets a TTL of 7 days after close. The hard case: a match that ended at 23:59:58 whose event arrives at 00:00:03. I assign windows by event_time from the game server, not arrival time, and hold the window open for a 5-minute allowed lateness. At close + 5 min, I snapshot W40 with its Kafka offset — that snapshot is the official result. Events later than that go to a dead-letter topic for manual review, because silently changing a published final ranking is worse than excluding one late event."
🎯 Staff Move: End each dive with who owns it: "Game-design owns the window definition and lateness; trust-and-safety owns quarantine rules; the platform owns update lag."
Phase 6: Wrap-Up (2–3 minutes)#
"Summary: scores as events, one writer per entity, exact top tier in a ZSET, histogram for approximate rank, windows as keys with snapshots at close, pre-aggregation for hot writes, and a 1s cache on the top page. Payouts come from batch recomputation off the event log, never the live board. Next I'd build friends leaderboards — that's a different problem, a small ZSET per user or on-read computation over ~200 friends — and I'd deliberately not build exact global rank for the tail unless someone can name who needs it."
Common Timing Mistakes#
| Mistake | Time Lost | Fix |
|---|---|---|
| Explaining skip-list internals | 5 min | "ZSET, O(log N). Moving on." |
| Designing the REST API in detail | 5 min | Three reads, one internal write, done |
| Debating Redis vs Memcached | 3 min | Memcached has no sorted structure — one sentence |
| Designing friends leaderboards first | 8 min | Park it; it's a different fan-out problem |
| Never reaching 500M members | Interview-ending | Transition at minute 12 |
1. The Staff Lens#
1.1 Why This Problem Exists in Staff Interviews#
The leaderboard looks like the easiest question in the catalog because the textbook answer is one data structure. That's exactly why it's used: it reveals whether a candidate stops at the data structure or asks what the numbers are for.
Real leaderboards are fought over. Players screenshot their rank. Streamers argue on air about who was first. Creator programs pay out based on view rankings. Sales contests pay bonuses. The moment a number is public or paid, it acquires stakeholders — and each stakeholder has a different correctness bar. The Staff skill being tested is allocating exactness: spending memory, latency, and engineering effort on precision where someone is harmed by imprecision, and nowhere else.
The counting half of the problem tests the other classic Staff skill — recognizing that a single key is a single-threaded bottleneck no matter how large the cluster, and designing the write path around key skew rather than average load.
1.2 The L5 vs L6 Contrast — Visual#
The L5 path is right for a 10M-member board. It stalls at step 5 because "exact rank for everyone" was assumed rather than chosen, so sharding breaks a promise that was never priced.
1.3 The Staff Question That Cuts Through Everything#
"Who is harmed if this rank is off by 1% — and do they get paid based on it?"
Asking it out loud forces:
- Correctness allocation: exact for the top, approximate for the tail
- Freshness: real-time for display, snapshot for anything consequential
- Ownership: who adjudicates disputes, and from which recorded state
2. Problem Framing & Intent#
2.1 The Three Intents — Explained#
Intent 1: Real-time competitive board. Games, fitness challenges, coding contests. Millions of members, frequent updates, heavy reads of the top page and of "my rank." The correctness bar is exact at the top, fast everywhere. Aggregation is usually max (best score) or sum (points). Windows matter: daily and weekly boards drive re-engagement. The hard problems: rank at scale, window boundaries, and cheating.
Intent 2: Engagement counting & trending. Likes, views, shares, "top posts this hour." Write volume dwarfs read volume at the key level, key popularity is power-law distributed, and the "board" (trending) is computed over an unbounded key space. The correctness bar is monotonic and eventually accurate for counts, approximate for trending. The hard problems: hot keys, pre-aggregation, heavy-hitter detection with bounded memory, and counts that must never visibly go down.
Intent 3: Money-bearing ranking. Creator payouts, sales contests, tournament prize pools, affiliate rankings. Lower volume, extreme scrutiny. The correctness bar is reproducible and auditable: given the event log up to cutoff T, anyone can recompute the ranking and get the same answer. The hard problems: cutoff semantics, fraud filtering that's applied consistently, and a dispute process.
🎯 Staff Move: "The live board is for engagement; the payout board is for accounting. I'll never pay anyone from a Redis ZSET — the payout job recomputes from the event log in batch with fraud filters applied, and publishes a signed snapshot. The live board can be wrong for 2 seconds; the payout board can't be wrong at all."
2.2 When NOT to Build a Real-Time Leaderboard#
| Situation | Use Instead | Why |
|---|---|---|
| < 1M members, low write rate | Postgres with ORDER BY score DESC LIMIT 100 + a materialized rank refreshed every minute | One store; rank staleness of 60s is invisible |
| Rankings refreshed daily | Batch job (Spark/SQL) writing a ranked table | No streaming infra to operate |
| Friends-only boards (~200 friends) | Compute on read: fetch friends' scores (multi-get), sort in memory | 200 lookups < 5ms; no per-user ZSET to maintain |
| Payouts | Batch recomputation from the event log | Reproducibility beats freshness |
| Trending over unbounded keys | Stream heavy-hitter sketch (Space-Saving / Count-Min + heap) | A ZSET of every post ever is unbounded |
| Analytics "top products this quarter" | OLAP (ClickHouse, BigQuery, Druid) | Ad-hoc slicing; minutes of latency acceptable |
"If nobody refreshes the page expecting it to change within seconds, it isn't a real-time leaderboard, and I shouldn't pay for one."
2.3 What the Interviewer Leaves Underspecified#
| Unstated Assumption | Why It Matters | What to Say |
|---|---|---|
Aggregation (max, sum, latest) | max is idempotent-ish; sum requires exact dedupe | "Best score per player, max — safer under replays." |
| Tie-breaking | Two players at 9,800 — who is #1? | "Earlier achiever wins; encoded in the score." |
| Board size | 1M vs 1B changes everything | "200M members, 20M active per week." |
| Freshness | 100ms vs 60s | "Top tier within 2s; tail within 30s." |
| Window definition | UTC vs local midnight; ISO week | "UTC, ISO weeks; game design signs off." |
| Who writes scores | Client-submitted scores invite cheating | "Game servers only." |
| Does money depend on it | Changes the correctness bar to auditable | "Payouts use batch snapshots, not the live board." |
| Segmentation | Global vs region vs friends vs league | "Global + per-region boards; friends is a separate path." |
2.4 Precise Terminology#
| Term | Precise Meaning |
|---|---|
| Dense rank | Ties share a rank; next rank increments by 1 (1, 2, 2, 3) |
| Competition rank | Ties share a rank; next skips (1, 2, 2, 4) — what most games show |
| Exact tier | Members whose score ≥ tier threshold; served from an exact ZSET |
| Tier threshold | Score of the Kth member (e.g., 100,000th); recomputed continuously |
| Approximate rank | Rank estimated from a histogram: members in higher buckets + interpolated offset within bucket |
| Rank error bound | Max difference between approximate and true rank; ≤ population of one bucket |
| Window | Time interval a board aggregates over: tumbling (fixed days/weeks) or sliding (last 24h) |
| Allowed lateness | How long after window close late events are still applied |
| Snapshot | Board state + the Kafka offset it reflects; reproducible |
| Hot key | A key whose write/read rate approaches the per-key or per-shard ceiling |
| Pre-aggregation | Folding many events per key into one update in the stream before hitting the store |
3. The Fault Lines#
3.1 Fault Line 1: Exact vs Approximate Rank#
The tension: Exact rank for every member is what users seem to want and what the naive structure provides — until the board outgrows one shard, at which point exact rank for the tail costs a scatter-gather per read.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Exact ZSET for everyone, one shard | Simple; exact; µs rank | Memory cap (~50M–100M members per comfortable shard); one thread for all reads | Platform when the board outgrows the shard mid-season |
| Exact, hash-sharded, scatter-gather rank | Scales memory | Every rank read hits all N shards — 64× read amplification; p99 = slowest shard | Latency budget; Redis fleet cost |
| Exact top tier + histogram tail | Top exact; tail O(1); memory tiny | Tail rank approximate; tier threshold must be maintained | Tail users see "Top 12%" instead of "#4,812,339" — product must sign off |
| Percentile only (no ranks) | Cheapest; honest | No "#1" bragging at the top | Top players — unacceptable for competitive games |
Staff default: Exact top tier (top 10K–100K) in one ZSET + histogram for everyone. Show #rank inside the tier and Top X% outside it.
When to deviate: Boards under ~50M members — keep one exact ZSET; it's simpler and fits. Boards where every member's exact position matters contractually (a 500-person sales contest) — exact, but that's small.
"Exact rank at position 4 million is a number nobody screenshots and nobody gets paid for. I'll spend exactness where there's a stakeholder."
3.2 Fault Line 2: Freshness vs Write Amplification#
The tension: Applying every event immediately gives the freshest board and the highest write load. Batching reduces load and introduces lag.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Write-through per event | Sub-100ms freshness | 200K writes/s; hot entities saturate one thread | Platform on-call during finals |
| Pre-aggregate 1s per key in stream | Writes drop to ≤ 1/s per active key; 10–1000× reduction for hot keys | +1s lag; stream job becomes critical path | Players see scores ~1s later — nobody notices |
| Micro-batch 30–60s | Minimal load | Visible lag; "I just scored, why didn't my rank change?" | Players; support tickets |
| Nightly batch | Cheapest; reproducible | Not real-time | Engagement — only acceptable for payouts |
Staff default: 1s pre-aggregation in Flink for all boards; 5s for engagement counters; batch for payouts.
"One second of lag buys a 1,000× reduction on hot keys. That's the best trade in this design."
3.3 Fault Line 3: One Sorted Set vs Sharded Board#
The tension: How do you split a board that doesn't fit, when rank is a global property?
| Strategy | Top-K | My Rank | What Breaks | Who Pays |
|---|---|---|---|---|
| Hash shard by entity (N shards) | Merge N per-shard top-Ks — O(N·K) | Σ ZCOUNT(my_score, +inf) across N — N calls | Rank read amplification ×N | Latency, fleet cost |
| Score-range shards | Read top shard only | Offset of higher shards (cached counts) + local ZREVRANK | Distribution drift → rebalance; top shard hottest; members move shards on score change | Platform — rebalancing is continuous |
| Top tier + histogram | Top tier ZSET | Tier: exact; tail: histogram | Approximate tail | Product accepts percentile |
| Per-segment boards (region, league of 100) | Per segment | Per segment exact | Global rank not provided | Product — "global rank" becomes "league rank" |
Staff default: Top tier + histogram for global; per-league boards (many small exact ZSETs) where game design allows — leagues of 50–100 players are how many games sidestep global rank entirely, and they're better for engagement.
3.4 Fault Line 4: Window Semantics#
The tension: "Weekly leaderboard" hides four decisions: tumbling or sliding, which time zone, event time or arrival time, and what happens to late events.
| Decision | Options | Staff Default | Who Signs Off |
|---|---|---|---|
| Shape | Tumbling (Mon–Sun) vs sliding (last 7×24h) | Tumbling — sliding needs per-event expiry and has no "final result" | Game design |
| Clock | UTC vs player-local | UTC — one global boundary; local-time boards fragment the ranking | Game design + product |
| Assignment | Event time vs arrival time | Event time from authoritative server | Platform |
| Lateness | Drop / apply / manual | Apply within 5 min allowed lateness; after close+5m → DLQ, manual review | Trust & safety + game design |
| Reset mechanism | Delete/zero keys vs new key per window | New key per window with TTL | Platform |
"A sliding window has no final result, so nobody can win it. If a board ever awards anything, it's tumbling."
3.5 Fault Line 5: Trust vs Throughput in Score Ingestion#
The tension: Every validation step adds latency and complexity; every skipped step is a door for cheaters and double-counting.
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Client-submitted scores | Simple; zero server cost | Trivial cheating; 999,999,999 at #1 within an hour of launch | Honest players; brand |
| Server-authoritative, no validation | Blocks naive cheating | Server bugs and replays double-count; exploits in game logic still land | Players who lose rank to a bug |
| Server-authoritative + dedupe + outlier quarantine | Idempotent; suspicious jumps held before public | +100–500ms for quarantined events; false positives delay legit scores | Trust & safety owns rule tuning |
| Full async review before any publish | Near-zero public cheating | Minutes-to-hours lag; kills real-time engagement | Engagement |
Staff default: Server-authoritative events with event_id dedupe, per-event plausibility checks (score delta ≤ max possible per match), statistical outlier quarantine (e.g., > 6σ above the player's history or above the current #1), and a retroactive removal path that updates the board and the snapshots.
"The public top-10 is an attack surface. I'll hold a score out of the top tier for review before I'll let a cheater sit at #1 on a stream for twenty minutes."
4. Failure Modes & Operational Reality#
4.1 The Hot Key That Stalled a Shard#
t=0: Celebrity posts; likes counter likes:{post_9} lives on Redis shard 14
t=+30s: 38K INCR/s on one key; shard 14's single thread at 95% CPU
t=+45s: Other 1.2M keys on shard 14 see p99 latency rise from 0.3ms to 40ms
t=+60s: Unrelated leaderboard reads on shard 14 time out at 50ms; API returns 5xx for 3% of rank calls
t=+2min: Client retries double load on shard 14
t=+4min: On-call identifies hot key via redis-cli --hotkeys / key-level sampling
t=+6min: Enable pre-aggregation path for the post; INCR rate → 1/s; shard recovers
Detection: redis.shard_cpu{shard} > 80%, hot_key.ops_per_sec top-K sampling, counter.write_rate{key} > 10K/s, p99 skew between shards.
Blast radius: Every key co-located on the hot shard — which is why the hot key harms other teams.
Mitigation: Route all counter writes through the stream pre-aggregator by default; auto-split keys that cross 5K writes/s into 16 sub-keys.
Prevention: Never allow direct INCR from request paths for public counters; client-side retry budgets.
Owner: Counter platform team; the posting product's team for retry behavior.
4.2 The Window-Boundary Herd#
A daily board naively resets at 00:00 UTC by a cron job that runs DEL on 40K per-region/per-mode keys, while the client app refreshes every player's rank at midnight to show "new day."
t=00:00:00 Reset job issues DEL on 40K ZSETs (some with 5M members)
t=+0.2s DEL of a 5M-member ZSET blocks the Redis thread ~2–3s (synchronous free)
t=+1s 20M clients refresh "my rank" simultaneously; half hit blocked shards
t=+3s Events for the new day arrive; some land before DEL, some after → new-day scores lost
Fix: New key per window (lb:{board}:{date}), no deletes; old keys get EXPIRE (Redis frees large values lazily with lazyfree-lazy-expire yes / UNLINK); clients refresh with jitter (0–60s); first read of the new window shows "Day started — play to rank."
Detection: redis.blocked_ms, board.events_dropped_total{reason="window_race"}, rank-read p99 at boundary.
Owner: Platform (key scheme), client team (refresh jitter).
4.3 Double-Counted Scores After a Consumer Rebalance#
sum boards (total points this week) are only correct if each event is applied once. A Flink job restart replays 90 seconds of events from the last checkpoint; if the sink applies ZINCRBY non-transactionally, those 90 seconds are counted twice.
Detection: Reconciliation job: recompute a sample of 10K players' weekly sums from the event log every 15 min; alert if > 0.1% mismatch. Metric board.reconcile_mismatch_ratio.
Mitigation: Idempotent sinks — write absolute values (Flink holds the per-entity sum in state; sink does ZADD with the absolute score, which is idempotent) rather than increments. For max boards, ZADD … GT is naturally idempotent.
Owner: Platform streaming team.
4.4 Cheater at #1 During a Live Final#
An exploit lets a player submit a score 40× the plausible max. It passes the per-match check (the check used a stale max-score constant from a previous season). The player sits at #1 on stream for 22 minutes.
Detection: board.top1_score_jump_ratio (new #1 score / previous #1 > 1.5× pages trust & safety), outlier quarantine hit rate.
Mitigation: Quarantine new entries into the top 100 whose score exceeds the previous #1 by more than 20% until reviewed (auto-release after 10 min if no action). Remove retroactively from ZSET, histogram, and snapshots.
Owner: Trust & safety (rules), game team (max-score constants per season), platform (quarantine mechanism).
4.5 Silent Staleness — The Board Stopped Moving#
The Flink job's Kafka consumer lags after a partition reassignment. The board is still served, still looks plausible, and is 40 minutes stale. Nobody notices until a player tweets that their rank hasn't changed all evening.
Detection: board.update_lag_seconds = now − event_time of latest applied event; alert at > 30s p99 for 5 min. Also synthetic canary: a test entity scores every 10s on a hidden board; alert if its rank isn't updated within 5s.
Owner: Platform.
4.6 Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Hot counter key | hot_key.ops_per_sec > 10K; shard CPU > 80% | Every key on that shard | Pre-aggregate; auto-split into 16 | Counter platform |
| Board outgrows shard | zset.members > 50M; shard memory > 70% | One board | Top tier + histogram; per-league boards | Platform |
| Window boundary herd | redis.blocked_ms; boundary rank-read p99 | All boards at boundary | Key-per-window; UNLINK/lazy expire; client jitter | Platform + client |
| Double counting on replay | board.reconcile_mismatch_ratio > 0.1% | sum boards | Absolute-value idempotent sinks | Streaming team |
| Cheater at top | top1_score_jump_ratio > 1.5 | Public trust | Quarantine + retroactive removal | Trust & safety |
| Silent staleness | board.update_lag_seconds > 30s; canary | All boards on the job | Lag alert; canary entity | Platform |
| Top-page cache stampede | Cache miss storm at TTL expiry | Top-tier shard | Single-flight refresh; stale-while-revalidate | Platform |
| Histogram drift | rank.approx_error_estimate > 0.1% | Tail rank accuracy | Adaptive bucket re-splitting | Platform |
| Payout dispute | Support escalation | One contest, reputational | Reproducible snapshot + dispute SLA | Contest owner + legal |
5. Evaluation Rubric#
5.1 Level-Based Signals#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Rank semantics | Exact for everyone via ZREVRANK | Exact top tier, approximate tail with a stated error bound | Exactness declared per product surface; payout boards governed as financial records |
| Sharding | "Shard the ZSET" | Explains why hash-sharding breaks rank; picks tier + histogram or score ranges | Chooses segmentation (leagues) with game design as a product lever that removes the scaling problem |
| Hot writes | "Redis is fast" | Single-key ceiling; pre-aggregation; split counters with factor justified | Org counter service with exact/sharded/sketched tiers and published per-key limits |
| Windows | Cron reset | Key-per-window, event-time assignment, allowed lateness, snapshot at close | One org definition of calendar windows and time zones; change control for rule changes mid-season |
| Integrity | Validates payload | Server-authoritative, dedupe, quarantine, idempotent sinks | Dispute process, adjudicator, audit retention, legal sign-off on contest rules |
| Operations | Redis replicas | Update-lag SLO, canary entity, reconciliation job | Cost per board and per 1M events; chargeback; capacity calendar for events |
5.2 Strong Hire Signals#
| Signal | What It Sounds Like |
|---|---|
| Allocates exactness | "Exact where someone screenshots or gets paid; percentile everywhere else." |
| Knows rank doesn't shard | "With 64 hash shards, 'my rank' is 64 ZCOUNTs. I won't pay that per request." |
| Single-key ceiling | "One key, one thread, ~100K ops/s. Cluster size doesn't help. Pre-aggregate." |
| Idempotent sinks | "Flink holds the sum; the sink writes the absolute value. Replays can't double-count." |
| Snapshots for money | "Payouts come from a batch recompute at a fixed offset, never from the live ZSET." |
5.3 Lean No-Hire Signals#
| Signal | Why It Misses the Bar |
|---|---|
SQL COUNT(*) for rank at scale | O(N) per request; doesn't recognize the core problem |
| Cron job that deletes boards at midnight | Herd, blocking deletes, and lost events at the boundary |
| Client-submitted scores | No trust boundary |
ZINCRBY from at-least-once consumers without dedupe | Silent double counting |
| Exact rank for 1B members "because users want it" | No cost awareness, no stakeholder analysis |
5.4 Common False Positives#
- Skip-list internals ≠ leaderboard design. Knowing why ZSETs are O(log N) says nothing about what to do when one doesn't fit.
- Naming Count-Min Sketch ≠ knowing when to use it. Sketches are for unbounded key spaces (trending), not for a player's own score.
- Elaborate caching layers ≠ read scaling judgment. The top page needs one 1s cache; per-user rank must not be cached per user at 20M users.
- Kafka + Flink everywhere ≠ Staff. For a 1M-member board with 100 writes/s, Postgres and a minute-level refresh is the mature answer.
6. Interview Flow & Pivots#
6.1 Typical 45-Minute Shape#
| Phase | Time | Goal |
|---|---|---|
| Framing | 0–3 min | Three intents; commit to competitive board; exact top, approximate tail |
| Entities + API | 3–5 min | ScoreEvent, Board, RankView; rank vs percentile contract |
| High-level design | 5–12 min | Events → Flink → ZSET + histogram → cache |
| Transition | 12 min | Offer: rank at scale, hot writes, windows |
| Deep dives | 12–38 min | Sharding → hot keys → windows/late events → integrity |
| Org & evolution | 38–43 min | Payout snapshots, trust & safety ownership, counter platform |
| Wrap-up | 43–45 min | Summary; friends boards next; what not to build |
6.2 How Interviewers Pivot — And What They're Testing#
| Interviewer Pivot | What They're Testing | Where to Go |
|---|---|---|
| "Now it's 1 billion players." | Sharding and rank semantics | Tier + histogram; leagues |
| "One post gets 2M likes a minute." | Hot key reasoning | Pre-aggregation; split counters; per-key ceilings |
| "Show my friends' leaderboard." | Fan-out choice | Compute on read over ~200 friends; no per-user ZSET |
| "What about ties?" | Detail precision | Composite score encoding; earliest achiever wins |
| "We pay the top 100 $10K each." | Correctness bar shift | Batch recompute, snapshot, dispute process |
| "Trending posts in the last hour." | Unbounded key spaces | Space-Saving / CMS + heap per window |
| "Redis shard with the top tier dies." | Recovery | Replica promotion; rebuild from snapshot + Kafka replay |
6.3 What to Deliberately Skip#
| Topic | Why L5 Goes Here | What L6 Says Instead |
|---|---|---|
| Skip list internals | Feels rigorous | "O(log N) ZSET. The structure isn't the problem." |
| Pagination API design | Concrete | "Cursor by rank offset for the top tier. Standard." |
| Websocket push of rank changes | Exciting | "Poll top page every 5s with a 1s cache; push only for the viewer's own rank change if product asks." |
| Avatar/profile joins | Completeness | "Top page hydrates 100 profiles from a cache. Trivial." |
6.4 Follow-Up Questions to Expect#
- "How do you encode tie-breaking in a double score without losing precision?"
- "How do you rebuild the top-tier ZSET after losing a Redis primary and its replica?"
- "How would you implement a sliding 24-hour board?"
- "How do you show 'players around me' when I'm outside the exact tier?"
- "What does it cost to keep daily boards for 365 days?"
- "How do you remove a banned cheater from all historical boards?"
- "How do you run a board across three regions with one global ranking?"
7. Active Drills#
Drill 1: The Opening#
Prompt: "Design a leaderboard for our game."
Staff Answer
"Three questions shape this: is it a live competitive board, an engagement counter, or does money depend on it? And what must 'rank' mean — exact for everyone, or exact at the top? I'll assume a live board for 200M players with daily, weekly and all-time windows; exact top 100K within 2s; everyone else gets a percentile within ±0.1%. Scores come only from game servers as events. Payouts, if any, come from batch snapshots. I'll walk through the event pipeline, the two serving structures, hot writes, windows, and integrity."
Why this is L6:
- Defines rank semantics before choosing structures
- Separates display from payout correctness
- Establishes the trust boundary up front
What L7 adds:
- Asks whether leagues (segments of ~100) would serve engagement better than a global board — a product lever that removes a scaling problem
- Notes other teams' counter/ranking needs and whether to build a shared primitive
❌ Common L5 Trap
"Redis sorted set: ZADD on each score, ZREVRANGE for top 10, ZREVRANK for my rank. Redis Cluster for scale."
Why this misses: Correct at 10M members. Redis Cluster places each key on one slot — a single board ZSET doesn't spread across the cluster at all. The candidate hasn't noticed that "scale" and "one key" are in conflict.
Drill 2: Core Mechanic — Ties#
Prompt: "Two players both have 9,800 points. Who's first, and how do you implement it in a sorted set?"
Staff Answer
"Game design picks the rule; I'd default to 'earliest to reach the score wins.' ZSETs tie-break lexicographically by member, which is arbitrary, so I encode the rule into the score: composite = score × 2^22 + (2^22 − 1 − seconds_since_season_start) — higher score first, then earlier timestamp. A double has a 53-bit mantissa, so if scores fit in 31 bits I have 22 bits for time — ~4M slots, which at 1-second resolution covers ~48 days, enough for a season-scoped board. For longer or larger ranges, use a second-level sort: keep the ZSET score as the raw score and resolve ties in the top page by fetching tied members and sorting by achieved_at from a hash."
Why this is L6:
- Names that the tie rule is a product decision
- Knows the double-precision budget and does the bit math
- Has a fallback when the encoding doesn't fit
What L7 adds:
- Publishes the tie rule in the contest terms, because ties at #1 with prizes become legal disputes
Drill 3: "Top Tier + Histogram" — Make It Concrete#
Prompt: "You said approximate rank from a histogram. Walk me through the data structure and the error."
Staff Answer
"Per board and window: an array of 10K counters over score ranges, plus a count of members per bucket. On score change from s₁ to s₂: decrement bucket(s₁), increment bucket(s₂) — done in Flink state, flushed each second to Redis as a hash. Rank of score s = Σ counts of buckets above s + interpolated position within s's bucket (assume uniform within bucket). Error ≤ count of s's bucket. With 200M members and 10K fixed-width buckets, low-score buckets could hold millions, so bucket edges are adaptive: re-split any bucket over 0.1% of members (200K), merge sparse ones, recomputed hourly from the score distribution. Max error: 200K ranks out of 200M — 'Top 37.1%' is exact to one decimal."
Why this is L6:
- Provides a concrete structure, update rule, and error bound
- Handles skewed distributions with adaptive buckets
What L7 adds:
- Makes the error bound a product-visible contract ("percentile accurate to 0.1") so UX never displays false precision
Drill 4: Dependency Down — Top-Tier Redis Lost#
Prompt: "The Redis primary holding the top-tier ZSET dies and the replica was 30 seconds behind."
Staff Answer
"Promote the replica — the board is now up to 30s stale for the top tier. Because the sink writes absolute scores from Flink state, re-driving is simple: Flink replays from its last checkpoint offset and rewrites affected entities; within ~1 minute the board converges. If both are lost, rebuild from the latest 5-minute snapshot plus Kafka replay from the snapshot's offset: 100K members load in seconds. During rebuild, the top page serves the cached last-good copy with a 'refreshing' flag; rank reads fall back to the histogram. No data is lost because Redis was never the source of truth — the event log is."
Why this is L6:
- Redis is a projection, not the system of record
- Idempotent absolute writes make replay safe
- Degraded mode is explicit (cached top page, histogram rank)
What L7 adds:
- Sets the RTO in the product SLO ("top page may be stale ≤ 2 min during infra failure") and game-days it before major tournaments
Drill 5: Hot Key — Viral Counter#
Prompt: "A post gets 2M likes per minute. Your counter is one Redis key."
Staff Answer
"2M/min is ~33K/s on one key — a third of one Redis thread, and it degrades every other key on that shard. Fix in order: (1) Pre-aggregate in the stream: 1s tumbling per post → one INCRBY 33000 per second. (2) Like events must be idempotent per (user, post) — a like is a set membership, not an increment; the count is derived. Dedupe with a per-post set for small posts, and for hot posts a Bloom filter or partitioned set keyed by user hash. (3) Display count may lag by 1–5s and must be monotonic — the client never shows a lower number than it showed before. Unlikes are applied, but UI smooths them."
Why this is L6:
- Quantifies against the single-thread ceiling
- Recognizes likes are set membership (idempotency), not blind increments
- Adds the product constraint of monotonic display
What L7 adds:
- The counter platform auto-detects keys over 5K writes/s and flips them to the pre-aggregated path without the product team's involvement
Drill 6: Multi-Tenant — Boards as a Platform#
Prompt: "Now 40 games in the studio want leaderboards. One team's event generates 10× everyone else's traffic."
Staff Answer
"Kafka topic partitions and Flink parallelism are shared; one game's spike can lag everyone's boards. I'd give each game a quota on events/s at ingestion (429 back to the game server's publisher, which buffers), separate Kafka topics for the top 3 games by volume, and Redis key placement by {game} hash tags into per-game shard groups so a hot board only affects its own game. Update-lag SLO is measured per game. Big events are pre-registered on a capacity calendar 2 weeks ahead."
Why this is L6:
- Isolates by tenant at ingest, compute, and storage
- Measures SLOs per tenant, not globally
What L7 adds:
- Chargeback per 1M events and per GB of board storage, so games see the cost of 365 retained daily boards
Drill 7: Build vs Buy#
Prompt: "Why not use a managed game-backend leaderboard service?"
Staff Answer
"For one game with standard boards, I would. Managed game-services leaderboards and a managed Redis with sorted sets cover the 80% case. I'd build when we need: approximate-rank tiers for 200M+ members, custom windows and lateness, integration with our anti-cheat quarantine, or payout-grade snapshots tied to our event log. The build cost is roughly 2–3 engineers for 2 quarters plus on-call forever; the break-even is when the studio has multiple titles with those needs."
Why this is L6:
- Buys by default; names specific capability gaps that justify building
What L7 adds:
- Keeps the event contract (
ScoreEventschema, snapshots) ours regardless of the serving engine, so switching vendors is a projection rebuild, not a migration of truth
Drill 8: Policy Change Without Outage — Changing Scoring Rules Mid-Season#
Prompt: "Game design wants to change a scoring rule halfway through the season."
Staff Answer
"Two options with different fairness implications. Forward-only: new events use the new rule; existing scores stand — cheap, but players who scored early are advantaged or disadvantaged. Retroactive: recompute every player's score from the event log with the new rule — fair, but ranks jump overnight and it costs a full replay (200M players, ~1–2 hours in batch). Either way: version the rule (rule_version on each board), shadow-compute the new ranking first and show game design the top-1,000 diff, announce before switching, and snapshot the board at the switch for dispute reference. Game design and community teams sign off, not the platform."
Why this is L6:
- Surfaces fairness, not just mechanics
- Uses the event log as the enabler of retroactive change
- Shadow → diff → announce → switch
What L7 adds:
- Establishes a rule-change policy for all competitive products: no retroactive changes in the final 7 days of a season with prizes
Drill 9: Cost#
Prompt: "We keep daily, weekly, and all-time boards for 20 regions and 5 modes. What does storage cost?"
Staff Answer
"Boards = 20 × 5 = 100 per window. If every daily board held all 20M daily actives at ~100 B, that's 2 GB per board-day — 200 GB/day of Redis, clearly wrong to retain. With top-tier-only ZSETs (100K members ≈ 10 MB) plus histograms (80 KB), it's ~1 GB/day across all boards. Retain 7 days hot in Redis (~7 GB), then snapshots to object storage at cents per GB-month. Full per-player history lives in the event log / OLAP, not in Redis. The design choice of exact-top-only is what makes retention cheap."
Why this is L6:
- Shows the tiered design is also the cost design
- Moves history out of RAM
What L7 adds:
- Presents cost per board per month so product can decide whether 5 modes × 20 regions is worth it
Drill 10: Multi-Region#
Prompt: "Players are in three regions. We want one global board."
Staff Answer
"Score events are produced in each region. Options: (A) replicate all events to one home region that computes the global board — simple, one source of truth, +100–200ms cross-region lag, and the global board is unavailable if the home region is down (regional boards still work). (B) Each region computes its partial top tier and histogram; a global merger combines top-Ks (merge 3 lists) and sums histograms — histograms are mergeable by addition, which is why I like them. Rank = sum of per-region higher-bucket counts. I'd pick B: it degrades to 'global board missing one region' rather than 'no global board.' Payout boards are computed in batch from all regions' event logs."
Why this is L6:
- Exploits mergeability of top-K lists and histograms
- Chooses a degradation mode explicitly
What L7 adds:
- Checks data-residency rules for player IDs crossing regions; may need pseudonymous IDs in the global merge
8. Deep Dive Scenarios#
Deep Dive 1: Peak-Traffic Incident — The World Final#
Context: During the championship final, 4.2M concurrent viewers open the leaderboard. Top-page p99 is 3s and rising; score updates are lagging 90s. Game servers for 12K concurrent matches keep writing.
Questions to Surface First:
- Is the pressure on reads (top page, rank) or writes (score updates)?
- Is the top-page cache effective, or are misses stampeding the top-tier shard?
- Which shard(s) are hot, and what else lives on them?
Typical L5 Approach: Adds read replicas and scales Redis. Replica sync adds more load on the stressed primary; the fix takes 20 minutes.
Staff Approach: Checks cache hit rate first: finds the 1s TTL expires simultaneously on 40 API nodes, each missing into the shard (40 ZREVRANGE/s — fine), but the "around me" endpoint is uncached and 4M viewers call it on page load. Disables "around me" via feature flag for non-participants (viewers aren't on the board), keeps it for players. Load drops 90% in 1 minute.
Principal Approach: Separates viewer traffic from player traffic as product surfaces with different SLOs — a spectator view served from CDN-cached snapshots every 2s — and makes tournament events a capacity-calendar item with a pre-event load test.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate (0–5 min) | Split read vs write latency; check per-endpoint QPS; top-tier shard CPU. |
| Triage | around endpoint: 1.1M QPS, 95% from non-participants who have no rank. |
| Quick fix | Flag off around for users without a score in this board; return top page only. |
| Guardrails | Single-flight cache refresh; stale-while-revalidate 5s; per-endpoint rate limits. |
| Post-mortem | Why did viewers call a player-only endpoint? Why no spectator mode? Why no load test? |
Metrics to Watch: api.qps{endpoint}, cache.hit_ratio{endpoint}, redis.shard_cpu{shard}, board.update_lag_seconds
Organizational Follow-up: Spectator surface owned by the esports product team; capacity review 2 weeks before every tournament.
Ownership Question: "Who decides to turn off 'around me' during the final?" Staff answer: The platform on-call, under a pre-approved degradation runbook listing which features can be shed in which order. Product agreed to that order in advance.
Key Takeaway: "At peak, your largest audience is often not your users of record. Find who's reading and whether they need what they're reading."
What clears the Staff bar:
- Finds the uncached endpoint rather than scaling everything
- Uses a pre-agreed shed order
- Converts the incident into a product surface split
Deep Dive 2: Silent Failure — Weekly Sums Drifting Up#
Context: A player reports that their weekly points are 12% higher than their match history adds up to. Spot checks show about 3% of players are over-counted, all this week.
Questions to Surface First:
- When did the drift start, and does it correlate with a deploy or a consumer restart?
- Does the sink use increments or absolute values?
- Is the reconciliation job running, and why didn't it alert?
Typical L5 Approach: Fixes the affected players' scores with a script.
Staff Approach: Correlates the drift with a Flink restart on Tuesday that replayed 4 minutes of events into a
ZINCRBYsink added by a new board type that bypassed the standard absolute-value sink. Recomputes the week's sums from the event log, rewrites the board, and makes the absolute-value sink mandatory (increment sinks fail the pipeline lint). The reconciliation job existed but sampled only all-time boards — expands it to every board type.
Principal Approach: Makes idempotent sinks a platform invariant enforced in CI and notes the governance gap: a new board type shipped without the platform's review. Adds board-type registration with required reconciliation coverage.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate | Freeze publication of this week's snapshot; notify contest owner if prizes attach. |
| Triage | Diff board vs event-log recomputation; affected set = players active in the replay window. |
| Quick fix | Batch recompute weekly sums; rewrite ZSET with absolute scores; regenerate snapshots. |
| Guardrails | Pipeline lint: sinks for sum boards must be absolute-value; reconciliation on all board types. |
| Post-mortem | Why could a board type bypass the standard sink? Why did reconciliation cover only one type? |
Metrics to Watch: board.reconcile_mismatch_ratio{board_type}, flink.restarts_total, sink.mode{board}
Organizational Follow-up: Player communication from community team; add "replay-safe" check to launch review.
Ownership Question: "Who owns the correctness of weekly sums?"
Staff answer: The platform owns the pipeline invariant (idempotent sinks, reconciliation). The game team owns choosing sum and reviewing the board-type registration.
Key Takeaway: "At-least-once delivery plus increment writes equals silent inflation. Write absolute values."
What clears the Staff bar:
- Links drift to replay semantics immediately
- Fixes the invariant, not the rows
- Extends detection coverage to where it was missing
Deep Dive 3: Large-Customer Onboarding — The Mega-Title Launch#
Context: A new title expects 80M players in the first month and a single global board — 4× the current largest board. Game design insists on "exact rank for everyone, it's core to the fantasy."
Questions to Surface First:
- What does exact rank at #52,000,000 give a player that "Top 65%" doesn't?
- Would leagues (divisions of 100) meet the design goal?
- Are prizes attached?
Typical L5 Approach: Provisions a large Redis node for an 80M-member ZSET (~8–10 GB) and hopes the thread holds.
Staff Approach: Shows the numbers: 80M-member ZSET fits memory but all rank reads (projected 150K/s at launch) hit one thread — over the ceiling. Proposes exact top 100K + histogram, plus exact rank inside leagues of 100 (where most players' competitive attention actually sits). Offers a UX test comparing "#52,114,209" vs "Top 65% · #3 in your league."
Principal Approach: Makes league-based ranking the studio's default pattern for new titles, backed by engagement data, so the platform never needs single boards beyond ~50M members.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Week 1 | Load model: rank reads 150K/s vs ~100K/s thread ceiling; memory OK; single point of heat. |
| Negotiate | Present options with cost and UX: exact-all (read replicas + scatter, 3× cost), tier+histogram, leagues. |
| Build | Leagues: 800K small ZSETs, each ~10 KB — spread across the cluster by hash tag; top tier global. |
| Test | Load test at 2× launch estimate; kill a shard during test. |
| Launch | Feature-flag exact-rank display for tier members only. |
Metrics to Watch: rank.read_qps{board}, redis.shard_cpu, league.assignment_lag
Organizational Follow-up: Game design signs off on rank display contract; community team prepares FAQ.
Ownership Question: "Who decides whether players see exact rank?" Staff answer: Game design decides the experience; the platform provides the cost of each option. If they choose exact-all, the title's budget pays the 3×.
Key Takeaway: "The cheapest scaling fix is often a product change that users prefer anyway."
What clears the Staff bar:
- Quantifies the thread ceiling, not just memory
- Reframes the requirement with UX evidence
- Attaches cost to the product choice
Deep Dive 4: Post-Mortem — The Wrong Winner Got Paid#
Context: A $50K sales-contest prize was paid to the rep shown #1 on the live board at the deadline. Two days later, refund events arrived that dropped them to #3. The rightful winner escalated to the VP.
Questions to Surface First:
- What were the published contest rules on cutoff and refunds?
- Was the payout taken from the live board or a snapshot?
- Is there a defined dispute window?
Typical L5 Approach: Adds refund handling to the live board so ranks update correctly.
Staff Approach: Identifies the real bug: payouts read the live board, and the contest had no rule for events that arrive after the deadline but refer to activity before it. Proposes: payout board computed in batch at deadline + N days (refund window, e.g., 14 days), from the event log, with a published provisional ranking at the deadline and a final ranking after the window. Snapshot signed and stored with the offset.
Principal Approach: Establishes a standard for any money-bearing ranking: provisional vs final results, a dispute window, a named adjudicator (sales ops, not engineering), and legal review of contest terms. The platform only publishes; it never decides.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate | Recompute from event log including refunds; confirm correct winner; hand to sales ops + finance. |
| Triage | Payout script queried ZSET at 23:59:59 on deadline day; refund events not in scope. |
| Fix | Provisional/final model; payout reads only final snapshots. |
| Guardrails | Payout API refuses to read non-final boards. |
| Post-mortem | Contest terms lacked refund treatment; engineering implemented an undefined rule. |
Metrics to Watch: contest.final_vs_provisional_rank_changes, payout.source{snapshot|live}
Organizational Follow-up: Contest template with required cutoff, refund, tie and dispute clauses.
Ownership Question: "Who owns the decision of who won?" Staff answer: The contest owner (sales ops) with finance, from the final snapshot. Engineering owns that the snapshot is reproducible.
Key Takeaway: "A live board is a display. A payout needs a cutoff, a settling window, and an owner who isn't the database."
What clears the Staff bar:
- Finds the missing rule, not just the missing code
- Separates provisional and final results
- Moves adjudication to the business owner
Deep Dive 5: Multi-Region Expansion — Global Board With Residency#
Context: The game launches in the EU and Asia. Player identifiers are personal data; EU counsel says EU player IDs shouldn't be replicated to the US. Product wants a single global top-100.
Questions to Surface First:
- What must the global board show — display names, IDs, both?
- Can we use pseudonymous board IDs?
- Is global rank needed, or just global top-100 + regional rank?
Typical L5 Approach: Replicates all score events to a US region to compute the board.
Staff Approach: Computes regional top tiers and histograms in-region. The global merger receives only pseudonymous board IDs (
hmac(region_key, player_id)) and scores for each region's top 1,000, plus histogram counts (aggregates, not personal data). Global top-100 = merge of 3 × 1,000. Display names are resolved in the viewer's request by calling the owning region's profile service.
Principal Approach: Takes the "aggregate-only crosses borders" pattern into the org's data-residency standard so every future global feature (trending, counts) reuses it.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Design | Regional pipelines; global merger consumes pseudonymous top-K + histograms. |
| Correctness | Global top-100 exact (each region's top 1,000 ≥ any global top-100 member from that region). |
| Failure | Region down → global board marked partial, still served. |
| Privacy | Counsel reviews pseudonymization; key rotation plan. |
| Cost | Merger traffic: ~3K entries/s — negligible. |
Metrics to Watch: global_merge.lag_seconds{region}, global_board.partial, residency.cross_region_pii_bytes (must be 0)
Organizational Follow-up: Privacy review sign-off; documentation of the partial-board UX.
Ownership Question: "Who approves sending pseudonymous IDs cross-region?" Staff answer: Privacy counsel and the data-protection officer. Engineering proposes; they approve.
Key Takeaway: "Mergeable aggregates — top-K lists and histograms — let rankings go global without moving personal data."
What clears the Staff bar:
- Uses mergeability as a privacy tool
- Explicit partial-board degradation
- Routes the residency decision to its owner
9. Level Expectations Summary#
After studying this case study, you should be able to:
- Separate competitive boards, engagement counting, and money-bearing rankings, and commit to one in the first 2 minutes
- Allocate exactness: exact top tier, approximate tail with a numeric error bound
- Explain why hash-sharding breaks global rank and choose between tier+histogram, score-range shards, and leagues
- Quantify the single-key ceiling and fix hot writes with pre-aggregation and split counters
- Implement windows as keys with event-time assignment, allowed lateness, and snapshots at close
- Make stream sinks idempotent by writing absolute values
- Define the trust boundary (server-authoritative scores, quarantine) and the payout boundary (batch snapshot, dispute process)
The Bar for This Question#
Mid-level (L4): SQL table with an index; ORDER BY for top-K; COUNT(*) for rank. Works at 100K members.
Senior (L5): Redis sorted sets, ZADD/ZREVRANK, caching the top page, sharding "by user." Competent and fast at 10M members. Stalls on "rank across shards," "one viral key," and "which number do we pay from."
Staff+ (L6): Treats rank as a contract with different bars for different parts of the board, designs the write path around key skew, makes windows and replays safe by construction, and names owners for rules, cheating, and payouts. The interviewer should learn something from the answer.
10. Staff Insiders: Controversial Opinions#
10.1 Nobody Needs Exact Rank Below the First Page#
| Position on board | What the player does with the number |
|---|---|
| #1–#100 | Screenshots, streams, disputes, gets paid |
| #101–#100K | Tracks progress; exact is nice |
| #100K–#200M | Wants "am I improving?" — percentile answers better |
The Staff position: Exact rank in the tail is a cost with no stakeholder. "Top 37%" is more motivating than "#74,112,093" and costs O(1).
Why this matters in interviews: Proposing it — with an error bound — is the clearest signal that you allocate correctness deliberately.
10.2 Global Leaderboards Are a Product Anti-Pattern at Scale#
The Staff position: At 100M players, a global board is demotivating for 99.99% of them. Leagues of ~100 are both better engagement and trivially scalable — many small exact ZSETs spread naturally across a cluster.
Why this matters in interviews: Changing the requirement is sometimes the best engineering move. Say it with data, and offer the global top tier as well.
10.3 Counters Should Never Be Increments at the Store#
| Write style | Replay-safe? | Hot-key safe? |
|---|---|---|
INCR per event | No | No |
INCRBY batched per second | No | Yes |
| Absolute value from stream state | Yes | Yes |
| Derived from set membership (likes) | Yes | Needs partitioned sets |
The Staff position: The stream owns the count; the store receives absolute values. Increments at the store are how counters drift silently.
Why this matters in interviews: It shows you've reasoned about at-least-once delivery end to end.
10.4 The Live Board Should Never Pay Anyone#
The Staff position: Real-time and reproducible are different products. Money flows from batch-recomputed, snapshot-backed, provisional-then-final results.
Why this matters in interviews: It separates display correctness from financial correctness — a distinction L5 answers almost never make.
10.5 Count-Min Sketch Is Overused in Interviews#
The Staff position: Sketches are for heavy-hitter detection over unbounded key spaces (trending). A like count for a known post should be exact and cheap with pre-aggregation. Reaching for a sketch where a HINCRBY batch works is complexity theater.
Why this matters in interviews: Picking the simplest structure that meets the bar is the Staff instinct.
11. The Principal Lens (L7)#
Why L7 Sees This Problem Differently#
At Staff level, a leaderboard is one system with an exactness budget. At Principal level, it's one instance of a pattern the whole company reinvents: like counts, view counts, trending, seller rankings, creator funds, sales contests, gamified fitness streaks. Each team builds its own counter on its own store, discovers the hot-key ceiling during its first viral event, and ships its own replay bug. Two of them pay people from live numbers. The L7 move is to recognize counting and ranking as a platform primitive with declared correctness tiers, and to separate "numbers we display" from "numbers we pay on" as an org-wide rule.
The Org-Level Fault Line#
A shared counting & ranking platform vs per-product implementations.
| Option | What It Buys | What It Costs |
|---|---|---|
| Shared platform with tiers (exact, pre-aggregated, sketched) | Hot-key handling, idempotent sinks, reconciliation, and snapshots built once | Platform team of 4–6; product teams give up schema freedom |
| Per-product | Autonomy; tailored semantics | Same incidents rediscovered by each team; inconsistent counts across surfaces |
| Shared library, team-run | Reuse without a central team | Version skew; no one owns reconciliation |
The Principal position: Shared platform for counting and display rankings; money-bearing rankings stay in the owning business domain but must consume the platform's event log and snapshot format.
Cost Model#
Assumptions: managed Redis ~$0.02/GB-hour-equivalent pricing tier, shared Kafka/Flink clusters with allocated share, engineer fully loaded ~$250K/year.
| Scale | Volume | Infra $/month | Headcount | On-call Load |
|---|---|---|---|---|
| Small | 1M members, 500 writes/s, Postgres + Redis ZSET | ~$1–2K | 0.25 FTE | Rare |
| Medium | 200M members, 200K events/s peak, Kafka + Flink + Redis tiers | ~$30–60K | 3–4 engineers | 1–3 pages/month; event-driven spikes |
| Large | Studio/company platform: 2M+ counter events/s, 1,000+ boards, multi-region | ~$200–350K | 6–8 engineers | Dedicated rotation; tournament war rooms |
Exact-rank-for-everyone at Medium scale roughly triples Redis cost (full-member ZSETs per window, scatter-gather replicas) — ~$60–100K/month extra. That's the number to show when product asks for it.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversibility Cost |
|---|---|---|
| Scores as an immutable event log | One-way (good) | Without it, no replay, no audit, no retroactive fixes |
| Publicly exposing exact rank for everyone | One-way | Taking it away later is a user-visible regression and a community backlash |
| Tie-break rule in a live contest | One-way mid-season | Changing it reorders winners |
| Window definition (UTC, ISO week) | One-way-ish | Historical boards become incomparable |
| Serving store (Redis vs other) | Two-way | Rebuild projection from the event log |
| Histogram bucket scheme | Two-way | Recompute from scores |
| Pre-aggregation window (1s vs 5s) | Two-way | Config |
The Standard I'd Write#
RFC: Counting and Ranking Standard (v1)
Scope: Any user-visible count or ranking, and any ranking used for compensation or prizes.
MUST:
- Source all counts and scores from an append-only event log with unique
event_ids.- Write stores with idempotent operations (absolute values or set membership); no store-side increments from at-least-once consumers.
- Declare a correctness tier per surface:
exact,exact-top-k + approximate, orapproximate, with the error bound shown in the API contract.- Rankings used for money MUST be computed from a final snapshot after a declared settling window, with a named business adjudicator.
SHOULD:
- Route counter writes through the platform's pre-aggregation path.
- Run reconciliation against the event log for every board type at least hourly.
- Use tumbling UTC windows unless product justifies otherwise.
Exceptions: Approved by the counting platform tech lead; money-bearing exceptions also require finance sign-off.
Success metrics: zero payouts from live boards;
reconcile_mismatch_ratio< 0.01% on all boards; no hot-key incidents affecting co-located tenants for two consecutive quarters.
What I'd Tell the VP#
Our leaderboards and counters work, but every team builds its own, and we've had three incidents this year from the same two bugs: viral traffic overwhelming a single counter and replays inflating scores. Last quarter we paid a contest prize to the wrong person because the payout read a live number. I'm proposing one shared counting platform and one rule: numbers we display can be approximate and fast, but numbers we pay on come from a settled, reproducible snapshot owned by the business. It's a team of four and roughly $40K/month in infrastructure, and it removes an entire class of incidents plus a real financial and legal exposure.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Sees the pattern across products | "Likes, views, contests, and game boards are the same primitive with different correctness tiers." |
| Prices exactness | "Exact rank for everyone is ~$60–100K/month more at this scale. Who's funding it?" |
| Separates display from money | "Displayed numbers and paid numbers are governed differently." |
| Uses product levers | "Leagues remove the scaling problem and improve engagement." |
| Identifies one-way doors | "Once we show exact rank to 200M people, we can't take it back." |
Staff answers that L7 interviewers find insufficient:
- "We'll build a great leaderboard service for this game" — no view of the other five teams building counters.
- "Payouts use a snapshot" — correct, but no settling window, adjudicator, or contest-terms ownership.
- "We'll shard by region for multi-region" — no residency analysis, no mergeable-aggregate reasoning.
Appendices
Appendix A: Mechanics in Depth#
A.1 ZSET Operations for a max Board#
# Update best score (Redis 6.2+): only raise, never lower
ZADD lb:{board}:2026-W40 GT <composite_score> <player_id>
# Top 100
ZREVRANGE lb:{board}:2026-W40 0 99 WITHSCORES
# My rank (0-based)
ZREVRANK lb:{board}:2026-W40 <player_id>
# Around me (±5)
r = ZREVRANK ...; ZREVRANGE lb:... max(0,r-5) r+5 WITHSCORES
# Trim top tier to 100K members
ZREMRANGEBYRANK lb:{board}:2026-W40 0 -100001
ZADD GT is idempotent for max boards — replays can't lower or inflate.
A.2 Score Histogram#
state: edges[0..B], counts[0..B-1] # B = 10,000
on_update(old, new):
if old != null: counts[bucket(old)] -= 1
counts[bucket(new)] += 1
approx_rank(s):
b = bucket(s)
higher = sum(counts[b+1..B-1]) # maintain suffix sums for O(1)
frac = (edges[b+1] - s) / (edges[b+1] - edges[b])
return higher + frac * counts[b]
rebalance hourly: split buckets with counts > 0.1% of N; merge adjacent sparse buckets
A.3 Sharded Counter#
write(key, n): INCRBY key:{hash(event_id) % S} n # S = 16 for hot keys
read(key): sum(MGET key:0 .. key:S-1) # or cached sum, 1s TTL
promote: if write_rate(key) > 5K/s for 30s → S = 16; demote when < 500/s for 1h
A.4 Heavy Hitters (Trending)#
Count-Min Sketch (w=2,000, d=5, ~40 KB) per window gives count estimates with error ≤ 0.14% of window total at 99.3% confidence; a min-heap of the top 1,000 candidates is updated when an item's estimate exceeds the heap minimum. Windows of 5 minutes merged into an hour by summing sketches cell-wise. Space-Saving with m = 10 × K counters is the alternative when you need guaranteed top-K membership with bounded error.
Appendix B: Keys and Data Model#
| Record | Key | Notes |
|---|---|---|
| Top tier | lb:{board}:{window_id} | Hash tag {board} groups a board's windows on one slot group; TTL = window close + 7d |
| Histogram | lbh:{board}:{window_id} | Redis hash or array of counters; 80 KB |
| Tier threshold | lbt:{board}:{window_id} | Score of the Kth member; updated each second |
| Player score (source) | Flink keyed state / Cassandra (board, window, player) | Absolute score, event offset |
| Counter | cnt:{entity}:{metric}[:{shard}] | Split when hot |
| Snapshot | object storage board/window/offset.parquet | Immutable; signed for payout boards |
Window IDs: all, d2026-09-30, w2026-W40, s2026-S3 (season).
Appendix C: Sharding Strategies — Quick Comparison#
| Strategy | Top-K Cost | Rank Cost | Rebalancing | Best For |
|---|---|---|---|---|
| Single ZSET | O(K log N), 1 call | O(log N), 1 call | None | ≤ 50M members |
| Hash shards (S) | Merge S lists | S × ZCOUNT | Rare | Top-K only, rank rarely needed |
| Score-range shards | Top shard | Cached offsets + 1 call | Continuous | Stable distributions |
| Top tier + histogram | 1 call | 1 call (exact or approx) | Hourly bucket re-split | Large boards, percentile OK |
| Leagues | Per league | Per league | League assignment | Engagement-first games |
Appendix D: API Contract#
rankis returned only when exact; otherwisenullwithpercentileandpercentile_error(e.g., 0.1).- Top page:
Cache-Control: max-age=1, stale-while-revalidate=5. as_oftimestamp andboard_version(Kafka offset) on every response — clients can detect staleness; support can reproduce disputes.- Payout endpoints return
status: provisional | final; onlyfinalis valid for payment.
Appendix E: Observability#
board.update_lag_seconds{board} # event_time → visible; SLO p99 < 2s for top tier
board.reconcile_mismatch_ratio{board_type}
redis.shard_cpu{shard}, hot_key.ops_per_sec{key}
rank.approx_error_estimate{board}
board.top1_score_jump_ratio{board}
canary.rank_update_seconds # synthetic player on a hidden board
cache.hit_ratio{endpoint}
| Alert | Threshold | Routes To |
|---|---|---|
| Update lag | p99 > 30s for 5 min | Platform |
| Reconcile mismatch | > 0.1% | Platform streaming |
| Hot key | > 10K ops/s for 1 min | Counter platform |
| Top-1 jump | > 1.5× previous #1 | Trust & safety |
| Canary stale | > 10s | Platform |
Appendix F: Scale Evolution#
| Scale | Design |
|---|---|
| < 1M members | Postgres + materialized rank every 60s, or one Redis ZSET |
| 1–50M | One ZSET per board-window; direct writes with dedupe; 1s top cache |
| 50M–1B | Event log + Flink; exact tier + histogram; leagues; snapshots |
| Multi-product / multi-region | Counter & ranking platform; mergeable aggregates; payout standard |
What you don't build on day one: sliding windows, sketches for trending, multi-region merge, exact rank for the tail, real-time push of rank changes.
Appendix G: Multi-Tenancy, Fairness, and Cost#
- Per-tenant ingestion quotas (events/s) and board count limits; 429 to publishers with buffering.
- Hash-tag placement groups each tenant's keys; a hot tenant heats only its slot group.
- Retention tiers: 7 days hot in Redis, snapshots in object storage for 13 months, event log in OLAP.
- Chargeback: per 1M events ingested, per GB-month of hot board storage, per board-window retained.
- Fairness signal:
ingest.share{tenant}> 50% of cluster for 10 min → notify tenant and platform.