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 write volume is not what you think. 500K events/sec, and each event carries several countable terms (hashtags, tokens, bigrams). Call it 5 - that is 2.5M counter updates per second spread across 10M distinct keys.
- Almost all of that work is wasted. Of 10M distinct terms a day, a few thousand have any chance of reaching a top-50 list. You are spending 99.9% of your write budget on rows nobody will ever read.
ORDER BY count DESCover 10M rows is a full sort. Do it every 10 seconds, times 20 regions, times 3 window sizes, and the database does nothing else.- There is no window. A single lifetime counter per term means the list shows what has been popular since launch, not what is happening now. Adding windows multiplies the row count by the number of live buckets.
- Raw count is the wrong ranking anyway. Sort 10M counters by count and you get the platformβs permanent vocabulary - the same generic words every single day. Volume is not trend.
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
- Count-Min Sketch (Cormode and Muthukrishnan, 2005) - βAn Improved Data Stream Summary: The Count-Min Sketch and its Applications.β The frequency-estimation primitive this whole design rests on: a fixed-size 2D counter array that answers βhow often have I seen this termβ with a bounded one-sided error. Cited by name; it is in the Journal of Algorithms, 2005.
- Space-Saving (Metwally, Agrawal and El Abbadi, 2005) - βEfficient Computation of Frequent and Top-k Elements in Data Streams,β ICDT 2005. Built for exactly this question rather than adapted to it: m counters, a guarantee that anything above
N/mof the stream is in the table, and an error bound tied to stream length. - Twitterβs trends pipeline and the Storm/Heron lineage - Twitter ran trend detection on Storm, then rewrote the engine as Heron for backpressure and per-component resource isolation while keeping the Storm topology API. The shape we borrow is the topology itself: parallel counting bolts feeding a low-parallelism ranker. (Apache Heron repo)
- Apache Flink event time and watermarks - The clearest production treatment of out-of-order streams: watermarks, allowed lateness, and side outputs for events that miss their window. We use all three. (Flink docs on time)
- HyperLogLog++ (Heule, Nunkesser and Hall, 2013) - βHyperLogLog in Practice,β EDBT 2013. Googleβs refinement of HLL with much better accuracy at low cardinality, which matters here because most terms have few distinct users and that is precisely the range where plain HLL is worst. (Google Research publications)
4. Functional Requirements
Core (Top 3)
- Return the top K trending terms for a time window - K = 50, ordered, with a score you can explain
- Support several windows - 5 minutes, 1 hour, 24 hours, all live at the same time
- Scope trends by region - a user in Mumbai and a user in Berlin see different lists
Below the Line
- Personalised trends based on who you follow
- Trend explanations (βwhy is this trendingβ)
- Spam and manipulation defence beyond the basics in Deep Dive 4
- Multi-language tokenisation, CJK segmentation, transliteration
- Historical trend browsing and βtrends on this day last yearβ
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
- Exact frequency counts
- Strong consistency between regions
- Sub-second trend detection
6. Scale Estimation (Back-of-Envelope)
- Events: 500K/sec peak. At ~5 countable terms per event, that is 2.5M term increments/sec.
- Distinct terms: 10M/day globally, with a long tail - roughly 300K distinct terms in any given minute, and most of them appear once.
- Candidate set: of those 10M, maybe 2,000-5,000 terms per region per day have the volume to be rankable. Three orders of magnitude of the key space is dead weight.
- Output size: 20 regions Γ 3 windows Γ 50 entries = 3,000 rows total. The entire answer the system exists to produce is a few hundred KB.
- Read side: 100K reads/sec Γ ~8KB response = ~800 MB/sec of egress if uncached. With a 10-second CDN TTL and 20 regions, origin sees 2 requests/sec.
- Sketch memory: a Count-Min Sketch at depth 5 and width 65,536 with 32-bit counters is 5 Γ 65536 Γ 4 = 1.3 MB, fixed, regardless of how many distinct terms arrive.
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
- Event - one tweet, search or post, carrying a timestamp, a region, a user ID and the terms extracted from it
- Term - the countable unit: a hashtag, a word, or an n-gram phrase
- Bucket - a fixed slice of time (1 minute, 5 minutes, 1 hour) holding counts for terms seen in that slice
- Sketch - the bounded-memory structure that holds approximate frequencies for a bucket
- Baseline - a termβs own historical expected rate for this region and this slot of the week
- Trend list - the materialised ordered 50 for one region and one window
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
userIdorregion- both are derived server-side from the authenticated session and the request IP. A client that can set its ownuserIdcan 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:
- Ingest API - terminates the client request, authenticates it, derives
regionfrom IP or account setting, tokenises the content into terms, and appends one record to the log. Returns202immediately.
π‘ 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. - 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.
- 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:
- A user posts βEngland just won the #worldcupfinalβ - the Ingest API authenticates, stamps
region=IN,tsEvent - API tokenises into terms:
#worldcupfinal,england,won,world cup final - API appends one record with all four terms to the
term-eventsKafka topic, partitioned round-robin - Flink consumes the partition, splits the record into one update per term
- For each term, Flink updates the count for
(term, region, current 1-minute bucket)in local RocksDB state - 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:
- 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.
- 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:
- 64 partial aggregators each consume their Kafka partitions and count terms locally
- Every 10 seconds, a processing-time timer fires on each aggregator
- Each aggregator walks its counters through a bounded heap and emits its top candidates
- The merger receives 64 candidate lists and sums counts per term across all of them
- One final heap runs over the union - a few thousand terms, not 10 million
- 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:
- 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.
- 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.
- 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:
- Merger finishes a refresh cycle and holds 60 ranked lists, one per region and window pair
- Merger writes each as a single Redis key:
SET trends:IN:5m <json>with a 120-second TTL - Reader opens the app, client calls
GET /v1/trends?region=IN&window=5m - Request lands on a CDN edge. Cache hit for the current 10-second slice, served in under 20ms
- On miss, the edge forwards to the Trends API, which does one Redis
GETand returns withCache-Control: public, max-age=10 - 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:
#worldcupfinalis tokenised at the Ingest API and written into one Kafka record among five terms- A partial aggregator increments it in the Space-Saving table for
(IN, 1m, 1760000000)and adds the user to the termβs HyperLogLog - The 10-second timer fires; the term is well inside the local top 1,000, so it is emitted as a candidate
- The merger sums the termβs count across 64 aggregators and joins the baseline row for
(IN, slot 412) - Its z-score against its own history is large, it clears the volume and distinct-user floors, and the spam filter passes it
- It is written into
trends:IN:5mat 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.
- Increment a term: for each row
i, computecol = h_i(term) mod wand add 1 tocounters[i][col].dwrites, no allocation, no key storage. The term string is never stored at all. - Query a term: compute the same
dpositions and take the minimum of thedcounters.
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):
- If the term is already tracked, increment its count.
- If not and a slot is free, install it with
count = 1, error = 0. - 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 + 1anderror = 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:
- 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.
- 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.
- 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.
Deep Dive 3: Trending Means Rate of Change, Not Volume
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.
- 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.
- 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
accountAgeDaysandaccountScoreare 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. - 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.
- 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.
- 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:
- 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
202without waiting for any counting - the client path never blocks on aggregation. - 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.
- A 10-second timer fires on each aggregator, which walks its bounded candidate set through a heap and emits its local top 1,000.
- The merger sums per-term counts across all 64 candidate lists - tens of thousands of terms, not millions.
- 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
SETper 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:
- 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.
- The CDN edge checks its cache for this region and window. During the current 10-second slice it almost always has it.
- On a miss, the edge forwards to the Trends API, which does exactly one Redis
GETand no computation. The response carriesmax-age=10, matching the pipelineβs refresh cadence, so the edge never serves a list more than one refresh behind. - 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):
- 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
- 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
- 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
- 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):
- 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
- Cache miss falls through to the Trends API β one Redis
GET, no computation, response taggedmax-age=10to match the refresh cadence - 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
- Trending is top-K over an unbounded stream, not a leaderboard. No known key set, no exact score anyone can verify, and approximation is the correct answer rather than a compromise.
- Count-Min Sketch gives fixed memory and one-sided error. 1.3 MB per bucket regardless of cardinality, and it only ever over-counts - so a real spike can never be silently lost. Pair it with Space-Saving, which holds the candidate set and the counts together under a hard ceiling of
mslots and guarantees anything aboveN/mof the stream is still in the table. - Two-tier aggregation is how you avoid sorting 10M counters, and keeping a local top-(K Γ f) is how you stop it from dropping a term that is 51st on every partition.
- Build sliding windows from tumbling buckets, or better, decay scores exponentially so nothing falls off a cliff at the window edge.
- Rank by rate of change, not volume. A z-score against the termβs own baseline, with floors on occurrences and distinct users, and smoothing so a brand-new term cannot win on an undefined denominator.
- Count people, not posts. Distinct-user counting via HyperLogLog caps any single accountβs contribution at 1 by construction, which is the cheapest large win against manipulation.
Related Designs
- Leaderboard - the exact-counting counterpart. Bounded known players, authoritative scores, Redis sorted sets. Read the two together to see why unbounded cardinality forces a completely different structure.
- Search Autocomplete - uses a trending signal as a ranking boost on top of a trie, with an EMA-based spike detector. Same velocity idea, different consumer.
- Twitter Feed - the system that produces most of these events, and the fan-out patterns behind them
- Metrics Monitoring - the same streaming aggregation and pre-computation machinery applied to numeric time series instead of string frequencies
Related Concepts
Understand the building blocks used in this design:
- Bloom Filters β β the probabilistic-structure intuition that Count-Min Sketch and HyperLogLog build on
- Batch vs Stream β β why the trend list is computed in-stream while baselines are rebuilt in batch
- Message Queues β β Kafka absorbs the 500K events/sec burst and provides the replay window for recovery
- Fan-Out β β the partial-aggregator tier is fan-out on write, with the merger collapsing it back to one answer
- Caching β β a regionβs trend list is identical for every user, which is what makes a 10-second CDN TTL absorb 100K reads/sec
Discussion
Newest first