Interview Setup
Interview Prompt
Design a trending topics system like Twitter's Trends. Detect hashtags and keywords with abnormal growth, rank by velocity, update every 1 to 5 minutes, and show top 20 per region.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| Velocity vs raw count: which metric drives ranking? | Taylor Swift consistently generates 500,000 mentions per hour but is not trending. Velocity surfaces genuine statistical spikes rather than baseline popular topics. |
| What event volume and unique topic cardinality must the system support? | Ingesting 500,000 events per second alongside 10 million unique topics per hour makes exact per-topic in memory counting cost-prohibitive. |
| How do you detect and handle bot or spam amplification? | Coordinated bot networks can simulate organic spikes. Inline filtering must evaluate accounts before scoring rather than relying on delayed post-hoc cleanup. |
| Are trends computed per region or globally? | Serving more than 200 geographic regions requires partitioned Flink state and regional Redis sorted sets, while viral national events create partition hotspots. |
Scope
In scope
- Sliding window counting
- Exponential decay and velocity scoring
- Count-Min Sketch for approximate counts
- Heavy hitters algorithm
- Spam and bot filtering with bounded candidate analysis
- Regional top-K serving from Redis
Out of scope (state explicitly)
- Full content moderation pipeline
- Personalized feed ranking (only trend re-ranking)
- Real time full text search on tweets
Functional Requirements
Clarify the core definition of trending early in the discussion, distinguishing velocity growth from raw volume, establishing sliding window boundaries, and confirming regional granularity. Key architectural requirements include token and hashtag extraction, stateful sliding window counters, and a low latency ranked serving API.
In an interview setting, clarify how the platform defines trending, because ranking first by total volume over the past hour requires an entirely different architecture from real time spike detection.
- Detect trends: Identify topics, hashtags, and keywords experiencing an abnormal surge in volume.
- Real time freshness: Ensure trending topics refresh every 1 to 5 minutes.
- Velocity ranking: Display top 10, 20, or 50 trending topics ranked by acceleration over baseline rather than raw volume.
- Regional boards: Support country and city level leaderboards, with personalized user ranking treated as a staff level extension.
- Explanatory context: Provide context alongside each trend, such as volume estimates and sample representative posts.
- Lifecycle tracking: Detect the emergence, peak velocity, and natural decay of topics over time.
- Abuse filtering: Suppress coordinated bot activity and artificial amplification before topics enter ranking sets.
- Privacy and data minimization: Avoid exposing private source events through public trend context and retain only the event fields required for ranking, auditing, and abuse detection.
Non-Functional Requirements
The system requires freshness within minutes and resilience against coordinated hashtag inflation, where approximate counting balances high ingestion throughput with bounded memory constraints.
- Low Latency: Detect emerging surges within 1 to 5 minutes of onset and serve trend rankings via API with p99 latency under 50 milliseconds.
- High Throughput: Ingest and process more than 500,000 events per second across posts, searches, and interactions.
- Horizontal Scalability: Scale to billions of daily events distributed across more than 200 distinct geographic regions.
- Statistical Accuracy: Maintain a displayed trend false positive rate under 1% for bot inflation and stale baseline anomalies.
- High Availability: Maintain 99.99% uptime for the trend serving API and streaming ingestion pipeline.
- Event time Correctness: Process events by event time with bounded lateness so delayed records do not corrupt already closed ranking windows.
Capacity Estimations
Inbound event velocity and unique topic cardinality determine stream processing state size, partition counts, and Redis memory footprints.
| Metric | Calculation | Value |
|---|---|---|
| Events / sec (tweets + searches + posts) | Derived from daily volume ÷ 86400 (+ peak factor) | 500K |
| Equivalent events / day at 500K/sec sustained | 500,000 x 86,400 | 43.2B |
| Unique topics / hour | Given | 10M |
| Trending topics displayed | Given | Top 20 per region |
| Regions | Given | 200+ countries |
| Topic velocity window | Given | 5-minute sliding window |
| Storage for trend history | Given | 10 GB/day |
Architecture Diagram
Inbound posts and interactions pass through an ingestion gateway into an event bus, where streaming workers extract topics, evaluate sliding window velocity against historical baselines, and publish regional top-K leaderboards into an in-memory datastore for fast query serving.
Component Deep Dives
Velocity Over Volume: The Core Algorithm
Velocity scoring, approximate frequency tracking, inline bot suppression, and rolling baseline calculation operate together inside the streaming pipeline before aggregated leaderboards reach Redis. The core mathematical model is shown first, followed by defensive controls and stream partitioning details.
Raw mention count favors topics with high constant volume, such as celebrity names, which fail to capture newsworthy spikes. By tracking velocity as the relative surge above a historical baseline, the system highlights emerging conversations regardless of their baseline scale.
Motivation: Why velocity takes precedence over raw mention volume:
- High-baseline topic ("Taylor Swift"): Consistently generates 500,000 mentions/hour.
Not trending because this represents normal expected volume.
- Low-baseline topic ("#ObscureEvent"): Typically generates 10 mentions/hour, surges to 5,000.
Trending because this represents a 500x statistical anomaly.
Interval Normalization:
The current measurement and historical baseline MUST use the same evaluation interval.
For a 5-minute scoring window, convert the hourly baseline to its equivalent 5-minute rate:
V_baseline_5m = V_baseline_hourly / 12
Mathematical Formulation:
velocity = (V_current - V_baseline) / max(V_baseline, min_floor = 10)
where V_current and V_baseline are counts for the same evaluation interval and min_floor is 10 mentions per interval
Apply a separate minimum-support threshold so tiny samples do not rank highly only because their baseline is small.
Concrete Numerical Walkthrough:
Topic A (Celebrity):
V_current = 200,000 mentions
V_baseline = 100,000 mentions
velocity = (200,000 - 100,000) / 100,000 = 1.0 (2x normal baseline)
Topic B (Breaking Event):
V_current = 5,000 mentions
V_baseline = 10 mentions
velocity = (5,000 - 10) / 10 = 499.0 (500x normal baseline)
Outcome:
Topic B ranks substantially higher than Topic A despite a 40x lower absolute volume.
This formula surfaces genuinely surprising, newsworthy events across any topic scale.
Optional Exponential Smoothing:
smoothed_velocity_t = lambda * velocity_t + (1 - lambda) * smoothed_velocity_(t-1)
where 0 < lambda <= 1
Natural Decay Dynamics:
Following a surge, the rolling baseline gradually incorporates the higher volume,
causing computed velocity to drop back toward zero as the topic transitions:
Emerging -> Peak -> Decaying -> NormalizedSpam and Manipulation Detection
Coordinated bot networks frequently attempt to game trending leaderboards using automated astroturfing campaigns. The streaming pipeline evaluates multi dimensional behavioral heuristics inline to neutralize inorganic surges before they enter scoring stages.
Inline Stream Filtering Heuristics in Apache Flink: 1. Account Age Weighting: If more than 30% of mentions originate from accounts younger than 7 days, penalize the final velocity score by multiplying it by 0.2. 2. Account Diversity Ratio: For promoted candidate topics, estimate distinct account count with HyperLogLog. If more than 50% of volume stems from fewer than 100 unique account IDs, suppress the topic completely to block coordinated bot farms. 3. Near-Duplicate Text Clustering (SimHash): Compute 64-bit SimHash fingerprints only for promoted candidate topics. Bucket fingerprints by leading bits before Hamming-distance comparison so the system avoids all-pairs comparison across the full 500,000 events/sec ingress stream. If more than 60% of candidate posts exhibit a Hamming distance <= 3, mark the volume as copy paste astroturfing and suppress. 4. Growth Velocity Curve: As a heuristic, organic human trends often ramp gradually over a 15 to 30 minute window. Bot campaigns can manifest as an instantaneous step function within a single minute, triggering automated rate-limiting flags. 5. Geographic Telemetry Validation: Detect mismatches between the targeted trending region and the timezone or IP origin of participating accounts as an anomaly signal rather than an automatic block, because VPNs, travel, and shared networks can produce legitimate mismatches.
Baseline Management and Seasonality
Topic volume naturally exhibits strong cyclical patterns across hours of the day and days of the week. Maintaining rolling seasonal baselines ensures regular cyclical spikes are recognized as normal traffic rather than false positive breaking news.
Seasonal Baseline Model: Day of Week combined with Hour of Day
Historical Source of Truth:
ClickHouse topic_hourly_counts retains hourly mention aggregates for the prior 35 days.
Redis Storage Key:
baseline:{region}:{topic_id}:{day_of_week}:{hour}
Cache TTL: 2,592,000 seconds (30 days), refreshed weekly.
day_of_week and hour use the region's canonical local timezone.
Batch Calculation (Apache Spark, scheduled weekly):
V_baseline_hourly = average(historical hourly volume for this day and hour across the prior 4 weeks)
V_baseline_5m = V_baseline_hourly / 12
Concrete Example ("#MondayMotivation"):
- Historical Monday 08:00 baseline: 50,000 mentions/hour (regular recurring morning spike)
- Historical Tuesday 08:00 baseline: 2,000 mentions/hour
Evaluation Scenarios:
The following worked examples use the same hourly interval for current and baseline counts so the original arithmetic remains directly comparable. Production scoring uses the 5-minute window after converting the hourly baseline with V_baseline_5m = V_baseline_hourly / 12.
- Scenario 1 (Monday at 08:00 with 55,000 mentions):
velocity = (55,000 - 50,000) / 50,000 = 0.1
Result: Not trending because volume matches expected Monday morning seasonality.
- Scenario 2 (Tuesday at 08:00 with 20,000 mentions):
velocity = (20,000 - 2,000) / 2,000 = 9.0
Result: Actively trending because a 10x surge is highly unusual for a Tuesday.
Cold-Start Provisioning:
Unseen topics with no historical record default to min_floor = 10 as their baseline.Event Bus Architecture (Kafka)
Partitioning strategy on the message broker distributes ingestion load, while Flink keys state by region and topic and uses event-time windows and watermarks to handle out-of-order arrival across Kafka partitions.
Topic: raw-events
Configuration:
Partitions: 64
PartitionKey: "region + event_hash % 32"
RetentionHours: 24
ReplicationFactor: 3
MinInSyncReplicas: 2
ProducerAcks: all
EnableIdempotentProducer: true
Idempotency: producer reuses the same eventId across retries
ProducerEventSchema:
eventId: "uuid"
source: "tweet | search | post | webhook"
text: "string"
region: "string"
userId: "string"
timestamp: "int64"
PipelineStages:
- topic: raw-events
consumer: "topic-extractor"
output: normalized-events
responsibility: "Named entity recognition, hashtag extraction, canonical topic_id assignment, and keyword normalization"
- topic: normalized-events
consumer: "spam-filter"
output: filtered-events
partitionKey: "(region, topic_id_hash % 32) for normal topics, whereas hot candidate shards may use (region, topic_id, event_hash % 100)"
responsibility: "Account-age checks, account-diversity estimation, bounded SimHash candidate analysis, and growth-anomaly detection"
- topic: filtered-events
consumer: "trend-detector"
outputs:
- "regional top-K"
- "ClickHouse trend_snapshots"
partitionKey: "(region, topic_id_hash % 32) for normal topics, or (region, topic_id, event_hash % 100) for hot topics"
responsibility: "5-minute event time sliding windows, seasonal-baseline comparison, ranking, lifecycle tracking, and periodic snapshot publication with secondary merge for salted hot topics"
- topic: filtered-events
consumer: "hourly-aggregator"
output: "ClickHouse topic_hourly_counts"
partitionKey: "(region, topic_id_hash % 32)"
responsibility: "Hourly topic mention aggregation for the 35-day regional seasonal baseline history"
HotRegionPartitioning:
rawKey: "region + event_hash % 32"
normalKey: "region + topic_id_hash % 32"
hotKey: "region + topic_id + event_hash % 100"
mergeStage: "Secondary Flink aggregation combines salted partial counts by (region, topic_id)"
note: "Raw ingress uses event hashing so a hot region does not bottleneck one Kafka partition. After topic extraction, normal topics are keyed by (region, topic_id) semantics, while detected hot topics use event salting and a secondary merge restores the canonical (region, topic_id) aggregate."
PipelineOutput:
destination: "Redis sorted set trending:{region}"
role: "Materialized serving view. Flink checkpoints are the recovery source, while ClickHouse aggregates and snapshots are durable analytics sources"
frequency: "Top 50 trends per region emitted every 60 seconds"
deadLetterQueue: "raw-events-dlq for malformed or permanently unprocessable events"
operationalAlert: "Alert if Flink checkpoint completion lag exceeds 120 seconds"API Design
Service Interface and Domain Types
The trending topics service provides low latency REST endpoints for fetching regional leaderboards and inspecting topic level velocity metrics. Non-personalized regional responses can use shared edge caching, while personalized responses require user scoped private caching and must not enter a shared public cache.
Domain interfaces specify strongly typed request and response contracts for leaderboard retrieval and velocity inspection.
// Trending Topics Service Interface and Domain Types
type TopicId = string;
type RegionCode = string;
type UserId = string;
type ISODateTime = string;
export interface TrendingTopic {
rank: number;
topicId: TopicId;
topic: string;
volume: string;
velocity: number;
category: string;
startedAt: ISODateTime;
context?: string;
representativePostIds?: string[]; // public, trend-eligible post IDs only
}
export type GetTrendingTopicsRequest =
| {
region: RegionCode;
count?: number; // default: 20, max: 50
category?: string;
personalized?: false;
viewerId?: never;
}
| {
region: RegionCode;
count?: number; // default: 20, max: 50
category?: string;
personalized: true;
viewerId: UserId;
};
export interface GetTrendingTopicsResponse {
trends: TrendingTopic[];
asOf: ISODateTime;
region: RegionCode;
personalized: boolean;
}
export interface GetTopicVelocityRequest {
topicId: TopicId;
region: RegionCode;
windowMinutes?: number; // default: 5
}
export interface GetTopicVelocityResponse {
topicId: TopicId;
topic: string;
region: RegionCode;
currentCount: number;
baselineCount: number;
windowMinutes: number;
velocity: number;
isTrending: boolean;
computedAt: ISODateTime;
}
export interface TrendingTopicsService {
getTrendingTopics(request: GetTrendingTopicsRequest): Promise<GetTrendingTopicsResponse>;
getTopicVelocity(request: GetTopicVelocityRequest): Promise<GetTopicVelocityResponse>;
}Get Trending Topics Endpoint
Fetches the top trending topics for a specified geographic region with category tags, volume estimates, and contextual summaries.
GET /api/v1/trends?region=US&count=20&personalized=false HTTP/1.1
Host: api.example.com
Accept: application/json
HTTP/1.1 200 OK
Content-Type: application/json
Cache-Control: public, max-age=60
# Public caching applies only when personalized=false or omitted.
# Personalized responses must use private user scoped caching.
{
"trends": [
{
"rank": 1,
"topic_id": "entity:fifa-world-cup-2026",
"topic": "#WorldCup2026",
"volume": "2.3M",
"velocity": 15.2,
"category": "Sports",
"started_at": "2026-07-10T08:00:00Z",
"context": "Quarter-finals underway",
"representative_post_ids": ["post_123", "post_456"]
}
],
"as_of": "2026-07-10T10:05:00Z",
"region": "US",
"personalized": false
}Topic Velocity Inspection Endpoint
Provides detailed acceleration and baseline metrics for a specific topic. Counts are measured over the requested window, which defaults to 5 minutes, and are intended for auditing trend calculations.
GET /api/v1/trends/velocity?topic_id=entity%3Afifa-world-cup-2026®ion=US HTTP/1.1
Host: api.example.com
Accept: application/json
HTTP/1.1 200 OK
Content-Type: application/json
{
"topic_id": "entity:fifa-world-cup-2026",
"topic": "#WorldCup2026",
"region": "US",
"current_count": 2300000,
"baseline_count": 141975,
"window_minutes": 5,
"velocity": 15.2,
"is_trending": true,
"computed_at": "2026-07-10T10:05:00Z"
}Common Error Responses
Standard error responses follow platform conventions for validation failures, missing parameters, and rate limiting.
400 Bad Request: invalid input, missing required fields, or malformed JSON payload 401 Unauthorized: missing or invalid authentication token or API key 403 Forbidden: authenticated caller lacks required permissions for this resource 404 Not Found: requested resource ID does not exist 409 Conflict: duplicate write or version conflict, retry with a unique idempotency key 422 Unprocessable Entity: syntactically valid request failed semantic business validation 429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header 500 Internal Error: unexpected server failure, retry safely with an idempotency key 503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff 504 Gateway Timeout: search index shard responded slowly, narrow query parameters or retry
Data Model
Redis In-Memory Serving Store
The data layer spans volatile in-memory sorted sets for low latency serving, stateful RocksDB checkpoints for stream processing, and columnar storage for long-term historical trend snapshots.
Redis sorted sets store regional leaderboards scored by velocity for logarithmic range queries, while associated hashes store topic metadata and volume strings.
# Redis Data Structures
Normal key: trending:{region}
Type: Sorted Set (ZSET)
Score: velocity (IEEE 754 floating point)
Member: topic_id
Tie-break: deterministic secondary ordering by topic_id when velocity scores are equal
TTL: None (periodically refreshed and pruned to top 50 every 60 seconds)
Query: ZREVRANGEBYSCORE trending:US +inf -inf WITHSCORES LIMIT 0 20
Hot-region keys: trending:{region}:{shard}
Use only when a regional leaderboard becomes a read hot spot. Query shards in parallel and merge their local top-K results before returning the regional board.
Key: topic_ctx:{region}:{topic_id}
Type: Hash (HSET)
Fields:
display_topic: string (for example, "#WorldCup2026")
volume: string (for example, "2.3M")
started_at: ISO 8601 timestamp string
category: string (for example, "Sports")
context: string (editorial summary or breaking snippet)
representative_post_ids: JSON array of post IDs
TTL: 86,400 seconds (24 hours)
Query: HGETALL topic_ctx:US:entity:fifa-world-cup-2026
Key: baseline:{region}:{topic_id}:{day_of_week}:{hour}
Type: Integer string
Value: Rolling 4-week historical average mention volume for this regional hourly slot
TTL: 2,592,000 seconds (30 days), refreshed weekly
Query: GET baseline:US:entity:monday-motivation:1:8Apache Flink Streaming State
Keyed operator state in Apache Flink maintains sliding window counters, circular historical buffers, and deduplication bloom filters backed by embedded RocksDB.
# Apache Flink Keyed State Specification: Keyed by (topic_id: String, region: String)
KeyedStateDescriptor:
minute_buckets:
type: "ListState<Long>"
capacity: 5
description: "Five one-minute event time buckets whose sum forms the active 5-minute sliding window"
history_buffer:
type: "ListState<Long>"
capacity: 12
description: "Circular buffer of previous 5-minute aggregates used for baseline comparison and velocity decay curves"
first_seen_ts:
type: "ValueState<Long>"
description: "First observed event time for topic lifecycle tracking"
peak_velocity:
type: "ValueState<Double>"
description: "Highest observed velocity used to identify peak lifecycle state"
lifecycle_state:
type: "ValueState<Enum>"
values: "Emerging | Peak | Decaying | Normalized"
last_emit_ts:
type: "ValueState<Long>"
description: "Timestamp of last emission to Redis to enforce the 60-second broadcast rate"
event_dedupe_filter:
type: "ValueState<BloomFilter>"
description: "Rotating duplicate-suppression hint for event IDs within the active processing horizon. Never drop an event solely because this Bloom filter reports a match."
distinct_account_estimator:
type: "ValueState<HyperLogLog>"
description: "Maintained only for promoted candidate topics to estimate unique account count"
simhash_candidate_buckets:
type: "ValueState<BoundedMap>"
description: "Bounded SimHash buckets maintained only for promoted candidate topics"
# State Backend and Checkpointing Strategy
OperationalConfig:
state_backend: "EmbeddedRocksDBStateBackend"
checkpoint_interval_ms: 60000 # Checkpoint snapshot persisted every 60 seconds
checkpoint_storage: "s3://trending-service-checkpoints/prod/"
consistency_mode: "EXACTLY_ONCE"ClickHouse Historical Snapshots
ClickHouse archives periodic snapshots of regional leaderboards and retains hourly topic aggregates produced by the streaming pipeline. These aggregates support historical analytics, auditability, and the weekly baseline calculation used for velocity scoring.
CREATE TABLE topic_hourly_counts (
region LowCardinality(String),
topic_id String,
hour_start DateTime,
mention_count UInt64
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(hour_start)
ORDER BY (region, topic_id, hour_start)
TTL hour_start + INTERVAL 35 DAY;
CREATE TABLE trend_snapshots (
region LowCardinality(String),
topic_id String,
velocity Float64,
volume UInt64,
rank UInt16,
snapshot_at DateTime
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(snapshot_at)
ORDER BY (region, snapshot_at, rank);Fault Tolerance
| Concern | Solution |
|---|---|
| Flink worker failure | Persist RocksDB checkpoints to Amazon S3 every 60 seconds and resume execution from the last validated checkpoint snapshot. |
| Redis cluster outage | Redis serving becomes unavailable and the Flink sink path may backpressure while internal state remains recoverable from durable checkpoints. When Redis recovers, the serving layer rebuilds regional top K sorted sets from the latest completed computation output or restored Flink state. |
| Kafka consumer lag | Trends experience proportional processing latency, triggering automated alerts whenever lag exceeds 2 minutes and activating capacity or partition remediation. |
| Coordinated spam surge | Inline filters evaluate account age, distinct-account concentration, bounded SimHash similarity, growth acceleration, and geographic anomalies before ranking. |
| Topic cardinality explosion | Track all topics in the fixed-footprint Count-Min Sketch and promote only a bounded candidate set above the heavy-hitter threshold. |
| Event-time watermark stall | Monitor partition watermarks and source idleness. Mark inactive partitions idle so one stalled partition cannot indefinitely hold back downstream window evaluation. |
Failure Modes and Edge Cases
Streaming trend detection must handle worker failover, partition hotspots, stale baseline fallbacks, event time stalls, and content safety overrides without disrupting leaderboard availability.
The following subsections analyze critical edge cases that arise under real time streaming operations, ranging from window boundary splits to crisis content moderation.
Window Boundary Splitting
When fixed tumbling windows are used, an event spike that spans the boundary between two windows gets split in half, preventing either window from detecting a significant surge. Implementing sliding windows with overlapping intervals ensures the full spike is captured within a single evaluation window.
Tumbling Window Split Boundary Condition: Window 1: [10:00, 10:05) -> Records 2,500 mentions Window 2: [10:05, 10:10) -> Records 2,500 mentions Real-world event: 5,000 mentions concentrated between 10:04 and 10:06 Failure mode: Both windows observe only half the surge volume, causing the spike to miss the velocity threshold. Sliding Window Resolution: Configuration: 5-minute event time window width with 1-minute slide interval Evaluation at T=10:06: Window covers [10:01, 10:06] and captures all 5,000 mentions intact. Allowed lateness: Events arriving after the configured lateness bound are no longer eligible to update the active window and are counted in late-data metrics or routed to a late-data side output. State overhead: Requires 5x more state tracking inside Flink RocksDB to maintain overlapping buckets, which is justified by reliable detection.
Hot Key Partition Skew
High-profile regional events, such as the Super Bowl in the United States, generate disproportionate traffic spikes on a single Kafka partition and Flink worker. Key salting and two-stage aggregation distribute this load across worker instances.
Regional Partition Hotspot:
Event: Super Bowl broadcast drives a 10x traffic surge concentrated in the US region.
Baseline failure mode: Without load spreading, the US region or a single hot topic can concentrate
disproportionate work on one Kafka partition or Flink task while other workers remain largely idle.
Remediation Patterns:
1. Multi-Stage Hierarchical Partitioning:
Split hot regions into geographic sub-keys (US-East, US-West, US-Central).
Each sub-partition aggregates independently, followed by a secondary Flink merge stage.
2. Key Salting with Event Shards:
Partition incoming events by compound key: (region, topic_hash, event_hash % 100).
Spreads a hot topic across 100 first-stage worker shards, followed by a second-stage merge by (region, topic_id).
3. Redis Read Sharding:
Store regional leaderboards as trending:{region}:{shard} and query each shard in parallel,
then merge the shard-local top-K results before returning the regional board.Stale Trend: Outdated Baselines
If upstream batch jobs that update historical baselines experience silent failures, normal recurring events get compared against outdated baselines and trigger false positive trending alerts. Active freshness monitoring and fallback baseline formulas prevent stale trends.
Stale Baseline False-Positive Anomaly:
Context: The hashtag "#MondayMotivation" experiences an expected surge every Monday morning.
Failure mode: The weekly Spark batch job fails silently, leaving a 2-week-old baseline in place.
Outcome: The normal Monday surge gets evaluated against an outdated, lower baseline,
causing the topic to register an enormous velocity score and rank falsely as breaking news.
Automated Detection and Fallback:
1. Freshness Monitoring: Emit a warning when baseline age exceeds 7 days and escalate to a hard failure at 14 days.
2. Conservative Fallback Heuristic:
When baseline age exceeds 14 days, evaluate:
fallback_baseline = max(V_baseline, V_current * 0.5)
This formula guarantees that velocity cannot exceed 1.0 solely due to stale baseline data.
3. Resilience: Trigger automated Spark job retries with exponential backoff and failure alerting.Flink Checkpoint Recovery: Replay Idempotency
When a Flink worker crashes between periodic checkpoints, restarting from the last committed offset causes intermediate Kafka events to be processed a second time. Using Kafka transactional producers and idempotent Redis sorted set operations prevents count doubling.
Worker Crash and Checkpoint Replay Scenario (simplified single-partition example):
T=10:00: Flink persists state checkpoint at Kafka offset 1000.
T=10:00 to 10:03: Worker ingests events from offset 1001 to 1500.
T=10:03: Worker process crashes before writing the next checkpoint.
Recovery: Flink restarts the task from offset 1000, causing events 1001 to 1500 to be re-read.
Naive risk: Replaying events can repeat downstream side effects that were emitted after the last completed checkpoint. Flink's restored state itself starts from the checkpoint boundary before replaying those records.
Idempotent Resolution:
1. Checkpointed State and Offsets:
Flink stores source offsets together with operator state in the checkpoint snapshot,
so recovery resumes from a consistent state boundary.
2. Exactly-Once Kafka Outputs:
Flink can use Kafka transactions so downstream Kafka consumers observe committed results
only after the checkpoint and transactional output are completed.
3. Idempotent Redis Materialization:
Rankings in Redis are written using idempotent ZADD commands with explicit scores,
so replaying an identical velocity calculation overwrites the score rather than incrementing.
Kafka transactions do not make Redis itself transactional.
4. Managed State Recovery:
Worker memory state restores exclusively from the RocksDB snapshot rather than re-reading Redis.Trend Suppression: Sensitive Content
During breaking tragedies or emergencies, raw velocity surges can surface sensitive, graphic, or exploitative topics to the top of public boards. Multi-stage automated filtering paired with editorial tooling ensures responsible contextualization.
Content Moderation and Crisis Response Pipeline:
Challenge: Breaking violent events or crises surge instantly in velocity,
creating the risk that graphic or sensitive content trends inappropriately.
Protection Stages:
1. Automated Keyword Blocklist:
Explicit, violent, and hate terms are dropped before reaching scoring operators.
2. Sentiment and Acceleration Correlation:
Surges combining high negative sentiment with sensitive keywords trigger urgent human review.
3. Trust and Safety Control Panel:
On-call safety personnel can instantly suppress abusive tags across all regional boards.
4. Contextual Story Banners:
Rather than displaying a raw hashtag, platform curators replace the entry with a verified
news headline and authoritative journalistic source link.Additional Considerations
Interview Walkthrough
Structuring the discussion around velocity scoring, streaming window mechanics, and memory-bounded estimation keeps the interview focused on high-signal architectural trade-offs.
- 25-minute cut
Skip deep dive extensions unless interviewing for staff-level scope.
- Define trending as velocity and acceleration rather than raw mention volume over the past hour (5 min)
- Diagram the Kafka ingestion, Flink sliding window aggregation, and Redis sorted set serving pipeline (6 min)
- Quantify 500,000 events per second and 10 million unique topics per hour using Count-Min Sketch estimation (5 min)
- Establish inline bot and spam filtering heuristics before topics enter velocity scoring sets (5 min)
- Staff extension: Address hot partition rebalancing, checkpoint replay deduplication, and cross language entity merging (4 min)
- Lead with velocity over volume because trending represents recent acceleration relative to historical baselines rather than total lifetime mention count.
- Ingest posts and events into Kafka Architecture and Guarantees using a routing key that spreads raw ingress load. Flink then keys state by region and canonical topic_id for normal topics, uses salted shards for detected hot topics, and relies on event-time windows and watermarks rather than Kafka partition ordering for cross-partition correctness.
- Compute windowed counts using principles from Stream Processing Basics with 5-minute sliding windows and 1-minute slide intervals to eliminate window boundary splitting.
- Incorporate a Count-Min Sketch with heavy-hitter promotion to bound Flink state memory when tracking 10 million unique topics per hour, reducing worker memory from 2 GB to 2.2 MB.
- Deploy multi-signal spam detection inline to suppress bot farms, SimHash near duplicate spam, and step function velocity spikes before ranking.
- Serve regional top-K boards with under 50 millisecond latency using Redis Patterns for Interview Systems sorted sets and hash metadata. Keep shared regional responses publicly cacheable, but keep personalized results user scoped.
- Mitigate partition hotspots during viral events by sub-partitioning hot geographic keys and combining partial aggregations in a downstream Flink stage.
- Quantify stream capacity and storage bounds using Back-of-the-Envelope Estimation to validate that 10 GB of daily snapshot storage fits comfortably on disk.
- Common pitfall: Ranking by lifetime total volume, which causes evergreen topics like music or sports teams to dominate permanently and prevents genuine breaking news from surfacing.
Related Problems and Concepts
Explore how real time stream aggregation, in-memory ranking, and distributed counter architectures connect across related system designs:
- Stream Processing Basics: Sliding vs tumbling windows, event time vs processing time, watermarks, and stateful stream joins.
- Redis Patterns for Interview Systems: Sorted sets for real time leaderboards, hashes for topic metadata, and atomic script execution.
- Kafka Architecture and Guarantees: Partitioning keys, consumer groups, log retention policies, and exactly-once transactional sinks.
- Back-of-the-Envelope Estimation: Sizing memory footprints for 500,000 events per second and 10 million unique topics per hour.
- System Design Interview Patterns: High-signal pacing, requirements scoping, and trade-off articulation in distributed systems.
- Distributed Stream Processing: Designing distributed stream engines with fault-tolerant checkpointing, state backends, and backpressure handling.
- Distributed Message Broker: Log storage structures, segment indexing, replication protocols, and high throughput ingestion.
- Like Count for High-Profile Posts: Handling high-concurrency counter updates, hot key mitigation, and write-behind reconciliation.
Engineering Trade-offs
Exact Counting vs Count-Min Sketch
Evaluating trade-offs between memory efficiency and ranking precision, streaming engine architectures, personalization layers, and cross language deduplication determines the operational scalability of the trending platform.
Tracking exact counts for 10 million distinct topics per hour exhausts worker memory, whereas an approximate sketch dramatically reduces memory at the cost of bounded overestimation error. A hybrid sketch-and-promote architecture balances both approaches.
Exact Counting with Hash Maps:
- Memory Complexity: O(N) where N represents unique topics per hour (10,000,000 topics).
- Storage Demand: Requires approximately 2 GB of RAM per Flink worker node under the stated workload assumption.
- Failure Risk: Unbounded cardinality spikes can trigger JVM garbage collection thrashing or OOM crashes.
Count-Min Sketch (Approximate Frequency Estimation):
- Fixed Footprint: The design reserves 200 KB for 4 hash functions and an array width of 2,000 slots, plus implementation overhead.
- Accuracy Bounds: Counts do not under-estimate true frequency, but additive overestimation depends on sketch width and traffic volume.
- Production Tuning: Choose sketch width and depth from an explicit error budget rather than treating 200 KB as a universal accuracy setting.
Hybrid Sketch and Promote Architecture:
1. Stream Ingestion: All 10,000,000 topics increment the compact 200 KB Count-Min Sketch.
2. Promotion Threshold: When a topic's estimated frequency crosses the heavy-hitter threshold,
it becomes eligible for promotion to an exact tracking hash map.
3. Bounded State: Only roughly 10,000 candidate topics are maintained in exact memory at a time,
with a hard cap that evicts the lowest estimated candidates when necessary.
4. Leaderboard Precision: Top-K ranking runs exclusively on promoted exact counters.
5. Memory Economy: Core sketch plus candidate counters are approximately 2.2 MB before auxiliary lifecycle,
SimHash, and metadata state, preserving the original design estimate of roughly 1,000-fold reduction.Flink vs Spark Structured Streaming vs Kafka Streams
Selecting the appropriate stream processing engine requires balancing latency requirements, windowing capabilities, and state management complexity across the pipeline.
Streaming Engine Architectural Comparison:
Feature Apache Flink Spark Structured Kafka Streams
────────────────────────────────────────────────────────────────────────────────
Processing Model True streaming Micro-batch by True streaming
execution default, while execution
Continuous
Processing is separate
Windowing Engine Rich sliding, Event-time windows Event-time windows
session, and tumbling and watermarks with grace periods
Keyed State Managed keyed state Stateful operators Local RocksDB state
with Embedded use state stores and stores with changelog
RocksDB option checkpointing topics
Exactly-Once Checkpointed state Replayable sources, exactly_once_v2 when
consistency, where checkpointing, and configured for
end-to-end depends idempotent or Kafka-native flows
on sink semantics transactional sinks
Latency Low latency, Workload and trigger Millisecond-scale
event-driven dependent, where event processing is
execution micro-batch can possible for common
trade latency for topologies
throughput
Best Suited For Complex event Unified SQL/data Kafka-centric
processing with processing and applications with
large keyed state existing Spark stack local state and joins
Selection Verdict:
Apache Flink is a strong fit for trending topics because it combines event time windows,
stateful streaming, and low latency execution with scalable keyed state.Personalized Trends (Staff Extension)
Reordering trending topics for individual users risks introducing filter bubbles and complicates regional caching. Personalization should operate as an optional post scoring re-ranking phase applied to the top regional candidates rather than requiring an independent counting pipeline. This design keeps streaming state unified while allowing personalized blending based on user interest embeddings.
Personalized Trends Architecture (Staff Level Extension):
Global and Regional Baseline:
Standard trending produces the top 20 topics for a geographic region, serving identical
rankings to all users in that territory from Redis cache.
Personalization Layer:
Blends broad regional news with individual user engagement history without fragmenting
the underlying stream aggregation pipeline.
Algorithm:
1. Profile Vectors: Offline Spark jobs generate 64-dimensional user interest vectors daily.
2. Topic Embeddings: Trending topics receive category vectors generated from entity analysis.
3. Real-Time Re-ranking:
relevance = cosine_similarity(topic_embedding, user_interest_embedding)
final_score = (w1 * velocity) + (w2 * relevance)
Example weights: w1 = 0.0930769231, w2 = 2.0384615385
Concrete Example:
User follows #Technology and #ArtificialIntelligence.
- Global topic #WorldCup: velocity = 15.0, user relevance = 0.1 -> blended score = 1.6
- Niche topic #GPT5Release: velocity = 5.0, user relevance = 0.9 -> blended score = 2.3
Result: The personalized feed elevates #GPT5Release while retaining #WorldCup in view.
Latency Guardrail:
Re-ranking evaluates only the top 50 regional candidates at query time, with a p99 target of under 10 ms.Trend Category Auto-Classification
Categorizing topics into verticals such as Sports, Politics, and Entertainment enables vertical-specific trend filtering and enriches candidate ranking context without relying solely on manual taxonomy mapping.
Automated Topic Categorization Pipeline:
Objective:
Assign broad taxonomy tags (Sports, Politics, Tech, Entertainment) to facilitate category-specific
trend boards and provide context to downstream recommendation engines.
Multi-Stage Classification Flow:
1. Named Entity Recognition (NER):
Extract recognized proper nouns from post text (for example, "Taylor Swift" -> PERSON).
2. Knowledge Graph Resolution:
Map extracted entities to Wikidata or knowledge graph nodes ("Taylor Swift" -> Musician).
3. Lightweight Topic Modeling:
Run pre-trained embeddings or lightweight BERT classifiers across sample post clusters.
4. Curated Hashtag Taxonomy:
Direct keyword lookups for major recurring tags (for example, #WorldCup -> Sports).
Fallback Handling:
Topics with ambiguous or unclassified signals default to the generic "Trending" category.Multi-Language Trend Merging
Global events spark concurrent discussions across multiple languages and distinct hashtags. Resolving different linguistic variants to a shared entity identifier prevents identical news stories from fragmenting attention and crowding regional leaderboards.
Cross-Lingual Trend Entity Merging:
Challenge:
A global event produces distinct language-specific hashtags:
- Spanish: "#CopaDelMundo"
- English: "#WorldCup"
- Japanese: "#ワールドカップ"
Without Entity Resolution:
The identical sporting event occupies three distinct positions on a global or bilingual leaderboard,
diluting engagement volume and crowding out other newsworthy topics.
With Entity Resolution:
1. Entity Dictionary: An offline dictionary maps multilingual synonyms to a canonical entity ID
("FIFA World Cup 2026").
2. Aggregation Key: Flink partitions and counts by canonical entity_id rather than raw text.
3. Localized Display: The serving layer displays the user's localized text variant while reporting
the aggregated global volume (for example, displaying "2.3M posts").Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.