System Design Problem

Design Top K Rankings System (App Store / Amazon Bestsellers)

Commonly Asked By:AppleAmazonGoogleNetflixTwitter

Interview Setup

Interview Prompt

Design a top-K rankings system like the App Store or Amazon Bestsellers. Track 1B purchase/download events per day across 5M items and serve top-10/50/100 lists by category and time window within minutes.

Clarifying Questions (ask before designing)

QuestionWhy it matters
Must rankings be exact, or is approximate within 1% acceptable?Exact accuracy requires dedicated per-item counters, whereas approximate counting unlocks Count-Min Sketch and suits trending boards where minor rank boundary differences do not impact user experience.
Which time windows are required, such as real-time, hourly, daily, weekly, or all-time?Each window requires a dedicated aggregation pipeline, because all-time rankings require cumulative storage whereas an hourly window uses a tumbling or sliding window.
Do we need historical rankings to answer queries such as what ranked number one last Tuesday?This requires persisting Cassandra snapshots per window, which increases storage requirements but represents a standard feature in app store leaderboards.
How many independent ranking lists exist across categories, regions, and metrics?Caching 1000 boards with 100 entries each requires only 20 MB of memory, whereas scaling to 100K boards requires a distributed multi-level merge strategy.

Scope

In scope

  • Min-heap per partition
  • Map-reduce merge
  • Approximate algorithms
  • Time-decay scoring
  • Multi-level aggregation
  • Capacity estimation with shown math

Out of scope (state explicitly)

  • Detailed frontend/UI pixel implementation
  • Org structure, staffing, and hiring plan

Functional Requirements

Start by asking your interviewer which metrics you rank by, which time windows matter, and whether approximate counts are acceptable. Historical snapshots and multi-dimensional boards are common follow-ups.

  • Compute and display top K items (K = 10, 50, 100) ranked by a chosen metric (sales, downloads, ratings, revenue), distinguishing between near-real-time approximate trending views and authoritative exact historical leaderboards
  • Support multi-dimensional ranking across category (e.g., games, productivity, overall), geography (e.g., US, UK, global), time period (hourly, daily, weekly, monthly, all-time), and metric (downloads, revenue, ratings, active users)
  • Near real-time updates ensuring rankings reflect recent activity within minutes
  • Support historical point-in-time rankings to resolve queries such as identifying what ranked number one on a past date
  • Handle large scale across millions of items and billions of events such as purchases and downloads
  • Expose pre-computed ranked lists via low-latency client APIs

Non-Functional Requirements

Your interviewer will care most about read latency on precomputed boards and event ingestion throughput. Frame this as a streaming problem early, because storing every item counter in memory is impossible at billion-event scale.

  • Low Latency: Return top K lists in under 50 ms
  • High Availability: 99.99% uptime for read paths
  • Scalability: Process billions of ranking events per day
  • Accuracy: Real-time streaming leaderboards use approximate Count-Min Sketch frequency estimation paired with Top-K heap state and may exhibit small boundary rank variance. In contrast, official daily, weekly, and historical rankings are reconciled using exact batch aggregation, with Redis continuously serving the precomputed ranking view for clients.
  • Freshness: Update rankings every 1 to 5 minutes
  • Read-Heavy: Millions of users viewing rankings with relatively few events generated per user

Capacity Estimations

Run this math before choosing between exact and approximate counting. Events per second and unique item cardinality determine whether a min-heap alone suffices or requires a Count-Min Sketch. Here, 1000 leaderboards serves as the baseline capacity assumption for initial back-of-the-envelope sizing (such as 200 categories across 5 time periods with default global geography and default downloads metric), while 12K and 20K represent expanded multi-dimensional deployment scales discussed later.

MetricCalculationValue
Total itemsGiven (assumption documented in value)5M
Events (purchases/downloads) / dayGiven (assumption documented in value)1B
Events / sec1B ÷ 86400~12K (peak 60K)
Event sizeGiven (assumption documented in value)100 bytes
Top K lists to maintain1000 (categories x time periods)1000 (categories x time periods)
List storage1000 x 100 entries x 200B20 MB
Event storage / day1B x 100B100 GB

Architecture Diagram

In the room: ask whether rankings must be exact or approximate, because that single requirement dictates the entire aggregation layer architecture.

Walk your interviewer through the diagram by write vs read paths. Events enter through Kafka and get aggregated in Flink using min-heaps and optional Count-Min Sketches. We flush precomputed top-K boards to Redis sorted sets so reads never sort at query time. An hourly Spark job reconciles approximate counts against raw events and writes historical snapshots to Cassandra.

Loading...

Component Deep Dives

Next we walk through each component on the diagram. The aggregation layer forms the central pivot of the architecture, requiring a structured discussion of event ingestion, stream processing, and read serving in sequence.

Event Ingestion (Purchase and Download Services)

Events enter the pipeline here, where keeping the ingestion path asynchronous ensures client writes and downstream reads remain fast.

  • Validate incoming events through verified purchase confirmation and deduplicate by device identifier within a 24-hour window for free downloads
  • Publish to the Kafka ranking-events topic using item_id as the partition key to guarantee item-level ordering
  • Return HTTP 202 Accepted immediately because ranking aggregation is fully asynchronous
  • Fraud filter: Velocity caps per IP and account are enforced in Flink before events qualify for heap updates

Stream Processing (Apache Flink): The Core Engine

This is where event aggregation happens, with the choice between exact and approximate counting serving as the critical architectural decision.

Approach 1: Exact Count Pipeline (for small-to-medium scale)

Flink Pipeline Topology (Dimension-Aware Multi-Board Aggregation):
  Source: Kafka topic 'ranking-events' (64 partitions distributing ingest throughput)
  KeyBy: (category, geography, metric, item_id) for item-level event accumulation
  Window: Tumbling or sliding window (e.g., 5-minute sliding or 24-hour tumbling)
  Aggregate: Incremental count accumulator per item within each dimension combination
  Keyed Process / Top K: Independent min-heap of size K per logical leaderboard key:
    Partition Worker
      ├── Dimension Key A (e.g., games:us:daily:downloads) → Min-Heap (size K)
      ├── Dimension Key B (e.g., productivity:global:hourly:revenue) → Min-Heap (size K)
      └── ...
  Merge Stage: Downstream reduce operator groups partial heaps by matching leaderboard key
  Sink: Redis Sorted Sets updated per board (ranking:{category}:{geography}:{period}:{metric})

Approach 2: Approximate Count (for massive scale): Count-Min Sketch + Heap

When tracking millions of unique items, maintaining exact counters for every entity in memory becomes cost-prohibitive.

Count-Min Sketch:

  • Probabilistic data structure consisting of d hash functions across a matrix of w counters
  • On an event for item X, hash the item key with each of the d functions and increment the corresponding d counters
  • To query the frequency of item X, evaluate all d mapped counters and take their minimum value
  • Space: O(w x d), requiring typically 2 MB of memory for under 0.1% error rate
  • Error: The sketch always overcounts and never undercounts, with error bounded mathematically by ε = e / w

Min-Heap for Top K:

TYPESCRIPT
interface RankedItem {
  itemId: string;
  count: number;
}

// Maintain a min-heap of size K over incoming streaming events for a specific leaderboard
function updateTopKHeap(heap: MinHeap<RankedItem>, item: RankedItem, k: number): void {
  if (heap.size() < k) {
    heap.insert(item);
  } else if (item.count > heap.peek().count) {
    heap.poll();
    heap.insert(item);
  }
}

Combined Streaming Approach:

  1. Count-Min Sketch tracks approximate counts across all items with minimal memory consumption
  2. A Min-Heap of size K maintains the running set of current top candidates
  3. When a new event arrives:
    • Increment the item count across sketch counters
    • Query the minimum counter estimate to check if the item qualifies for the top-K heap
    • If the count exceeds the heap root, insert or update the item in the heap
  4. Periodically flush the local heap state to Redis sorted sets

Multi-Dimensional Keyed State Architecture:

Because the system supports multiple dimensions such as category, geography, time period, and metric (potentially yielding 12K to 20K independent ranking lists), streaming aggregation is keyed by the specific leaderboard dimensions:

Partition Worker
  ├── Dimension Key A (e.g., games:us:daily:downloads) → Top-K Min-Heap State
  ├── Dimension Key B (e.g., productivity:global:hourly:revenue) → Top-K Min-Heap State
  └── Dimension Key C (e.g., education:uk:weekly:ratings) → Top-K Min-Heap State

Kafka partitioning provides distributed ingest throughput across 64 partitions. Flink keyed state using KeyBy(category, geography, metric, item_id) maintains independent aggregation state for each relevant leaderboard. Crucially, the system does not maintain a single heap that incorrectly mixes unrelated leaderboards together. Downstream, local Top-K heap results from partition workers are merged exclusively for their corresponding category, geography, metric, and time window.

Time-Windowed Rankings

The choice between sliding and tumbling windows fundamentally changes how trending content feels to users, making it important to justify your selection.

YAML
# Time-windowed aggregation configurations
hourly_rankings:
  window_type: "sliding"
  duration: "1 hour"
  slide_interval: "5 minutes"
  use_case: "Near-real-time trending leaderboards"

daily_rankings:
  window_type: "tumbling"
  duration: "24 hours"
  use_case: "Official daily top lists"

weekly_rankings:
  window_type: "tumbling"
  duration: "7 days"
  use_case: "Weekly chart summaries"

all_time_rankings:
  window_type: "cumulative"
  duration: "unbounded"
  use_case: "All-time popular items without window resets"

Exponential Decay for Trending Calculations:

Instead of rigid fixed windows, exponential decay weights recent activity higher than historical events:

PYTHON
# Exponential decay scoring formula:
# score = sum(event_value * math.exp(-decay_rate * age_in_hours))
# decay_rate (lambda) = 0.1 yields a half-life of approximately 7 hours
# Recent events contribute substantially higher weight than historical events
import math

def compute_decay_score(event_value: float, age_in_hours: float, decay_rate: float = 0.1) -> float:
    return event_value * math.exp(-decay_rate * age_in_hours)

Redis (Ranking Cache)

Pre-computed rankings reside in memory so client reads never require on-the-fly sorting at query time.

Each ranking leaderboard is backed by a Redis Sorted Set with atomic update operations:

REDIS
# Structure: Key = ranking:{category}:{geography}:{period}:{metric}, Score = metric score, Member = item_id

# Fetch top 10 items in games by downloads in the US for today
ZREVRANGE ranking:games:US:daily:downloads 0 9 WITHSCORES

# Look up an individual item's current weekly rank across global overall downloads
ZREVRANK ranking:overall:global:weekly:downloads app123

# Atomically update or insert an item score from Flink flush
ZADD ranking:games:US:daily:downloads 15432 app123

Ranking Service (Read API)

The read API delivers pre-computed leaderboards from cache, keeping computationally heavy sorting off the critical hot path.

  • Stateless API instances serve read requests directly from Redis cluster replicas
  • Querying ZREVRANGE ranking:{category}:{geography}:{period}:{metric} 0 K-1 returns the top-K leaderboard in under 5 milliseconds
  • ZREVRANK executes fast single-item rank lookups across multiple leaderboards concurrently
  • Historical point-in-time queries route to Cassandra snapshots generated by the hourly Spark batch process

Batch Pipeline (Spark): Reconciliation

An hourly batch reconciliation pipeline corrects any mathematical drift introduced by the approximate real-time stream processing layer.

  • Hourly Spark batch jobs read raw, immutable event logs from Cassandra bucketed partitions in parallel (and alternatively from cloud object storage such as S3 in Parquet format), avoiding single monolithic partition scans
  • The job computes exact rankings across all dimensions without probabilistic approximations
  • It reconciles existing Redis sorted sets, correcting any Count-Min Sketch overcount drift
  • The pipeline persists immutable historical ranking snapshots into Cassandra to support point-in-time queries

Event Bus Design (Kafka)

The event bus decouples ingestion producers from downstream processing consumers and absorbs burst traffic spikes.

YAML
topic: "ranking-events"
partitions: 64
partition_key: "item_id" # guarantees ordering per item within a partition
retention: "7 days"      # batch reconciliation replays from Cassandra when needed
replication_factor: 3
min_insync_replicas: 2

producer:
  idempotence: true # enable.idempotence=true to prevent duplicate counts
  payload:
    event_id: "string (UUID)"
    event_type: "purchase | download"
    item_id: "string"
    category: "string"
    value: "number"
    timestamp: "string (ISO 8601)"
    metadata: "object"

consumer_groups:
  flink_aggregator:
    description: "Maintains dimension-keyed Top-K min-heaps and Count-Min Sketches per leaderboard, flushing to Redis every 60 seconds"
  event_archiver:
    description: "Performs batch inserts into Cassandra ranking_events partitioned by date, hour, and deterministic bucket"
  fraud_detector:
    description: "Evaluates velocity and per-IP thresholds before events count toward rankings"

operational_paths:
  sync_path: "Validate payload, publish to ranking-events topic, and return 202 Accepted immediately"
  async_path: "Flink stream aggregation, Redis board update, and Spark hourly reconciliation"
  dead_letter_queue: "ranking-events-dlq with retry limit of 3, alerting when consumer lag exceeds 5 minutes"

API Design

Sketch the primary read endpoints first, focusing on top-K leaderboard retrieval and individual item rank lookups because interviewers frequently prioritize these contracts.

Ranking API Domain Signatures

TypeScript interfaces defining operational signatures and parameters across services:

TYPESCRIPT
interface RankingEntry {
  rank: number;
  itemId: string;
  name: string;
  score: number;
  change: string;
}

interface TopKRequest {
  category?: string; // defaults to "overall"
  geography?: string; // defaults to "global"
  period: "hourly" | "daily" | "weekly" | "monthly" | "all-time";
  metric?: "downloads" | "revenue" | "ratings" | "active_users"; // defaults to "downloads"
  limit?: number; // defaults to 50, maximum 100
}

interface TopKResponse {
  category: string;
  geography: string;
  period: string;
  metric: string;
  asOf: string; // ISO 8601 timestamp
  rankings: RankingEntry[];
}

interface ItemRankRequest {
  itemId: string;
  category?: string; // defaults to "overall"
  geography?: string; // defaults to "global"
  period?: "hourly" | "daily" | "weekly" | "monthly" | "all-time";
  metric?: "downloads" | "revenue" | "ratings" | "active_users"; // defaults to "downloads"
}

interface ItemRankResponse {
  itemId: string;
  rank: number; // Unambiguous rank in the requested leaderboard (category + geography + period + metric)
  category?: string;
  geography?: string;
  period?: string;
  metric?: string;
  ranks?: {
    overall?: number;
    category?: number;
    daily?: number;
    weekly?: number;
    [key: string]: number | undefined;
  };
}

interface HistoricalRankingRequest {
  category?: string; // defaults to "overall"
  geography?: string; // defaults to "global"
  period?: "daily" | "weekly" | "monthly";
  metric?: "downloads" | "revenue" | "ratings" | "active_users"; // defaults to "downloads"
  date: string; // Format: YYYY-MM-DD
  limit?: number;
}

// Client and internal service contract for leaderboard queries
interface RankingService {
  getTopK(request: TopKRequest): Promise<TopKResponse>;
  getItemRank(request: ItemRankRequest): Promise<ItemRankResponse>;
  getHistoricalRankings(request: HistoricalRankingRequest): Promise<TopKResponse>;
}

Get Top K Leaderboard

HTTP
GET /api/v1/rankings?category=games&geography=US&period=daily&metric=downloads&limit=50 HTTP/1.1
Host: api.example.com
Accept: application/json

HTTP/1.1 200 OK
Content-Type: application/json

{
  "category": "games",
  "geography": "US",
  "period": "daily",
  "metric": "downloads",
  "as_of": "2026-03-13T10:00:00Z",
  "rankings": [
    {"rank": 1, "item_id": "app123", "name": "Puzzle Master", "score": 152000, "change": "+2"},
    {"rank": 2, "item_id": "app456", "name": "Word Rush", "score": 148500, "change": "-1"}
  ]
}

Get Item Rank

HTTP
GET /api/v1/rankings/item/app123?category=games&geography=US&period=weekly&metric=downloads HTTP/1.1
Host: api.example.com
Accept: application/json

HTTP/1.1 200 OK
Content-Type: application/json

{
  "item_id": "app123",
  "category": "games",
  "geography": "US",
  "period": "weekly",
  "metric": "downloads",
  "rank": 3,
  "ranks": {
    "overall": 15,
    "category_games": 3,
    "daily": 1,
    "weekly": 5
  }
}

Get Historical Rankings

HTTP
GET /api/v1/rankings/history?category=games&geography=US&period=daily&metric=downloads&date=2026-03-01&limit=10 HTTP/1.1
Host: api.example.com
Accept: application/json

HTTP/1.1 200 OK
Content-Type: application/json

{
  "category": "games",
  "geography": "US",
  "period": "daily",
  "metric": "downloads",
  "snapshot_date": "2026-03-01",
  "rankings": [
    {"rank": 1, "item_id": "app123", "name": "Puzzle Master", "score": 142000},
    {"rank": 2, "item_id": "app789", "name": "Speed Racer", "score": 139800}
  ]
}

Common 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

Your data model separates active leaderboards in Redis from cold historical records in Cassandra so live reads never query the raw event log.

Kafka Topic: ranking-events

JSON
{
  "event_id": "uuid",
  "event_type": "purchase",
  "item_id": "app123",
  "category": "games",
  "timestamp": "2026-03-13T10:00:00Z",
  "value": 1,
  "metadata": {"country": "US", "price": 2.99}
}

Redis: Ranked Lists

YAML
# Global daily top ranking leaderboard across all categories by downloads
key: "ranking:overall:global:daily:downloads"
type: "Sorted Set (ZSET)"
member: "item_id"
score: "event_count (or revenue / weighted score)"

# Category-specific regional weekly leaderboard by revenue
key: "ranking:games:US:weekly:revenue"
type: "Sorted Set (ZSET)"
member: "item_id"
score: "event_value"

Cassandra: Historical Rankings

Historical rankings are partitioned by composite key (category, geography, period, metric) to cleanly isolate historical snapshots across every supported multi-dimensional leaderboard. Within each partition, rows cluster by snapshot_date DESC, rank ASC to provide immediate point-in-time leaderboard retrieval.

SQL
CREATE TABLE historical_rankings (
    category      TEXT,
    geography     TEXT,          -- 'global', 'US', 'UK', etc.
    period        TEXT,          -- 'daily', 'weekly', 'monthly'
    metric        TEXT,          -- 'downloads', 'revenue', etc.
    snapshot_date DATE,
    rank          INT,
    item_id       TEXT,
    score         BIGINT,
    PRIMARY KEY ((category, geography, period, metric), snapshot_date, rank)
) WITH CLUSTERING ORDER BY (snapshot_date DESC, rank ASC);

Cassandra: Raw Events (for batch reconciliation)

At ~1B events per day (~41.7M events per hour), composite partitioning by (event_date, event_hour, bucket) splits hourly events across 64 deterministic buckets derived via abs(hash(event_id)) % 64. This limits each partition to ~650K rows (~65 MB to 90 MB), maintaining partition sizes safely below Cassandra's 100 MB operational ceiling. Spark reconciliation reads these bucketed partitions concurrently in parallel rather than scanning an unmanageable monolithic hourly partition. Cloud object storage (such as S3 with Parquet formatting) also serves as an alternate or complementary immutable event store for long-term batch processing.

SQL
CREATE TABLE ranking_events (
    event_date  DATE,
    event_hour  INT,
    bucket      INT,           -- Deterministic bucket: abs(hash(event_id)) % 64 to keep partition sizes under 100 MB (~650K rows)
    event_id    UUID,
    item_id     TEXT,
    event_type  TEXT,
    category    TEXT,
    value       INT,
    timestamp   TIMESTAMP,
    PRIMARY KEY ((event_date, event_hour, bucket), event_id)
);

Fault Tolerance

Flink worker failover, Count-Min Sketch drift, and fraud-inflated counts require explicit architectural protections.

ConcernSolution
Flink failureCheckpoint state to S3 every 30 seconds and restart workers from the latest verified checkpoint
Redis data lossBatch pipeline can reconstruct rankings from raw events
Count-Min Sketch driftHourly reconciliation via Spark batch job
Event lossKafka uses RF=3 and min.insync.replicas=2 to preserve acknowledged events across tolerated broker failures, while the immutable event log supports replay and downstream recovery
Stale rankingsServe rankings directly from Redis with staleness bounded by the window interval, which remains fully acceptable for clients

Specific: Handling Rank Manipulation and Fraud

Bad actors frequently attempt to manipulate top charts through automated bot farms and simulated download activity.

  • Only count verified financial transactions and completed downloads, discarding refunded or cancelled events
  • Employ device fingerprinting to detect coordinated bot farm installations
  • Apply velocity anomaly checks to flag sudden download spikes for review
  • Filter out events originating from flagged fraudulent accounts before updating ranking heaps

Additional Considerations

Multi-dimensional leaderboards, time-decay scoring, and fraud detection represent advanced system design evaluation topics.

Lambda Architecture (Real-time and Batch)

YAML
speed_layer:
  engine: "Apache Flink"
  latency: "Real-time (~5 seconds)"
  precision: "Approximate (Count-Min Sketch + Min-Heap)"
  sink: "Redis Sorted Sets"

batch_layer:
  engine: "Apache Spark"
  schedule: "Hourly batch reconciliation"
  precision: "Exact raw event aggregation"
  storage: "Cassandra (cold snapshots) and Redis overwrite"

serving_layer:
  store: "Redis"
  access: "Pre-computed Sorted Sets queried via ZREVRANGE in under 5ms"

Multi-Dimensional Rankings

Note on leaderboard scale scenarios: 1000 lists is the baseline capacity assumption used for initial sizing; 12K lists is a larger staff-level multi-dimensional scenario (1000 categories x 4 time periods x 3 metrics); and 20K lists represents an expanded deployment across all orthogonal dimensions (50 categories x 20 geographies x 5 periods x 4 metrics).

YAML
dimensions:
  category: ["games", "productivity", "education", "utilities"]
  geography: ["US", "UK", "IN", "DE", "global"]
  time_period: ["hourly", "daily", "weekly", "monthly", "all-time"]
  metric: ["downloads", "revenue", "ratings", "active_users"]

calculation:
  formula: "total_lists = categories * geographies * periods * metrics"
  example: "50 categories * 20 geographies * 5 periods * 4 metrics = 20,000 ranking lists"
  storage_strategy: "Each ranking list is maintained independently as a Redis Sorted Set"

Heavy Hitters Problem

The Top-K problem connects to the classic Heavy Hitters problem in streaming algorithms:

  • Misra-Gries Algorithm: Maintains at most K-1 candidates using O(K) space
  • Space-Saving Algorithm: Keeps top K counters and replaces the minimum counter when a new item arrives
  • Lossy Counting: Maintains approximate counts with a mathematically guaranteed error bound

For system design interviews, pairing a Count-Min Sketch with a Min-Heap is the most practical and widely recognized pattern.

Related Problems and Concepts

Windowed aggregation patterns connect directly to Distributed Stream Processing for Flink state management and Trending Topics for velocity scoring models. Probabilistic sketch structures and hash matrix foundations are detailed in Bloom Filter.

Min-Heap Throughput Example at 60K Events/Sec

At an ingestion rate of 60K events per second with K = 100, each partition worker maintains a min-heap of size 100 over Count-Min Sketch frequency estimates across its assigned leaderboard dimension keys. In this illustrative configuration, 16 Flink subtasks consume the 64 Kafka partitions in parallel, averaging approximately 4 Kafka partitions per subtask. Kafka partition count and Flink operator parallelism are decoupled: Kafka partitions provide storage-level parallel ingress, while Flink task parallelism can scale independently based on compute load.

For every incoming event, incrementing the sketch consumes approximately 200 nanoseconds, while comparing the resulting estimate against the corresponding leaderboard heap root requires roughly 50 nanoseconds. Across the 16 Flink subtasks, this represents approximately 750K operations per subtask each second, totalling around 12M operations per second across the cluster. A periodic full heap rebuild occurs every 10 seconds from a sketch snapshot to catch items that entered the top-K thresholds late in the window. Storing 100 entries at 64 bytes each requires only 6.4 KB of memory per leaderboard window per partition, which is negligible compared to retaining counters for all 5M items.

Interview Walkthrough

  • 25-minute cut

    Focus on core streaming architecture unless the interviewer prompts for staff depth.

    • Requirements and scope clarification between exact and approximate counting (3 min)
    • Count-Min Sketch combined with min-heap aggregation (7 min)
    • Kafka ingestion and Flink stream windowing (7 min)
    • Redis Sorted Set serving path (5 min)
    • Hourly batch reconciliation sketch (3 min)
  • Frame the scenario as a streaming problem early, explaining that storing all item counts in memory is infeasible at scale and bounded memory requires approximate algorithms.
  • Present the Count-Min Sketch for constant-time frequency estimation, paired with dimension-keyed min-heaps of size K to track candidates per leaderboard without cross-board interference.
  • Compare exact hash map aggregation against approximate counting, explaining that exact counts work for small catalog cardinalities whereas sketches scale seamlessly to billions of events.
  • Discuss the end-to-end update flow, where events arrive via Kafka, stream processors increment sketch counters, and periodic flushes refresh the Redis leaderboard cache for API reads.
  • Specify the mathematical error tolerance, emphasizing that Count-Min Sketch overcounts but never undercounts, which suits trending discovery leaderboards rather than billing systems.
  • Clarify the dual-path accuracy model: real-time streaming rankings rely on approximate sketches for speed, whereas authoritative historical leaderboards are reconciled via exact batch jobs over partitioned event logs.
  • Apply Back-of-the-Envelope Estimation to show that 1M events per second with 8-byte counters across hash functions requires memory in megabytes rather than terabytes.
  • Address the common pitfall of attempting periodic full batch sorts, noting that sorting billions of raw events cannot complete before subsequent streaming windows close.

Engineering Trade-offs

Your interviewer will evaluate your trade-off analysis between exact and approximate counting and stream window selection. Walk through each trade-off and state your architectural choice for an App Store-style leaderboard.

Exact vs Approximate Counting: When Approximate Is Good Enough

Exact counting (for small-to-medium scale):
  Each item maintains an explicit counter in a hash map
  Space complexity: O(N) where N represents unique item cardinality
  At 10M unique items * 8 bytes = 80 MB (viable for in-memory state)
  At 1B unique items * 8 bytes = 8 GB, making in-memory aggregation prohibitive
  
  Recommendation: Use when unique items N < 10M per window and exact counts are required
  Production example: App Store top 100 apps (approximately 5M total catalog items)

Approximate counting (Count-Min Sketch for massive scale):
  Space complexity: O(w * d), remaining constant regardless of item count N
  Typical parameters: w = 10,000, d = 5 yielding 50,000 counters * 8 bytes = 400 KB
  Error bound: epsilon = e / w = 2.71 / 10,000 = 0.027% maximum overcount
  
  Recommendation: Use when unique items N > 10M and exact memory allocation is prohibitive
  Production example: Trending social hashtags with billions of distinct terms

Architectural insight: Top-K ranking is significantly more tolerant of approximation than exact counting:
  Suppose true counts are: Item A = 1000, Item B = 999, Item C = 998
  Count-Min Sketch might estimate: Item A = 1001, Item B = 1000, Item C = 999 (uniform slight overcount)
  The relative ordering remains preserved, so the resulting Top-K set remains accurate
  
  Boundary condition: If two items share nearly identical counts (such as A = 1000 and B = 999):
    The sketch may swap adjacent ranks, which is acceptable for trending discovery
    This approximation is unacceptable for official financial accounting or competitive prize payouts

Flink Tumbling vs Sliding vs Session Windows

For ranking systems, three primary stream window types apply:

Tumbling Windows (non-overlapping fixed intervals):
  Intervals: [0-5 min], [5-10 min], [10-15 min]
  Advantages: Computationally simple and efficient because each event is processed exactly once
  Disadvantages: Rankings exhibit boundary jumps rather than continuous trends
    At minute 59, a 1-hour tumbling window holds only 5 minutes of data for the next window
    At minute 60, the window resets abruptly, causing sharp rank transitions
  Best for: Formal batch snapshot intervals such as official hourly charts

Sliding Windows (overlapping continuous intervals):
  Configuration: Window size = 5 minutes, slide interval = 1 minute
  Intervals: [0-5 min], [1-6 min], [2-7 min], [3-8 min]
  Advantages: Produces smooth, continuous updates that reflect immediate momentum
  Disadvantages: Increases processing overhead because each event contributes to multiple active windows
  Resource cost: Requires maintaining state across all active sliding slices in memory
  Best for: Live "Trending Now" leaderboards updating every minute

Session Windows (inactivity-based intervals):
  Configuration: Groups activity bursts separated by inactivity gaps greater than 30 minutes
  Advantages: Natural grouping for individual user navigation sessions
  Disadvantages: Unpredictable duration and cardinality make it unsuitable for global top-K lists
  Best for: User behavioral analytics rather than global rankings

Architectural recommendation for rankings:
  Real-time client display: Sliding window (5-minute window, 1-minute slide) to show continuous momentum
  Official hourly leaderboards: Tumbling window to provide stable, non-overlapping historical records
  Viral trend detection: Sliding window combined with exponential time decay

Reconciling Real-Time vs Batch Rankings

Lambda Architecture for Top-K Ranking Systems:

Speed Layer (Apache Flink with approximately 5-second latency):
  Operational strengths: Delivers near-real-time updates reflecting the latest 5 minutes of activity
  Operational trade-offs:
    Count-Min Sketch introduces bounded overcounting and can occasionally swap items near rank K
    Flink operator state resides in memory and relies on checkpoint recovery during node restarts
    Late-arriving events exceeding the watermark threshold are dropped, leading to slight undercounts

Batch Layer (Apache Spark running on hourly intervals):
  Operational strengths: Computes exact counts across all historical events, including late arrivals
  Operational trade-offs: Data reflects ground truth but is stale by up to 1 hour

Serving Layer (Redis Sorted Sets):
  Unified view: Serves low-latency reads directly from pre-computed sorted sets
  Correction model: Hourly batch jobs overwrite approximate values with exact reconciled counts

Why real-time cannot run in batch-only mode:
  Users expect immediate feedback for viral events, whereas hourly updates feel static and unresponsive

Why streaming alone is insufficient without batch reconciliation:
  Without an authoritative batch reconciliation job, streaming inaccuracies and dropped late events accumulate over days

Practical production pattern:
  Redis serves user reads with under 5-millisecond latency from the speed layer
  Every hour, the Spark batch job computes exact figures and overwrites the Redis keys
  Maximum inaccuracy is kept below 1% during active windows and reset hourly

The Count Inflation Problem in Ranking Systems

Challenge: Fraudulent actors artificially inflate event counts to boost product rankings.

Detection and Mitigation Strategies:

1. Velocity Anomaly Detection:
   Baseline: An application normally averages 10K downloads per day across 30 days
   Anomaly: The same application receives 500K downloads in 1 hour (a 50x spike)
   Action: Flag the application for review and apply a fraud discount factor to its ranking score

2. Device and Account Deduplication:
   Rule: Multiple downloads from the same physical device count only once per 24-hour window
   Implementation: SETEX dedup:{device_id}:{item_id} 86400 "1" in Redis
   Execution: Increment ranking counters only when the deduplication key is newly created (SETNX)

3. Source Velocity Throttling:
   Enforce a maximum of 1,000 downloads per IP address per hour
   Enforce a maximum of 10 downloads per user account per day for standard consumers
   Implementation: INCR downloads:{ip}:{hour} with automated TTL expiration in Redis

4. Model-Based Post-Hoc Weighting:
   Machine learning pipelines generate a fraud_score between 0 and 1 for each event
   Events with suspicious velocity from newly created accounts receive up to a 70% count discount
   Implementation: Flink computes fraud_score per event and applies a weighted count increment:
     actual_count += (1.0 - fraud_score) * event_value
   Result: Raw fraudulent events are logged, but ranking scores reflect clean, verified demand

💬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...