Skip to content
Reliable Data Engineering
Practice problem medium streamingtop-kcount-min-sketchsliding-windowheavy-hitters
Practise with timer, notes and rubric

Design Real-Time Top-K Trending

Problem

Show the top 100 trending songs (or hashtags/products) in the last 1 hour, last 24 hours, and last 7 days, globally and per country. Updated every minute. Events: plays (or posts/views).


Clarifying questions

QuestionAssumed answer
Volume?1 B play events/day, 3× peak; 50 M distinct songs, 200 countries
Freshness?Lists refresh every 1 min
Exact or approximate?Approximate ranking acceptable; counts within ~1%
“Trending” = most played or fastest rising?Start with most played; discuss velocity as extension
Anti-gaming?Exclude plays < 30 s, bots, repeated plays by the same user beyond N/hour

1. Estimates

1B/day → 12k/s avg → ~35k/s peak; 200 B/event → 7 MB/s: modest
Naive exact state: 50M songs × 200 countries is the theoretical bound; actual (song, country) pairs active per hour ~20M
Per-minute partial counts: (song, country, minute) rows ≈ 20M/hour → fine for a streaming job
Output: 3 windows × 201 scopes × 100 rows ≈ 60k rows per refresh → tiny; serve from Redis/KV

2. Architecture

flowchart LR
    APP[Clients] --> K[[Kafka: plays<br/>key = song_id]]
    K --> F1[Stage 1: filter valid plays,<br/>per-minute counts per song × country]
    F1 --> MIN[(minute_counts<br/>Delta / state)]
    F1 --> F2[Stage 2: rolling window sums<br/>1h / 24h / 7d per scope]
    F2 --> TOP[Stage 3: top-K per scope<br/>min-heap of size K]
    TOP --> R[("Redis sorted sets<br/>trending:window:country")]
    R --> API[Trending API] --> UI[Apps]
    MIN --> BATCH[Hourly/daily batch: exact recompute<br/>for 24h / 7d, anti-fraud]
    BATCH --> R

3. Deep dives

3.1 Sliding windows without recomputing everything

Store per-minute buckets (a “pane” approach). A 1-hour window = sum of the last 60 buckets. Each minute:

count_1h(song) += bucket[now](song) − bucket[now − 60 min](song)

Only songs that changed are touched. For 24 h and 7 d, use hourly buckets (24 and 168 of them): coarser granularity is fine for longer windows. This is how you avoid a 7-day sliding window with a 1-minute slide (10,080 overlapping windows per event!).

3.2 Top-K maintenance

3.3 Approximate counting at very large cardinality

If the key space is enormous (all hashtags, all URLs), use Count-Min Sketch + heap (heavy hitters):

3.4 Anti-gaming

Filter in stage 1: plays ≥ 30 s, cap counted plays per (user, song, hour), drop known bots, weight by account age. These need per-user state (TTL 1 h). Batch recompute applies heavier fraud logic and overwrites the 24h/7d lists.

3.5 Serving

Redis sorted sets: ZADD trending:1h:DE score song_id, written as a full replacement per refresh (write to a new key, RENAME atomically) so readers never see half-updated lists. Cache in CDN for 60 s: every user sees the same list per country.

Popular = highest count. Trending = acceleration: e.g. score = count_last_1h / (avg hourly count last 7d + smoothing), or a z-score vs a baseline. Requires both short and long windows; already available from the bucket design. Add a minimum volume threshold to avoid “3 plays vs 0” spikes.

5. Trade-offs

DecisionChoiceAlternative
WindowingPer-minute/hour buckets + incremental sumsNative sliding windows (huge overlap, state explosion for 7 d)
CountingExact keyed counts (key space manageable)CMS for unbounded key spaces
ServingRedis + CDNOLAP query at request time (unnecessary: results are tiny and shared)
CorrectionBatch recompute for long windowsStreaming only (fraud filtering weaker)

6. What separates a senior answer

7. Follow-up questions

How do you handle a viral song creating a hot key?

Stage 1 pre-aggregates per partition (local combiner) before keying by song, so one song’s 10k plays/s become a few partial counts per task per minute. Salting with a second merge stage if needed.

How would you personalise trending (trending among people like me)?

Compute trending per segment (country × age band × top genre) as extra scopes, still cheap since output per scope is K rows; blend with personal recommendations at serving time. Cap the number of segments to control state.


Self-assessment rubric