Limited time: AI code review, hints, mock interviews, whiteboard analysis, and all Pro features are unlocked. Enroll
⏱️ 58 min read

Designing Trending Topics

Difficulty: Advanced Topics: Count-Min Sketch, Heavy Hitters, Stream Processing, Sliding Windows Asked at: Twitter, Google, Meta, Amazon Prerequisites:Message Queues, Batch vs Stream, and Bloom Filters


1. Understanding the Problem

Trending topics is the short list in the sidebar: the 50 terms that are spiking right now, scoped to where you are. Behind it is a firehose of events - tweets, searches, hashtags - and the job is to pick the 50 most interesting terms out of millions, refresh that pick every few seconds, and do it per region and per time window.

The thing to say out loud early, because it changes every decision that follows: this is not a leaderboard. A leaderboard tracks exact scores for a bounded, known set of players. You know every player ID up front, there are maybe 10 million of them, each has one authoritative number, and getting a rank wrong by one position is a bug. Trending is top-K over an unbounded, high-cardinality stream of strings you have never seen before. New terms appear continuously, 99.9% of them will never matter, and nobody can tell whether the term at rank 23 was counted as 41,200 or 41,350. One problem wants a Redis sorted set and exact arithmetic. The other wants approximate counting, and approximation is the correct answer rather than a shortcut.

Real examples: X/Twitter Trends, Google Trends and the Trending Searches carousel, YouTube Trending, Reddit’s r/popular, Spotify Viral 50.


2. Naive First Cut

flowchart LR
    APP["Client<br/>tweets and searches"]:::client
    API["Ingest API"]:::service
    DB[("Counters DB<br/>term and count")]:::data
    UI["Trends UI"]:::client

    APP -->|"1. Send event"| API
    API -->|"2. Increment per term"| DB
    UI -->|"3. ORDER BY count DESC LIMIT 50"| DB

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0
Color Meaning
🟣 Purple Clients
🟒 Green Services
🟑 Yellow Data stores
πŸ”΅ Blue Edge / CDN

One row per term with a counter. Increment on every occurrence, sort the table when someone opens the page.

Why this breaks:

The rest of the doc evolves this into a streaming pipeline with bounded-memory sketches, two-tier top-K aggregation, and a ranking function that measures rate of change instead of volume.


3. Prior Art We’re Drawing From


4. Functional Requirements

Core (Top 3)

  1. Return the top K trending terms for a time window - K = 50, ordered, with a score you can explain
  2. Support several windows - 5 minutes, 1 hour, 24 hours, all live at the same time
  3. Scope trends by region - a user in Mumbai and a user in Berlin see different lists

Below the Line


5. Non-Functional Requirements

NFR Target
Ingest throughput 500K events/sec at peak, no backpressure onto the client path
Freshness Trend list recomputed every 10-30 seconds
Read latency P99 < 50ms for the trend list, served almost entirely from cache
Read throughput 100K reads/sec, but the answer is identical per region so it caches perfectly
Accuracy Approximate counts acceptable; a genuine spike must never be silently dropped

Below the Line


6. Scale Estimation (Back-of-Envelope)

The shape of the problem: enormous write amplification into a tiny read answer, with all the interesting engineering on the write side.


7. Core Entities


8. API / System Interface

GET /v1/trends?region=IN&window=5m&limit=50
  Response: {
    region: "IN",
    window: "5m",
    generatedAt: 1760000012000,
    trends: [
      { rank: 1, term: "#worldcupfinal", score: 41.8, count: 182400,
        distinctUsers: 96300, deltaRank: 4 },
      ...
    ]
  }

GET /v1/trends/regions
  Response: [{ code: "IN", name: "India" }, { code: "DE", name: "Germany" }, ...]

POST /internal/events          (server-to-server only, from the ingest edge)
  Body: { eventId, tsEvent, userId, region, terms: [...], lang, source }
  Response: 202 Accepted

Security note: The read endpoint is public and anonymous, which is why it caches so well. The ingest endpoint must never be client-callable with an arbitrary userId or region - both are derived server-side from the authenticated session and the request IP. A client that can set its own userId can forge 100,000 distinct users and manufacture a trend for free.


9. High-Level Design

Three requirements, three passes. Each one adds the smallest set of components that makes it work.

FR1: Count Terms in a Window

Start where the naive cut did and do the arithmetic properly, because the arithmetic is the argument.

500K events/sec, ~5 terms each, so 2.5M increments/sec against a key space of 10M terms. A single Redis node handles roughly 100-200K simple commands/sec, so even the fast path needs 15-25 nodes doing nothing but INCR - before any read load, before windows, before regions. A relational table with row-level updates is not in the conversation. And the brutal part: the overwhelming majority of those increments land on terms that will never be read by anybody, because a top-50 list only ever surfaces a few thousand candidates.

Two moves fix the shape. First, decouple the client write from the counting work with an append-only log, so a spike in events never becomes a spike in database contention. Second, do the counting inside a stateful stream processor rather than in a shared store, so an increment is a local memory write instead of a network round trip.

New components:

  1. Ingest API - terminates the client request, authenticates it, derives region from IP or account setting, tokenises the content into terms, and appends one record to the log. Returns 202 immediately.
    πŸ’‘ Tokenising at ingest rather than downstream means the heavy per-event string work scales with your stateless edge tier, which is the cheapest tier to add machines to.
  2. Ingestion log (Kafka) - an append-only, partitioned, replayable buffer. Absorbs the 500K/sec burst, gives the processor at-least-once delivery, and lets you re-run the pipeline from 24 hours ago when you change the ranking function.
  3. Stream processor (Flink) - reads the log, maintains per-bucket counts in local keyed state, and emits results on a timer. State lives on the task manager’s own disk via RocksDB, so a term increment never leaves the process.
flowchart LR
    APP["Client"]:::client
    API["Ingest API<br/>tokenise and tag region"]:::service
    KAFKA["Kafka<br/>term events"]:::async
    PROC["Stream Processor<br/>per bucket counts"]:::service
    STATE[("Processor State<br/>RocksDB backed")]:::data

    APP -->|"1. Post event"| API
    API -->|"2. Append record"| KAFKA
    KAFKA -->|"3. Consume partition"| PROC
    PROC -->|"4. Read modify write counts"| STATE

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. A user posts β€œEngland just won the #worldcupfinal” - the Ingest API authenticates, stamps region=IN, tsEvent
  2. API tokenises into terms: #worldcupfinal, england, won, world cup final
  3. API appends one record with all four terms to the term-events Kafka topic, partitioned round-robin
  4. Flink consumes the partition, splits the record into one update per term
  5. For each term, Flink updates the count for (term, region, current 1-minute bucket) in local RocksDB state
  6. Nothing is sorted, nothing is queried, nothing leaves the machine. One event costs a handful of local writes.

Why round-robin partitioning and not partition-by-term? Keying by term would make each term’s count exact on exactly one subtask, which is tempting. It also means one viral hashtag at 50K events/sec pins a single subtask while the other 63 idle. Round-robin keeps ingest flat and pushes the aggregation problem to FR2, where it is cheaper to solve.


FR2: Get the Top K

The counts exist. Now produce an ordered 50 every 10 seconds.

The obvious move is to sort the counters and take the head. That is the wrong shape by a wide margin: you are asking for 50 items and paying to order 10 million. Sorting is O(n log n) in the size of the key space when the answer size is a constant.

Use a bounded min-heap of size K instead. Walk the counters once, keep a 50-element heap ordered smallest-first, and for each term compare against the heap root. If the term beats the root, pop and push; otherwise skip. That is O(n log K) with K = 50, so log K β‰ˆ 6 instead of log n β‰ˆ 23, and - the part that actually matters - memory is 50 entries rather than a full sorted copy of the key space.

One heap per aggregator still leaves a problem: with round-robin partitioning, every aggregator sees a slice of every term, so no single aggregator’s heap is the global answer. This is the two-tier aggregation pattern.

New components:

  1. Partial Aggregators - the parallel tier, one per Flink subtask. Each maintains counts and a local top-K heap over the slice of the stream it sees, and emits its local candidates on a 10-second timer.
  2. Top-K Merger - a single low-parallelism operator. Receives candidate lists from every partial aggregator, sums the per-term counts across them, and runs one final heap to produce the global top 50.
flowchart LR
    KAFKA["Kafka<br/>term events"]:::async
    A1["Partial Aggregator 1<br/>local top-K heap"]:::service
    A2["Partial Aggregator 2<br/>local top-K heap"]:::service
    A3["Partial Aggregator N<br/>local top-K heap"]:::service
    MERGE["Top-K Merger<br/>sum and re-rank"]:::service

    KAFKA -->|"1. Fan out partitions"| A1
    KAFKA -->|"2. Fan out partitions"| A2
    KAFKA -->|"3. Fan out partitions"| A3
    A1 -->|"4. Emit local candidates"| MERGE
    A2 -->|"5. Emit local candidates"| MERGE
    A3 -->|"6. Emit local candidates"| MERGE

    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0

Step-by-step flow:

  1. 64 partial aggregators each consume their Kafka partitions and count terms locally
  2. Every 10 seconds, a processing-time timer fires on each aggregator
  3. Each aggregator walks its counters through a bounded heap and emits its top candidates
  4. The merger receives 64 candidate lists and sums counts per term across all of them
  5. One final heap runs over the union - a few thousand terms, not 10 million
  6. That merged top 50 is handed to the ranking stage

The error this introduces, stated plainly. Each aggregator only reports its local top-K, so a term is invisible to the merger unless at least one aggregator ranked it. Consider a term that is genuinely 10th globally but, because its traffic is spread evenly across 64 partitions, sits at position 51 on every single one. It appears in zero candidate lists. The merger never sees it. It is missing from the trend list entirely, and no amount of correct merging recovers it.

The standard mitigation: have each aggregator keep and emit a local top-(K Γ— f) rather than top-K, with f around 10 to 20. At K = 50 and f = 20, each aggregator emits 1,000 candidates and the merger handles 64,000 - still trivial. For a term to be lost now it must rank below 1,000 on every partition while ranking in the global top 50, which requires its traffic to be both high and almost perfectly uniform across partitions. A genuinely spiking term is spiking everywhere, so it clears position 1,000 on essentially every partition. The cost of f is linear in network and merger CPU, both of which are nowhere near saturated, so buy a large f and stop worrying.


FR3: Multiple Windows and Regions

Three window lengths and twenty regions means sixty lists. The naive reading is sixty pipelines. It should be one pipeline that writes sixty small values.

The trick is that windows do not need independent counting. Count once into fixed buckets, then build each window by summing the buckets it covers. One-minute buckets give you the 5-minute window as a sum of 5 and the 1-hour window as a sum of 60. The 24-hour window uses coarser 1-hour buckets, because nobody needs the daily list to slide in 60-second steps.

Region is just part of the key. Count (term, region, bucket) rather than (term, bucket) and every region’s numbers are maintained in the same pass.

New components:

  1. Serving Store (Redis) - holds the materialised trend list per region and window. The merger overwrites the whole value on each refresh. Total data: 3,000 entries.
  2. Read Cache / CDN - fronts the read API with a 10-second TTL. Since every user in a region gets a byte-identical response, the hit rate is effectively 100% after the first request per region per refresh.
  3. Trends API - a thin read service. One Redis GET, serialise, set cache headers. No computation.
flowchart LR
    MERGE["Top-K Merger"]:::service
    REDIS[("Redis<br/>trend lists per region")]:::data
    TAPI["Trends API"]:::service
    CDN["CDN<br/>10s TTL"]:::edge
    READER["Reader"]:::client

    MERGE -->|"1. Overwrite list per region and window"| REDIS
    TAPI -->|"2. Read list on cache miss"| REDIS
    CDN -->|"3. Forward on miss"| TAPI
    READER -->|"4. Get trends for region"| CDN

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Step-by-step flow:

  1. Merger finishes a refresh cycle and holds 60 ranked lists, one per region and window pair
  2. Merger writes each as a single Redis key: SET trends:IN:5m <json> with a 120-second TTL
  3. Reader opens the app, client calls GET /v1/trends?region=IN&window=5m
  4. Request lands on a CDN edge. Cache hit for the current 10-second slice, served in under 20ms
  5. On miss, the edge forwards to the Trends API, which does one Redis GET and returns with Cache-Control: public, max-age=10
  6. The next 99,999 readers in that region hit the edge cache

How sliding windows come out of fixed buckets. Keep the last 60 one-minute buckets in a ring buffer indexed by minute mod 60. The 1-hour window is the sum of all 60. Advancing one minute means adding the new bucket and subtracting the one it overwrote - two operations, not 60. The honest limitation: the window slides in whole-bucket steps, so with one-minute buckets your β€œlast hour” is really β€œthe last 60 completed minutes” and can be up to a minute stale at the trailing edge. For a list refreshed every 10 seconds that is fine, and Deep Dive 2 shows the version that removes the cliff entirely.


10. Technology Choices

Tier What it stores Access pattern Primary pick Alternatives
Ingestion log Raw term events, 24-48h retention Sequential append, replayable sequential read Kafka Kinesis, Pub-Sub, Pulsar
Stream processor Windowed aggregation and top-K logic Per-event local update, timer-driven emit Flink Kafka Streams, Spark Structured Streaming
Sketch state store Count-Min Sketch and Space-Saving slots per bucket Keyed read-modify-write at 2.5M ops/sec Flink embedded state, RocksDB backed In-heap state, external Redis
Serving store 60 materialised trend lists Whole-value overwrite, tiny keyed read Redis DynamoDB, Memcached
Baseline store Per-term mean and stddev per region and weekly slot Batch write nightly, range scan by term ClickHouse or Druid BigQuery, Snowflake, Cassandra
Read cache Serialised trend list responses Read-only, identical per region CDN with 10s TTL Redis replicas, Varnish

Why Flink over Kafka Streams or Spark Structured Streaming? This design lives and dies on event-time semantics. Flink treats watermarks, allowed lateness and side outputs for late events as first-class concepts, and its incremental RocksDB checkpointing is built for exactly the kind of large keyed state the sketches produce. Kafka Streams has event-time windows and RocksDB state too, and for a simpler topology it is the lighter operational choice - but its parallelism is bolted to topic partition count, and its late-data story is a single grace period rather than a routable side output. Spark Structured Streaming is micro-batch; it supports watermarks, but the batch boundary fights a 10-second refresh with sub-window granularity, and the two-tier aggregation becomes two shuffle stages per batch.

Why a separate baseline store and not the same Redis? The baseline is a different access pattern and a different lifetime - written once nightly by a batch job, read as a range scan keyed by term, and holding hundreds of millions of rows across the term-by-slot cross product. A columnar store answers β€œgive me the mean and stddev for these 5,000 terms at slot 412” in one scan. Redis would need 5,000 round trips or a Lua script, and would be holding cold data in RAM for no reason.

Why Approximate Counting Is the Right Answer, Not a Compromise

Worth being blunt about this, because candidates often present sketches apologetically, as the thing you settle for when you cannot afford the real answer.

The output of this system is a list of fifty strings shown to a human. If #worldcupfinal was seen 182,400 times and the sketch reports 182,650, nothing observable changes: the term occupies the same rank, the UI shows the same text, and the number is rounded to β€œ182K” before it reaches a screen. No downstream consumer multiplies this count by a price, reconciles it against a ledger, or pays out on it. The entire value of the count is ordinal. Exact counting, meanwhile, costs one or two orders of magnitude more - a hash entry per distinct term per region per live bucket, memory that grows with whatever the stream throws at you, and a recovery time that grows with it. Approximate means 1.3 MB per bucket, decided in advance.

Paying 100x for precision that cannot be perceived is not rigour. Picking the structure whose error budget matches the product’s tolerance is the engineering, and here the tolerance is enormous. The one thing you must preserve is direction of error: a sketch that could under-count would let a real spike vanish, and that is a visible product failure. Count-Min Sketch only ever over-estimates, which is why it is the right primitive and a plain hash with collisions is not.


11. Data Modeling

Kafka (term-events topic - the raw event):

{
  "eventId": "01JAV7K2M9QN3P", "tsEvent": 1760000000123, "tsIngest": 1760000000461,
  "userId": "u_8831204", "region": "IN", "lang": "en", "source": "tweet",
  "terms": ["#worldcupfinal", "england", "world cup final"],
  "accountAgeDays": 1820, "accountScore": 0.91
}

Partitioned round-robin on eventId so ingest stays flat under a viral hashtag. Retention 48 hours, which is the replay budget: long enough to re-run the pipeline after a ranking change, short enough to stay cheap. accountAgeDays and accountScore are stamped at ingest because the merger needs them for Deep Dive 4 and cannot afford a user-service lookup per event.

Flink keyed state (the per-bucket aggregate):

Key:   term | region | bucketWidth | bucketStart
Value: events         uint64   approximate occurrence count
       distinctUsers  bytes    HyperLogLog register array, 12 KB
       firstSeenTs    int64    used for cold-start and velocity
       contentHashes  bytes    sampled shingle sketch for diversity scoring

Keyed by term so increments are local. bucketWidth sits in the key rather than a separate state namespace, which keeps the 1-minute and 1-hour buckets in one structure and lets a single scan feed both windows.

Count-Min Sketch serialisation (per region and bucket):

Key: "cms:{region}:{bucketWidth}:{bucketStart}"

Header:  { depth d=5, width w=65536, seeds[5], totalWeight N, bucketStart }
Payload: d Γ— w Γ— uint32  =  5 Γ— 65536 Γ— 4  =  1,310,720 bytes

The seeds must be serialised with the payload. Restore a sketch with freshly generated hash seeds and every term maps to different columns - the counters are still there, they just no longer mean anything, and nothing in the system will tell you. Treat the seed set as part of the data. Space-Saving state alongside it is a bounded list, which is the entire point of it: ss:{region}:{bucketWidth}:{bucketStart} holding [(term, count, error)] Γ— m with m = 1000.

ClickHouse (baseline store - what β€œnormal” looks like for each term):

CREATE TABLE term_baseline (
    term          String,
    region        LowCardinality(String),
    slot_of_week  UInt16,     -- 0..671, one 15-minute slot per week
    mean_count    Float32,
    stddev_count  Float32,
    sample_weeks  UInt8,      -- how much history backs this row
    updated_at    DateTime
) ENGINE = ReplacingMergeTree
ORDER BY (region, term, slot_of_week);

slot_of_week rather than a plain hour-of-day because traffic has a weekly rhythm, not just a daily one. Sunday 20:00 and Tuesday 20:00 are different baselines for sports terms, and treating them as the same slot makes every weekend look like a trend. Only terms with real history get a row; everything else falls through to the smoothing prior in Deep Dive 3.

Redis (serving store - the materialised answer):

Key: "trends:{region}:{window}"        e.g. "trends:IN:5m"
Value: JSON array of 50 objects
  { rank, term, score, count, distinctUsers, deltaRank, firstSeenTs }
TTL: 120 seconds
Written by: the merger, whole-value SET per refresh cycle

The TTL is deliberate. If the pipeline dies, the list expires in two minutes and the API returns β€œtrends unavailable” rather than confidently serving a list from an hour ago. A stale trend list is worse than no trend list, because a user cannot tell it is stale.

Access Patterns:

Query Data Source How
Increment a term Flink local state Keyed read-modify-write into RocksDB, no network hop
Local top-K per aggregator Flink local state Single scan of the bucket’s Space-Saving slots through a size-1000 heap
Global top-K Merger memory Sum ~64K candidates by term, one size-50 heap
Fetch baselines for candidates ClickHouse One range scan WHERE region=? AND slot_of_week=? AND term IN (...)
Read trend list CDN then Redis Edge hit for the current 10s slice; on miss one GET trends:{region}:{window}
Rebuild after processor loss Kafka replay Reset consumer group to the start of the window, reprocess from the log

How a Single Term Travels the Whole Pipeline:

  1. #worldcupfinal is tokenised at the Ingest API and written into one Kafka record among five terms
  2. A partial aggregator increments it in the Space-Saving table for (IN, 1m, 1760000000) and adds the user to the term’s HyperLogLog
  3. The 10-second timer fires; the term is well inside the local top 1,000, so it is emitted as a candidate
  4. The merger sums the term’s count across 64 aggregators and joins the baseline row for (IN, slot 412)
  5. Its z-score against its own history is large, it clears the volume and distinct-user floors, and the spam filter passes it
  6. It is written into trends:IN:5m at rank 1 and reaches the next reader through the CDN within 10 seconds

12. Deep Dives

Deep Dive 1: Counting 10M Distinct Terms Without Unbounded Memory

Problem: You need a frequency estimate for any term in the current bucket. The term space is unbounded and adversarial - anybody can invent a new hashtag, and 500K events/sec of invented hashtags is a cheap attack. Whatever structure you pick, its memory has to be decided by you, not by the stream.

Bad: A hash map from term to count, one per region and bucket.

Do the arithmetic. A term string averages ~20 bytes. In a JVM, one map entry costs the String object (header, length, char array - call it 56 bytes), a boxed or primitive counter, and the entry node itself. Realistically 100-150 bytes per distinct term. A single global one-minute bucket holding 300K distinct terms is ~36 MB. Survivable.

Now multiply by what the design actually needs:

Dimension Factor Running total
One global 1-minute bucket, 300K terms 36 MB 36 MB
20 regions, each keeping its own map ~8x after dedup of the shared tail ~290 MB
60 live 1-minute buckets for the hourly window 60x ~17 GB
24 live 1-hour buckets for the daily window +24 coarse buckets ~25 GB

Twenty-five gigabytes of live keyed state, checkpointed every few seconds, with a recovery time proportional to its size. And the number is a function of the input: a hashtag-generator bot doubles it on command. That last property is the disqualifier, not the absolute size.

Good: Count-Min Sketch. Fixed memory, chosen up front, immune to cardinality.

The structure is a 2D array of counters, d rows by w columns, with d independent hash functions - one per row.

Why the minimum. Collisions only ever add. If another term happens to share a cell with yours, that cell is inflated by the other term’s count - so every row gives you at least your true count, never less. Taking the minimum picks the row with the least foreign contamination, which is the tightest over-estimate available. Count-Min Sketch over-estimates but never under-estimates.

The sizing relationship is clean: with w = ⌈e/Ξ΅βŒ‰ and d = ⌈ln(1/Ξ΄)βŒ‰, the estimate exceeds the true count by more than Ρ·N with probability at most Ξ΄, where N is total stream weight. In plain terms - width controls how much error, depth controls how often you exceed it.

Depth d Width w Memory (uint32) Over-estimate bound
4 2,048 32 KB ~0.13% of N, missed ~1.8% of the time
5 16,384 320 KB ~0.017% of N, missed ~0.7% of the time
5 65,536 1.3 MB ~0.004% of N, missed ~0.7% of the time

At 1.3 MB per region-bucket, the whole fan of 20 regions Γ— 60 buckets is ~1.5 GB - and it stays 1.5 GB whether 300K or 300M distinct terms arrive.

The limit of CMS alone, and it is a real one: the sketch answers β€œhow often have I seen this term,” but it cannot tell you which terms to ask about. It stores no keys. You still need a candidate set from somewhere, and if you maintain that candidate set as a full hash map you are back to the memory problem you just solved.

Great: Either pair the sketch with a bounded candidate heap, or use a structure built for this exact question.

Option A - CMS plus heap. Keep the sketch for frequencies and a size-m min-heap of candidate terms. On each event, increment the sketch, query the new estimate, and if it beats the heap root, insert the term and evict the root. The heap holds m term strings - a few thousand - so memory is the sketch plus a small bounded set. This works and is widely deployed.

Option B - Space-Saving (Metwally, Agrawal and El Abbadi). Designed for heavy hitters from the start, and it carries a stronger guarantee. Keep exactly m slots, each (term, count, error):

  1. If the term is already tracked, increment its count.
  2. If not and a slot is free, install it with count = 1, error = 0.
  3. If not and all slots are full, find the slot with the minimum count, evict that term, and install the new term with count = min + 1 and error = min.

Step 3 is the whole idea. The new term inherits the evicted term’s count as its own upper bound of uncertainty, recorded explicitly in error. True count lies in [count - error, count]. The guarantee that follows: any term whose true frequency exceeds N/m of the stream is guaranteed to still be in the table, where N is stream length. At m = 1000, anything above 0.1% of the stream cannot be missed - and a top-50 trend in a region is far above 0.1%. Implemented as a Stream-Summary (counters bucketed by count in a doubly-linked list), increment and find-minimum are both O(1).

Why over-estimation is the safe direction here. Consider the two failure modes. Under-counting means a term that genuinely spiked reports a low number, drops out of the candidate set, and never appears - a real trend silently missing, which is a visible product bug with no way for anyone to notice it happened. Over-counting means a term reports slightly high and may enter the candidate set when it marginally should not have. That false candidate then faces two more filters: the baseline comparison in Deep Dive 3, which rejects anything not deviating from its own history, and the spam checks in Deep Dive 4. A spurious candidate gets killed downstream. A missing one is gone forever. Both CMS and Space-Saving err in the recoverable direction, which is why they are the right family of structure and a lossy hash map is not.

In simple terms: Instead of writing down every word and its tally, keep a fixed grid of counters and bump a few cells per word using several different hash functions. Looking up a word means checking its cells and trusting the smallest one, because collisions can only ever make a cell too big, never too small. The grid is the same size forever, no matter how many new words show up - and if it tells you a count is high, the real count is never higher than that.

flowchart LR
    EV["Term event"]:::client
    CMS[("Count-Min Sketch<br/>5 rows by 65536 cols")]:::data
    SS[("Space-Saving table<br/>1000 bounded slots")]:::data
    HEAP["Bounded top-K heap<br/>size 1000"]:::service
    OUT["Local candidates"]:::service

    EV -->|"1. Increment 5 cells"| CMS
    EV -->|"2. Increment or evict min"| SS
    CMS -->|"3. Query estimate as min of rows"| HEAP
    SS -->|"4. Supply candidate terms"| HEAP
    HEAP -->|"5. Emit every 10 seconds"| OUT

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

Deep Dive 2: Sliding Windows Over an Unbounded Stream

Problem: The product promises a 5-minute, a 1-hour and a 24-hour list, each refreshed every 10-30 seconds. A window that advances every 10 seconds but covers 24 hours is 8,640 overlapping windows a day per region, and events do not arrive in order.

Bad: Recompute the window from raw events on every refresh.

The 24-hour window at 500K events/sec covers 43.2 billion events. Refreshing every 10 seconds means replaying all of them 8,640 times a day. Even at an absurd 10M events/sec of replay throughput, one refresh takes 72 minutes to produce a 10-second-fresh answer. The approach is not slow, it is non-terminating.

Good: Tumbling buckets plus a ring buffer.

Count into fixed, non-overlapping buckets and build windows by summing buckets. One-minute buckets: the 5-minute window is a sum of 5, the hourly window a sum of 60. Store them in a ring buffer of 60 slots indexed by minute mod 60. Advancing one minute is two operations - add the arriving bucket’s counts to the running total, subtract the counts of the slot being overwritten. Not 60 reads, two.

ring[0..59], head = minute mod 60

on bucket close:
    running_total -= ring[head]        # the minute falling out of the window
    ring[head]     = new_bucket
    running_total += new_bucket
    head = (head + 1) mod 60

For the 24-hour window use 1-hour buckets and a 24-slot ring. Nobody needs the daily list to advance in 60-second increments, and the coarser bucket cuts live state by 60x.

Two limits to name. The window slides in whole-bucket steps, so β€œthe last hour” is really β€œthe last 60 completed minutes” and lags by up to one bucket width. And more annoying in practice: a term falls off a cliff at the trailing edge. A hashtag that spiked 59 minutes ago contributes fully at minute 59 and exactly zero at minute 61. The list visibly jumps, and a term can oscillate in and out of rank 50 as its spike bucket crosses the boundary.

Great: Exponentially decayed counts, plus watermarks for out-of-order events.

Replace the hard window with a continuously decaying score. Each term keeps (score, lastUpdatedTs). On an event:

Ξ”t    = now - lastUpdatedTs
score = score Γ— exp(-Ξ» Γ— Ξ”t) + weight
lastUpdatedTs = now

Pick Ξ» from a half-life rather than a window length: Ξ» = ln(2) / halfLife. A 5-minute half-life gives the fast list, a 4-hour half-life gives the daily-ish list. Recency is weighted smoothly, there is no boundary to oscillate across, and a term’s score decays toward zero on its own when the conversation moves on. Decay is applied lazily on touch, and at read time for the few thousand candidates you actually rank - you never sweep the whole key space to age it.

The decayed form has its own cost: the score is no longer a count, so you cannot show β€œ182K posts” from it. Production systems usually run both - decayed scores for ordering, bucketed counts for the number on screen.

Watermarks, and what you do with a late event. Events arrive out of order: a phone was offline, a region’s ingest lagged, a retry fired 40 seconds later. If a window fires the instant wall-clock time passes its end, you systematically undercount every window.

A watermark is the stream processor’s assertion about event time: a watermark of T means β€œI believe I have now seen all events with event time ≀ T.” It is not a clock reading. It is derived from the data - typically max observed event time - allowed out-of-orderness, so a 30-second allowance means the watermark trails the newest event by 30 seconds. A window fires when the watermark passes the window’s end, not when wall-clock does.

An event that arrives after the watermark has already passed its window is late. Three things you can do, and the right answer differs by window:

  1. Drop it, and count the drops. Correct for the 5-minute list. An event 40 seconds late cannot change which 50 strings a human sees, and the drop counter is the metric that tells you if lateness is growing into a real problem.
  2. Allowed lateness. Keep the window in state for an extra period after firing - say 2 minutes - and re-fire with an updated result when a late event lands. You pay for retained state and emit a correction, so downstream must tolerate the same window being emitted twice.
  3. Side output and reconcile. Route late events to a separate stream and fold them in with a slower batch job. This is the right call for the 24-hour list, where a user who was offline for an hour represents real volume you should not throw away.

Watch out: watermarks are generated per source partition and the operator takes the minimum across its inputs. One idle Kafka partition emits no events, so its watermark never advances, so the minimum never advances, so no window anywhere fires. The whole pipeline looks healthy and silently stops producing. Configure idle-source timeouts and alarm on watermark lag, not just on consumer lag.

In simple terms: Do not re-add a day of events every ten seconds. Count into small fixed time slots and add up the slots you care about, keeping them in a circular buffer so advancing one slot is one add and one subtract. Better still, give every term a score that fades on its own over time, so old activity shrinks away smoothly instead of disappearing the moment it crosses a boundary. And since events show up late, use a watermark - the system’s statement about how far along in event time it is - to decide when a window is finished.


Problem: This is the deep dive most candidates skip, and skipping it means the system they designed does not do the thing it was asked to do. A term with a permanently huge baseline is not trending. It is just common. If your ranking cannot tell the difference, you have built a popularity list with a refresh timer.

Bad: Rank by raw count in the window.

Sort the window’s counters and take the top 50, and you get the platform’s permanent vocabulary: its own name, love, lol, the dominant sports league, whatever word the language uses most. These terms win every window of every day forever, because they have genuinely high volume. The list is correct and completely useless - a user who checks it on Monday and Friday sees the same thing, so they stop checking.

Worse, it is inert during exactly the moments it should be loud. An earthquake generates 50,000 mentions in five minutes. The platform’s name generates 2 million. The earthquake does not place.

Good: Rank by the ratio of current-window count to the term’s historical baseline.

score(term) = count_now / baseline(term, region, slot)

Baseline is the term’s own average count over an equivalent window in recent history. Key it by slot_of_week, not just hour-of-day, so a Sunday-evening sports baseline is not compared against a Tuesday-evening one.

This works, and the arithmetic is satisfying:

Term Baseline per 5 min Now Ratio Verdict
#earthquake 100 50,000 500x Trending
platform’s own name 2,000,000 2,100,000 1.05x Not trending, correctly
weather 8,000 9,100 1.14x Not trending, correctly

The specific limit: ratio is scale-blind. A term that went from 2 occurrences to 20 scores 10x and outranks a term that went from 50,000 to 300,000 at 6x. Nobody cares about the first term. At 10M distinct terms a day, the long tail is enormous and almost all of it is capable of a large ratio from a tiny base, so the list fills with noise that is technically spiking.

Great: A standardised deviation from the term’s own expected rate, with floors and smoothing.

Use a z-score against the term’s own history rather than a bare ratio:

z(term) = (count_now - ΞΌ(term, region, slot)) / Οƒ(term, region, slot)

This asks the right question: how surprising is this, measured in units of this term’s own normal variation? A term that routinely swings between 8,000 and 12,000 needs a much bigger jump to be surprising than one that sits dead flat at 100 - and the ratio form cannot express that distinction at all. Three additions make it usable:

1. Absolute floors, applied as a gate before ranking. A term must clear a minimum before it is eligible at all:

eligible(term) =
      count_now        >= 500       # minimum occurrences in the window
  and distinct_users   >= 200       # minimum independent people
  and distinct_users   >= 0.25 Γ— count_now   # not one person posting 400 times

The first condition kills the 2 β†’ 20 case outright, regardless of how extreme its z-score is. The floors are the single most load-bearing part of the ranking function and the first thing to tune when the list looks wrong.

2. Smoothing for terms with no history. A brand-new hashtag has Οƒ undefined, so its z-score is infinite, so every new term wins. Shrink the estimate toward a global prior with weight set by how little history the term has:

ΞΌΜ‚ = (n Γ— ΞΌ_term + k Γ— ΞΌ_prior) / (n + k)
ΟƒΜ‚ = (n Γ— Οƒ_term + k Γ— Οƒ_prior) / (n + k)

n is the term’s sample weeks, k is a smoothing constant (around 3), and the priors come from the global distribution of terms at similar volume. With no history at all, a term is scored as if it were a typical term of its size - so it has to genuinely earn its rank rather than winning by having an undefined denominator. As history accumulates, n dominates and the term’s own numbers take over.

3. A recency penalty for terms already trending. Without one, the biggest story of the day squats rank 1 for eighteen hours. Multiply the score by a decay based on how long the term has already been on the list, so a term that has held rank 1 for two hours has to keep accelerating to stay there. New stories get oxygen.

Put together: score = smoothedZ Γ— recencyPenalty Γ— spamWeight for eligible terms, and nothing at all for terms that fail the gate.

Why z-score beats ratio: both compare against history, but the ratio implicitly assumes every term has the same variance, which is false by orders of magnitude. A term’s standard deviation is the natural unit of β€œunusual” for that term. Using it means a stable term’s modest jump and a volatile term’s big jump are compared on a scale where they are genuinely comparable - which is exactly what a single ranked list requires.

In simple terms: A word being used a lot does not make it trending; a word being used a lot more than it usually is does. So for every term, remember what a normal five minutes looks like for that term at this time of week, and rank by how far above normal it is right now. Then insist that a term clear a real floor of occurrences and distinct people before it can be ranked at all, so a word going from 2 uses to 20 does not beat a genuine news story. And treat a word with no history as an average word of its size, so a hashtag invented thirty seconds ago does not win just because it has no past to compare against.


Deep Dive 4: Spam and Manipulation Resistance

Problem: Being on the trend list is free distribution to millions of people, so the list is a target. A coordinated group - a marketing team, a political campaign, a botnet for hire - pushes a hashtag specifically to get it ranked. And the feedback loop is positive: landing on the list produces real organic traffic, which reinforces the rank the attack bought.

Bad: Count every event equally.

Work out the attack cost. 200 accounts posting 500 times each over five minutes is 100,000 events. That clears a 500-event volume floor two hundred times over. The term is new, so under a naive z-score its baseline is near zero and the score is enormous. Total cost: 200 throwaway accounts and a script. The list is now whatever the cheapest actor wants it to be, and your ranking function actively rewards the attack because manufactured spikes look exactly like news.

Good: Count distinct users, not events.

This single change collapses the attack surface. Rank on the number of distinct accounts using a term, and 200 accounts contribute 200 no matter how many times each one posts. A single account’s contribution is capped at 1 by construction - no per-user bookkeeping, no rate-limit logic, it falls out of the definition. The attacker’s cost shifts from β€œrun a script” to β€œacquire accounts,” which is thousands of times more expensive.

Distinct counting across millions of terms is its own memory problem - an exact set per term would dwarf the frequency counters. HyperLogLog solves it: a fixed register array (~12 KB) per term estimates cardinality within a couple of percent, and the registers merge, so per-partition HLLs combine at the merger without re-reading anything. The HLL++ refinements matter here specifically because most terms have low cardinality and that is where plain HLL is least accurate. Cheaper still: keep HLLs only for the candidate set coming out of Space-Saving - a few thousand terms per region - rather than all 10M.

The limit: distinct-user counting raises the price of an attack, it does not end it. Account acquisition is a functioning market, and a few thousand accounts is well within reach of a funded actor.

Great: Layered defences, none of which is sufficient alone.

  1. Per-user contribution caps. Distinct counting is the cap at 1. The weighted version is more useful: let a trusted account contribute 1.0 and a day-old account 0.05, so the effective distinct count is a reputation-weighted sum rather than a headcount.
  2. Account age and reputation weighting. A three-hour-old account with two followers contributes a fraction of a five-year-old account with real history. This is why accountAgeDays and accountScore are stamped onto the event at ingest - the merger cannot afford a user-service lookup per event, so the weighting inputs must already be in the stream.
  3. Content diversity signals. Organic trends are textually diverse - tens of thousands of people writing different sentences around a term. A coordinated push collapses: near-duplicate text, copy-paste chains, the same shortened URL, the same emoji run. Keep a sampled shingle sketch per candidate term and measure distinct content clusters against occurrence count. A term whose 100,000 posts reduce to eleven templates is not a conversation.
  4. Graph and timing signals. Participants in an organic trend are weakly connected. Participants in a bought one tend to co-follow heavily, to have been created in the same few days, and to post in a suspiciously tight temporal pattern. Burst entropy - how evenly the posting is spread across the window - is cheap to compute in-stream and surprisingly discriminating.
  5. Human review on the top of the list. Only ~50 terms per region matter, and they are read by millions. Gating that tiny set behind a review queue, with the pipeline able to suppress or hold a term pending review, is affordable precisely because the set is small. Large platforms do this, and a candidate who proposes it is describing production rather than theory.

Be honest about the ceiling. This is adversarial, and it is never finished. Every signal you publish becomes a specification the attacker builds against: cap by distinct users and they buy more accounts, weight by age and they age accounts before use, score content diversity and they generate varied text with a language model. The goal is not a solved problem, it is a cost curve - make manipulation expensive enough that it is not worth doing at scale, detect the attempts you can, and keep a human in the loop for the fifty strings that actually reach people’s screens.

In simple terms: Count people, not posts. Once a term’s score depends on how many different accounts used it, one person shouting a thousand times counts the same as shouting once, and the attacker has to buy accounts instead of running a loop. Then make accounts unequal - a brand-new one counts for almost nothing - and look at whether the posts around a term actually say different things, because a real conversation is messy and a coordinated push is a hundred copies of the same sentence. None of this is ever finished; it just makes cheating expensive.


13. Design Self-Audit

Question Answer
Can a term be missed entirely by two-tier aggregation? Yes, and it is the design’s real accuracy risk. A term below the local cutoff on every partition appears in zero candidate lists. With round-robin partitioning and local top-1000 per aggregator, a term must rank below 1000 on all 64 partitions while being globally top-50 - which requires high volume spread almost perfectly uniformly. A genuinely spiking term spikes on every partition and clears 1000 easily. The backstop is a slower exact batch job over the Kafka log that recomputes the window every few minutes and can inject a missed term.
What survives a stream processor restart? Flink checkpoints keyed state (sketches, Space-Saving tables, HLLs, ring buffers) to durable storage incrementally, so a restart restores the last checkpoint and replays Kafka from that offset. The 5-minute window is fully rebuilt within one window length. The 24-hour window is the expensive one - if checkpoint state is lost rather than stale, the daily list needs a replay of up to 48 hours of log, which is why Kafka retention is set to 48h and not 6h. Critical detail: Count-Min Sketch hash seeds must be checkpointed with the counters, or the restored sketch is silently meaningless.
How are late and out-of-order events handled? Watermarks are set to max event time - 30s, so windows fire on event time rather than wall clock. For the 5-minute list, events later than that are dropped and counted - a 40-second-late event cannot change which 50 strings a human sees. For the 24-hour list, late events go to a side output and are folded in by the reconciliation job. Watermark lag is alarmed on separately from consumer lag, because an idle Kafka partition can freeze the global watermark and silently stop all window firing.
Is the trend list identical for everyone in a region? Yes, and that is the single most valuable property on the read path. The response is byte-identical for every user in a region for the current 10-second slice, so a CDN with a 10-second TTL absorbs effectively all of the 100K reads/sec and origin sees a couple of requests per second. The moment personalised trends enter scope (explicitly below the line) this property dies and the read side becomes the hard part of the design.
What does cold start look like for a brand-new term? It has no baseline row, so Οƒ is undefined and a naive z-score is infinite. Smoothing shrinks it toward a global prior for terms of its volume, so it is scored as a typical term of that size and must earn its rank. It must also clear the absolute floors (500 occurrences, 200 distinct users, and a distinct-user ratio), which a brand-new term with real momentum clears within a window or two and a manufactured one does not clear cheaply.
What happens if the pipeline dies entirely? Redis keys carry a 120-second TTL, so the lists expire rather than lingering, and the API returns β€œtrends unavailable” instead of confidently serving an hour-old list a user cannot tell is stale. Kafka keeps accepting writes throughout, so no data is lost - the pipeline catches up from its last committed offset. There is no hot-key risk on the way back in, because round-robin partitioning spreads a viral hashtag across every partition rather than pinning one subtask.

14. Core Flows

Flow: Ingest to Trend List

sequenceDiagram
    participant C as Client
    participant API as Ingest API
    participant K as Kafka
    participant A as Partial Aggregators
    participant M as Top-K Merger
    participant B as Baseline Store
    participant R as Redis

    C->>API: POST event with text and region
    API->>API: Tokenise into terms and stamp account signals
    API->>K: Append record to term-events
    API-->>C: 202 Accepted
    K->>A: Consume partitions round-robin
    A->>A: Increment sketch and Space-Saving slots
    A->>A: Add user to term HyperLogLog
    Note over A: 10s timer fires
    A->>M: Emit local top 1000 candidates
    M->>M: Sum counts per term across aggregators
    M->>B: Fetch mean and stddev for candidates
    B-->>M: Baseline rows for region and slot
    M->>M: Compute z-score apply floors and spam weight
    M->>R: SET trends per region and window

Step-by-step walkthrough:

  1. The client posts content. The Ingest API authenticates, derives region server-side, tokenises, and stamps account age and reputation onto the record. One Kafka record carries all terms from the event, and the API returns 202 without waiting for any counting - the client path never blocks on aggregation.
  2. Partial aggregators consume their partitions and do purely local work: increment the Count-Min Sketch, update Space-Saving slots, add the user to the term’s HyperLogLog. No network hop per event.
  3. A 10-second timer fires on each aggregator, which walks its bounded candidate set through a heap and emits its local top 1,000.
  4. The merger sums per-term counts across all 64 candidate lists - tens of thousands of terms, not millions.
  5. Baselines for the surviving candidates come back in one range scan. The merger computes z-scores, applies the eligibility floors and spam weights, takes the final 50, and overwrites the materialised list with one Redis SET per region and window pair.

Non-obvious failure: the baseline store is unavailable when the merger needs it. Without baselines there are no z-scores, and the tempting fallback - rank by raw count - produces exactly the useless generic list from Deep Dive 3. The right fallback is to hold the previous ranking and keep serving it while the TTL allows, marking the list stale internally, because a slightly old correct list beats a fresh wrong one. Pair that with a small in-merger LRU of recent baseline rows so a short outage is invisible.

Flow: Read the Trend List

sequenceDiagram
    participant U as Reader
    participant CDN as CDN Edge
    participant T as Trends API
    participant R as Redis

    U->>CDN: GET trends for region IN and window 5m
    alt Cache hit within the 10s slice
        CDN-->>U: 200 trend list in under 20ms
    else Cache miss
        CDN->>T: Forward request
        T->>R: GET trends IN 5m
        R-->>T: Serialised list of 50
        T-->>CDN: 200 with Cache-Control max-age 10
        CDN-->>U: 200 trend list
        CDN->>CDN: Cache for the rest of the slice
    end

Step-by-step walkthrough:

  1. The reader requests the trend list for their region. No user identity is needed - the response does not depend on who is asking, which is the whole reason this path is cheap.
  2. The CDN edge checks its cache for this region and window. During the current 10-second slice it almost always has it.
  3. On a miss, the edge forwards to the Trends API, which does exactly one Redis GET and no computation. The response carries max-age=10, matching the pipeline’s refresh cadence, so the edge never serves a list more than one refresh behind.
  4. Every subsequent reader in that region for the rest of the slice is served from the edge. At 100K reads/sec across 20 regions, origin handles single-digit QPS.

Non-obvious failure: a Redis key expires at the moment the pipeline is lagging, so the API has nothing to return - and 100K reads/sec of cache misses all arrive at once. That is a cache stampede against an empty store. Two mitigations: request coalescing at the edge so only one request per region reaches origin, and a stale-while-revalidate directive so the edge keeps serving the last good list for a bounded grace period while the refresh is attempted.


15. Final Architecture

Everything from the deep dives in one picture: bounded-memory sketches in the partial aggregators (Deep Dive 1), bucketed and decayed windows (Deep Dive 2), baseline-relative scoring in the merger (Deep Dive 3), and the spam layer plus review queue gating the output (Deep Dive 4).

flowchart TD
    APP["Clients"]:::client
    API["Ingest API<br/>tokenise and tag"]:::service
    KAFKA["Kafka<br/>term events 48h retention"]:::async
    AGG["Partial Aggregators<br/>Count-Min Sketch and Space-Saving"]:::service
    LATE["Late event side output"]:::async
    MERGE["Top-K Merger<br/>sum and re-rank"]:::service
    SCORE["Scoring and Spam Layer<br/>z-score floors and weights"]:::service
    BASE[("Baseline Store<br/>mean and stddev per slot")]:::data
    RECON["Reconciliation Batch Job"]:::async
    REDIS[("Redis<br/>60 materialised lists")]:::data
    TAPI["Trends API"]:::service
    CDN["CDN<br/>10s TTL"]:::edge
    READER["Readers"]:::client

    APP -->|"Post event"| API
    API -->|"Append record"| KAFKA
    KAFKA -->|"Consume round-robin"| AGG
    AGG -->|"Route events past watermark"| LATE
    AGG -->|"Emit local top 1000"| MERGE
    MERGE -->|"Hand candidates to scorer"| SCORE
    SCORE -->|"Fetch term baselines"| BASE
    SCORE -->|"Write ranked lists"| REDIS
    LATE -->|"Fold in late volume"| RECON
    RECON -->|"Refresh baselines nightly"| BASE
    RECON -->|"Inject missed terms"| REDIS
    TAPI -->|"Read list on miss"| REDIS
    CDN -->|"Forward on miss"| TAPI
    READER -->|"Get trends for region"| CDN

    classDef client fill:#4c3a5e,stroke:#818cf8,color:#e2e8f0
    classDef edge fill:#1e3a5f,stroke:#60a5fa,color:#e2e8f0
    classDef service fill:#1a3a2a,stroke:#4ade80,color:#e2e8f0
    classDef async fill:#3b1f5e,stroke:#c084fc,color:#e2e8f0
    classDef data fill:#3b3520,stroke:#fbbf24,color:#e2e8f0

How it works end-to-end (write path):

  1. Clients post events into Kafka β€” the Ingest API tokenises content, derives region server-side, stamps account age and reputation, and appends one record. Round-robin partitioning keeps ingest flat under a viral hashtag, and 48-hour retention is the replay budget for recovery and ranking changes
  2. Partial aggregators count in bounded memory β€” Count-Min Sketch for frequency, Space-Saving for the candidate set, HyperLogLog for distinct users, all in local RocksDB-backed state, with each aggregator emitting its local top 1,000 every 10 seconds
  3. The scoring layer decides what β€œtrending” means β€” z-score against each term’s own baseline, gated by absolute floors and weighted by spam signals, then the final heap takes 50
  4. Ranked lists land in Redis β€” one whole-value write per region and window, 120-second TTL so a dead pipeline expires rather than serving stale confidence

How it works end-to-end (read path):

  1. Readers hit the CDN first β€” the response is byte-identical for everyone in a region, so a 10-second TTL absorbs effectively all of the 100K reads/sec
  2. Cache miss falls through to the Trends API β€” one Redis GET, no computation, response tagged max-age=10 to match the refresh cadence
  3. Late events and reconciliation close the loop β€” the side output feeds a batch job that folds late volume into the daily window, rebuilds baselines nightly, and acts as the backstop that can inject a term the two-tier merge missed

Key Technologies

Term What it is
Count-Min Sketch Fixed-size 2D counter array with one hash function per row. Increment all rows, query the minimum. Over-estimates, never under-estimates. Width sets error size, depth sets error frequency.
Space-Saving Heavy-hitter algorithm with exactly m slots. Evicts the minimum counter and lets the newcomer inherit its count as an error bound. Anything above N/m of the stream is guaranteed present.
HyperLogLog Approximate distinct-count in ~12 KB per set, within a couple of percent. Registers merge, so per-partition sketches combine for free. Used here to count people rather than posts.
Watermark The processor’s assertion that all events with event time ≀ T have arrived. Derived from observed event times minus an allowed lateness, not from a clock. Windows fire when the watermark passes their end.
Two-tier aggregation Many parallel aggregators each emit a local top-K, one merger sums and re-ranks. Bounds network and merger cost; introduces the risk of missing a term below every local cutoff.
Tumbling bucket A fixed, non-overlapping time slice. Sliding windows are built by summing buckets, so advancing the window is one add and one subtract against a ring buffer.
Exponential decay score = score Γ— e^(-λΔt) + weight. Replaces a hard window with a half-life, so recency is weighted continuously and nothing falls off a cliff at the boundary.

What’s Expected at Each Level

Mid-level

Get the pipeline shape right: events into Kafka, a stream processor counting into time buckets, a bounded min-heap instead of sorting the key space, and a tiny serving store the read API hits through a cache. Do the write-amplification arithmetic out loud - 500K events/sec times several terms each is 2.5M increments/sec - and explain why that rules out a counter row per term in a database. Know that a top-50 answer should cost O(n log K), not O(n log n), and that the list caches perfectly because it is identical for everyone in a region.

Senior

Name Count-Min Sketch and explain the mechanism without hand-waving: several hash functions, one row each, increment all rows, query the minimum, over-estimates but never under-estimates, width versus depth controlling error size versus error probability. Design the two-tier aggregation and volunteer its error - a term below every local cutoff is invisible to the merger - then size the local top-(K Γ— f) mitigation. Build sliding windows from tumbling buckets with a ring buffer rather than recomputing from raw events. Most importantly, catch that raw count is the wrong ranking and move to a baseline-relative score, because a design that ranks by volume does not answer the question it was asked.

Staff+

Pick between CMS-plus-heap and Space-Saving on their guarantees and say why over-estimation is the safe direction for this product. Move past ratio-to-baseline to a standardised z-score, and bring the parts that make it survive contact: absolute floors on occurrences and distinct users, smoothing so a term with no history does not win by having an undefined denominator, and a recency penalty so the day’s biggest story does not squat rank 1. Treat spam as a cost curve rather than a solved problem - distinct-user counting via HyperLogLog, reputation weighting fed from signals stamped at ingest, content-diversity detection, and a human gate on the fifty strings that reach millions of screens. Cover the operational edges that bite in production: sketch hash seeds checkpointed with the counters, watermark lag alarmed separately from consumer lag because one idle partition freezes every window, late events dropped for the 5-minute list and reconciled for the 24-hour one, and a Redis TTL short enough that a dead pipeline goes dark instead of serving stale confidence.


🎯 Key Takeaways



Understand the building blocks used in this design:

Discussion

Newest first
You

Free system design + DSA prep. If it helped you crack an interview, consider supporting.

SensAI SensAI
Beta
Listening...
Tap mic to stop voice mode

Shape what we build next

Every piece of feedback is read by the team and directly influences our roadmap.

What type of feedback?

Install SystemCraft

Add to your home screen for instant access, offline reading, and a distraction-free experience.

Offline reading Faster loads No browser tabs App-like feel

Unlock AI Features

One click to activate - no payment, no credit card. Just sign in and you're in.

AI code review and hints
SensAI chat assistant
AI mock interviews
Whiteboard analysis
100% free during early access