Trending Topics — a worked solution
Sliding windows + Count-Min Sketch + top-K.
Try it yourself first.
You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.
Open the Trending Topics workspaceThe problem
Build the real-time Trending Topics service that powers Twitter Trends, TikTok "Trending", YouTube Trending, Reddit "Hot", Twitch "Top Streams", or Pinterest's trending pins. Engagement events stream in at scale; the system maintains the top-K most-mentioned entities over multiple sliding time windows (1 m, 5 m, 1 h, 24 h), per locale, and exposes a read API that the product surface polls every app cold-start. The algorithmic core is Count-Min Sketch (Cormode & Muthukrishnan 2005) for the long tail plus HeavyKeeper / Space-Saving (Gong 2018 / Metwally 2005) for the top-K maintainer, wired into Apache Flink's sliding-window keyed-state machinery and the Kappa architecture (Kreps 2014) for replay.
The hard part is not the sketches — it's everything around them: viral hot keys breaking partitioning, watermark stalls under regional Kafka producer lag, abuse-driven brigading inflating fake trends, UI jitter from counter-noise swaps, and the 100 ms read-SLO that has to survive a Redis shard restart. This canonical is what a staff SRE team would actually deploy: 15 components, full long-form internals, transaction-aware emit pipeline, and an on-call ergonomics story.
The reference architecture
What each component is for
- Client
Two distinct shapes share this entry: (a) end-user app surfaces — mobile / web — emitting engagement events (post, view, search, like, dwell-impression sampled at 1–3%) and reading the current top-K via GET /v1/trends?country=US&window=5m; (b) internal services (search-rank, feed, ads) fetching the same top-K to bias their own re-rankers. Both shapes hit the same edge. Read:write skew is ~5:1 by request count but >100:1 by byte volume — every app cold-start fetches a trends payload, every engagement event is tiny.
Why it exists. Considered modelling internal-readers and end-user-readers separately. Rejected because their cache semantics are identical (per-locale, per-window TTL-bounded JSON), and the only divergence is auth (mTLS + JWT for internal vs anon bearer cookie for end-users) which the gateway already handles. Considered splitting engagement-producers (mostly mobile SDK batched POSTs) from trends-readers. Rejected: same observation — the divergence is auth + rate-limit, both downstream.
When it fails. Mobile SDK retry storms (e.g. after a brief gateway 503) can amplify ingest 3–5× — indistinguishable from a viral spike. Detection: distinct_session_ids_per_event_p99 collapses (a few clients producing many events). Mitigation: client-side exp backoff with jitter + per-session token-bucket; reject batches with >100 events at the gateway. The SDK MUST persist unsent events to local storage so a network blip doesn't lose them — but with a hard cap (300 s of buffered events) to prevent a recovered-from-airplane-mode device from dumping 8 hours of stale events into a current window.
- CDN · Edge Trends CacheCloudflare (anycast, 300+ POPs)
Caches the public trends payload at the edge keyed by (country, segment, window) with a 30 s TTL and stale-while-revalidate. Absorbs the bulk of read traffic — trends is the #1 polled endpoint on every app cold-start, so >95% of requests should never hit origin. Edge also runs Bot-Management to drop scraper traffic before it consumes origin capacity, and DDoS-mitigates the read path during news events.
Why it exists. Considered fronting the trend-api with a regional load-balancer + Redis read-replicas. Rejected because the payload is identical across all users for a given (locale, window), so caching at the EDGE (one cached blob serves 10M users from the nearest POP) collapses ~150K peak reads/sec into ~7.5K origin reads/sec. CDN also gives DDoS absorption + Bot-Management for free — trends is a frequent scraper target and the only way to keep origin from being scraped to death is at the edge. YouTube's 30-min refresh cadence and Twitter's WOEID-keyed cache are both edge-cached this way.
When it fails. CDN cache poisoning is the canonical risk — a bad payload propagates to every POP for the TTL window. Mitigation: short TTL (30 s for the hot tier), signed response bodies (HMAC over (locale, window, generation_id)), Surrogate-Key purge wired to the emitter so a bad emit can be retracted in <10 s. POP-level outages handled by anycast failover (within 30 s usually). Detection: edge_hit_ratio_per_country, edge_5xx_rate, origin_egress_bps. Cloudflare June 2022 cross-region routing retro reminds us anycast misconfig can route a country's traffic to a stale POP — we monitor edge_pop_distribution_per_country and page on >5σ deviation.
- API Gateway · Public EdgeEnvoy + WAF + TLS terminator
L7 ingress for /v1/events (engagement ingest) and /v1/trends/* (reads). Terminates TLS, applies WAF (OWASP top-10 + scraper rules), enforces per-tenant + per-IP token-bucket rate limits BEFORE routing so retry storms can't poison the ingest pipeline, calls auth for /v1/events (signed SDK key) and /v1/trends/for-you (user-bearer) but lets anonymous reads through for /v1/trends/place. Routes reads to the trend-api pool; routes writes to the ingest-svc pool; routes admin/op-flip endpoints to a separate admin pool (not modelled here).
Why it exists. Considered inline rate-limit + auth at the trend-api/ingest-svc. Rejected because (1) a misbehaving tenant must be circuit-broken WITHOUT 503'ing everyone — that needs to happen at the EDGE, not at the service tier whose only response to a runaway tenant is more 503s; (2) trend-api would have to ship the rate-limit Redis cluster to every replica; (3) Cloudflare-style WAF + Bot-Management is a separate competency — it doesn't belong inline in a trends-emitter codebase. Stripe's 2019-07-10 retro showed what happens when rate limits live in the wrong layer.
When it fails. Edge pool exhaustion under a coordinated scraper attack is the load-bearing failure — trends is the #1 scraper target and the gateway is the chokepoint. Detection: gateway_active_streams, downstream_5xx_rate, per-tenant rate-limit-hit. Mitigation: Bot-Management challenges, per-IP circuit-breaker, drain-mode + per-zone failover. Cert expiry is the OTHER classical failure (Microsoft Teams 2020-02-03 lesson) — cert-manager + alert at T-30 days. Token-bucket state lives in the gateway's process-local LRU + co-located Envoy ratelimit sidecar (so a downstream cache outage cannot disable the rate limiter on the entry path).
- Auth / IdentitySPIFFE/SPIRE service-mesh JWT + scoped API keys
Validates the Authorization header for protected endpoints. Three flavors: (1) SDK API keys for /v1/events ingestion — hashed lookup against a sharded Postgres keys table + per-key revocation list; (2) user-bearer JWTs for /v1/trends/for-you (personalized trends) signed by the platform identity provider; (3) service-mesh SPIFFE JWTs for internal service calls. Returns a signed (tenant_id, user_id, scopes, allowed_locales) envelope back to the gateway, valid for 60 s. Anonymous reads of /v1/trends/place do NOT go through here — they're the bulk of read traffic and would create a hot path through auth that doesn't earn its keep.
Why it exists. Considered baking auth into the gateway. Rejected because (1) a CVE in the gateway should not break authn for everyone (blast-radius isolation); (2) the keys table is hot in cache and doesn't belong on the trends hot path; (3) scoped API keys (e.g. ingest-only keys for SDKs vs admin keys for the operator console) require independent validation that's awkward to inline. Trends-read latency budget (100 ms p99) cannot absorb a 30 ms auth round-trip on every request — hence anonymous reads bypass auth entirely.
When it fails. Auth down → ingest writes fail-CLOSED 503 (better than fail-OPEN, which would let untrusted events into the top-K); /trends/for-you fails over to the anonymous /trends/place fallback (degraded UX, not outage). Cached scopes survive 60 s then 401. Detection: auth_5xx_rate, scope_cache_miss_rate, jwt_signature_fail_rate. Microsoft Teams 2020-02-03 cert-expiry outage (4 h global) is the cautionary tale — cert-manager + alert at T-30 days. SPIRE node-attestor failures during cluster expansion are the other gotcha — documented in SPIFFE issues.
- Service · Engagement IngestGo + librdkafka idempotent producer + Cloudflare-style geo enrich
Hosts POST /v1/events. For each batch: (1) schema-validate against the registered Avro schema (Confluent Schema Registry) — unknown field → reject with a structured error so the SDK can update; (2) enrich each event with geo (locale, country, region) derived from the producer-IP via MaxMind GeoIP2 in-process; (3) tokenize / normalize entity references (hashtag NFC-fold + lowercase, search-query Unicode-fold, URL-canonicalize); (4) classify the event-type into a topic (engagement / search / impression / dwell); (5) publish to Kafka with the idempotent producer (acks=all, enable.idempotence=true, max.in.flight=5) keyed by the *content* token (hash(entity_id, locale-bucket)) so all events for the same entity land on the same partition for downstream keyed aggregation. Returns 202 in <50 ms; durability is via Kafka's RF=3 sync replication, not via any synchronous downstream write.
Why it exists. Considered fronting Kafka directly with a thin proxy. Rejected because (a) schema validation, geo-enrichment, and tokenization are all stateful concerns the stream-processing tier shouldn't own — every downstream consumer would have to repeat them; (b) the partition-key choice (entity-id-hash so keyed aggregation is correct downstream) is set HERE — if a producer chose its own key, hot-key salting becomes impossible to enforce. The acceptor is the trust boundary between the SDK and the pipeline. Considered Flink's source connector handling enrichment. Rejected because Flink restarts (savepoint restore) replay enrichment for everything since the savepoint, including geo lookups that may have changed — doing the enrichment at INGEST captures the geo at event-time and locks it.
When it fails. Kafka producer unreachable → fail-CLOSED 503 + SDK retries with backoff; lost-event count is bounded by the SDK's local 300 s buffer (see client failure mode). Schema-registry outage → cached schemas survive 60 s; after that, reject with a soft error code so the SDK can bail to a versioned-static schema. Geo-enrichment DB stale → events are misattributed to wrong locale; mitigation is daily auto-refresh + canary against a held-out test set. Detection: kafka_produce_ack_p99 (alert <200 ms P2), schema_registry_5xx_rate, geo_enrich_p99, ingest_per_replica_qps_skew (alert >2.5× mean, indicates partition imbalance). The path NEVER blocks on a downstream service — Kafka durability is the only sync guarantee.
- Stream · Engagement SpineApache Kafka 3.7 (KRaft, 18 brokers, 3 AZ)
Single Kafka cluster carries every engagement event AND every downstream emission through the system. Inbound topics: events.engagement (post-ingest, partitioned 256 ways by hash(entity_id) so keyed aggregation downstream is correct). Emission topics: trending.candidates (Flink top-K emit per (locale, window)), trending.suspect (abuse-detector signals for gating), trending.correction (late-event correction stream that overwrites a published top-K entry), engagement.cold (long-haul archival to ClickHouse). Per-topic retention tuned: events.engagement = 72 h for Kappa replay (Kreps 2014 — a replayable log IS the batch layer); trending.* topics = 24 h. Compacted? No — we want the full log for replay, not the latest value.
Why it exists. Considered SQS / RabbitMQ / Kinesis. Rejected because (a) Kappa replay needs offset-addressable retention which SQS does not offer; (b) multi-consumer fan-out (Flink + abuse-svc + cold-archival all reading the same engagement stream) wants Kafka's pub-sub semantics; (c) 28 MB/s peak ingress is comfortably within a single Kafka cluster (LinkedIn flips multi-cluster at ~2 GB/s per the published numbers). Considered Apache Pulsar. Rejected because the operator-tooling story for Pulsar (BookKeeper + ZooKeeper) is heavier than Kafka KRaft and our SRE team owns Kafka ops at the Heron-paper level. Twitter Heron's break-out-trend topology is fed by Kafka-class logs; same here.
When it fails. Single broker loss → controller re-elects partition leaders, ~5 s gap on the affected partition's writes; min.ISR=2 means single-broker loss never blocks producers. Two-broker loss blocks writes on the affected partitions — correct (better than risking data loss). Hot partition (10× QPS on one entity, e.g. a viral hashtag) is the load-bearing failure: detection is per-partition lag skew >5× mean; mitigation is the application-layer hot-key salt + two-stage aggregation (see ingest-svc and flink). Consumer-group rebalance during deploys: KIP-429 cooperative-sticky assignor + static membership (group.instance.id) so a rolling restart triggers seconds of pause per partition, not minutes. Offset-reset on cluster failover is the OTHER lurking failure — auto.offset.reset=none is mandatory on all consumer groups; manual operator advance via runbook if a reset is unavoidable (Confluent MM2 docs warn about exactly this).
- Service · Flink Trending TopologyApache Flink 1.18 + RocksDB state + DataSketches (CMS, HeavyKeeper)
The algorithmic core. Consumes events.engagement and produces trending.candidates. Topology: (1) Map: parse event, extract (entity_id, locale, ts); (2) Read salted key: the ingest-svc already appended the salt (for entity_ids in the etcd-published HotKeySet) before hashing the Kafka partition key, so a viral key already arrives fanned across 16 sub-shards / source subtasks — Flink consumes salted_entity as-is; (3) KeyBy(salted_entity, locale): partitioned by (salted_entity, locale) so the keyed state lands on a Flink slot; (4) Sliding windows: 4 hopping windows (1m hop / 5m size, 1m hop / 5m size, 5m hop / 1h size, 15m hop / 24h size) with grace=30s for late events; (5) CMS update + HeavyKeeper update: per (locale, window) sketch state in RocksDB, updated per event with conservative-update CMS variant + HeavyKeeper for top-K; (6) Salt-merge: a SECOND KeyBy on (un-salted entity, locale) re-merges the salted sub-counts back into a single global count per entity; (7) Top-K emit: every window-hop fires a ProcessFunction that walks HeavyKeeper, picks top-500, applies EWMA rank-smoothing (α=0.3), and emits to trending.candidates keyed by (locale, window). Watermark: event-time minus 30 s bounded-out-of-orderness; late events route to a side-output that publishes trending.correction events to overwrite the affected (locale, window, entity).
Why it exists. Considered Kafka Streams. Rejected because: (a) Kafka Streams's state store rebalance during deploys is non-trivial at our shard count (256 keyed partitions × 4 tiers → 1024 logical state stores) — Flink's checkpointed savepoints with cooperative rebalance handle this better; (b) DataSketches-on-Flink has a mature Yahoo-published integration; on Kafka Streams it's a roll-your-own. Considered Storm/Heron. Rejected: the Heron paper is the historical precedent but Flink has won the open-source landscape; we own ops at the Flink level, not Heron. Considered RisingWave / Materialize (incremental view maintenance via SQL). Rejected because the abuse-gate + hot-key salt logic is imperative — expressing it in SQL is awkward and the IVM compilers don't yet support custom UDF state on the scale we need. Two-stage aggregation (salt+rekey) is the canonical Flink answer for hot keys; without it ONE viral hashtag pins ONE slot's CPU and stalls watermark advance for everything else on that slot.
When it fails. Watermark stall is the load-bearing failure: a region's Kafka producer batch latency spikes → watermark holds → windows don't tick → top-K goes stale (or under-counts if we drop late events). Mitigation: bounded-out-of-orderness watermark + late-events side-output → trending.correction stream that the emitter overwrites with. Checkpoint storm to S3 → increase retries + jitter prefixes + bound checkpoint duration to 60 s, on fail abandon and try again (don't block the job). Hot-key meltdown without salting: detection is per-slot CPU pin + per-key state-size skew; the hot-key detector job (separate Flink job) emits to etcd within 60 s; the ingest-svc salts that key on the next event and the topology's salted KeyBy fans it across 16 sub-shards. JobManager (Flink master) failover: HA mode with ZooKeeper / KRaft-equivalent + savepoint on the most recent checkpoint, RTO ~30 s. RocksDB write-amplification under viral-spike load = checkpoint size doubles — alert at +100% baseline. JVM GC pause cascade: ZGC + off-heap RocksDB; full-GC pauses >10 s trigger a Kafka rebalance — keep them under 1 s. Detection: watermark_event_time_lag_seconds (P1 >2× grace), checkpoint_duration_seconds, per-slot CPU + RocksDB key count, late_events_per_window.
- Object · Flink Checkpoints + Cold ReplayS3 (us-east-1 + us-west-2, cross-region replication)
Two concerns under one bucket: (1) Flink's RocksDB checkpoint store — each TaskManager async-uploads incremental RocksDB SST files every 10 s, plus a metadata file describing the checkpoint manifest; (2) Long-haul cold archive of events.engagement Kafka log dumped via the engagement.cold consumer to Parquet for >72h-old replay (Kappa horizon extends beyond Kafka retention via S3). Bucket structure: s3://trending-state/checkpoints/{job-id}/chk-{n}/ for Flink, s3://trending-state/engagement-cold/year=Y/month=M/day=D/hour=H/ for Parquet partitions. Cross-region replicated so a region loss doesn't lose the Flink savepoints.
Why it exists. Considered local-disk-only checkpoints (Flink's filesystem state backend). Rejected because a TaskManager-host loss + RocksDB-state-loss means re-bootstrapping from Kafka offset — fine for replay but 72 h of state to rebuild = ~4 h of catch-up at our throughput. Considered EFS / EBS-backed durable checkpoint. Rejected: cross-region durability is what we actually need for the disaster-recovery story (region loss), and S3 gives 11-nines durability cross-region; EBS does not. Considered HDFS. Rejected: another distributed system to operate. The Pinterest Flink retros, Lyft Flink-at-scale talks, and AWS Flink docs all converge on S3 as the production answer — with the caveat that the request-rate ceiling forces prefix-jitter.
When it fails. S3 SlowDown / 503 throttling during a checkpoint burst is the canonical failure: detection is checkpoint_duration_seconds > 60 s + s3_5xx_rate climbing; mitigation is incremental checkpoints + prefix jitter + exponential backoff in the Flink S3 client. Cross-region replication lag >15 min opens a DR-failover gap — if us-east-1 dies and the latest savepoint hasn't replicated, we restore from the second-latest in us-west-2 and re-process the gap from Kafka. Bucket-policy misconfig is the operational risk (a CIDR change locks Flink out from its own checkpoints) — mitigation is IAM-policy CI tests against a held-out 'Flink-reader' role. AWS S3 us-east-1 has had multiple regional events (2017, 2021); cross-region replication is the only defense.
- Service · Abuse / Trend-Manipulation GateFlink + Python user-defined-features + XGBoost classifier
Separate Flink job consuming the SAME events.engagement topic (different consumer-group). For each (entity, locale, 1m-window) candidate that the trending topology produces, computes anomaly features: velocity z-score against a 7-day baseline, account-age median for contributors, ASN/IP-block entropy, novel-account ratio (accounts <7d old), geo-cluster vs platform baseline, content-similarity to known spam templates. Scores via an XGBoost model loaded from S3; high-suspect entities are published to trending.suspect. The trend-emitter checks this stream BEFORE promoting an entity into the public top-K — first-time-trending entities are shadow-listed until the classifier signs off OR a human reviewer approves via the admin console (not modelled).
Why it exists. Considered inlining the abuse classifier into the main Flink topology. Rejected because (a) the classifier is heavy (XGBoost inference 5–20 ms per candidate; at peak that's 50 K classifier calls/sec on the main hot path — blows the watermark budget); (b) the classifier model rolls forward weekly with a separate deploy cadence — coupling its lifecycle to the trending topology forces correlated outages; (c) the abuse-signal store is QUERIED by other surfaces (account-quality system, ban-evasion detection) — it earns its own service boundary. Twitter's published Q1-2018 spam-suspension figure (>142 K applications, >130M spammy tweets) is the precedent; brigading is the #1 trending failure mode by user-reported volume.
When it fails. Abuse-svc down → first-time-trending entities BLOCK at the emitter pending human review; existing trends keep flowing. Detection: trending_first_time_promotion_lag_p99 (P2 alert >10 min). Classifier-model drift → false-positive rate climbs, real trends fail to promote; mitigation is the weekly canary against a held-out validation set + a manual-override admin endpoint. Adversarial: a sophisticated actor pre-warms accounts to defeat account-age features — mitigation is the ensemble (no single feature is load-bearing). Detection: false_positive_rate_24h, classifier_skew_drift, suspect_rate_per_locale (high = brigade campaign). Twitter's published 'state of the platform' integrity reports document this is an arms race; the gate is a tax, not a solution.
- Service · Top-K EmitterGo + Redis client + idempotent ZADD pipeline
Consumes three Kafka topics in lock-step: trending.candidates (the raw top-K from Flink), trending.suspect (abuse signals), trending.correction (late-event corrections). For each (locale, window) candidate emission: (1) join against trending.suspect — entities currently flagged are filtered out OR demoted; (2) apply hysteresis — a new entrant must beat the bottom-of-list by ≥0.15× AND sustain for ≥2 consecutive windows; (3) apply EWMA rank smoothing on top of Flink's already-smoothed score — second-order smoothing because the user-visible jitter is what product complains about; (4) write the smoothed top-50 list to Redis as a versioned key trending:{locale}:{window}:{generation_id} (atomically swap via SET); (5) issue a CDN surrogate-key purge for the affected (locale, window) so the edge re-fetches. Backfills to the analytics store happen via a separate cold-archival consumer, not here.
Why it exists. Considered making the trend-api do the abuse-join + hysteresis on every read. Rejected because (a) the hysteresis state (which entity was on the list LAST window) is per-(locale, window) and would have to live in a separate state store the trend-api reads on every request — blowing the 100 ms read SLO; (b) abuse-suspect joins are 200 ms-ish XGBoost classifier scores, doable async at emit-time but not sync on read; (c) decoupling EMIT-time work from READ-time work is the core architectural decision: the read path is a dumb Redis ZRANGE, the emit path bears the algorithmic complexity. Considered making the Flink topology own this. Rejected because the abuse-suspect stream and the candidate stream are produced by separate Flink jobs with separate watermarks; the join lives outside Flink to keep the algorithmic core focused.
When it fails. FALLING-BEHIND-GLOBALLY is the load-bearing failure: trend-emitter is the only writer to the read path, so global lag = global staleness. Detection: emitter_global_p99_emit_age (P1 >90s — beats per-locale alerts which lag a global outage by minutes). Cache write failure (Redis cluster down) → emitter fail-CLOSED on primary-down (better to serve stale than wrong); fallback is dual-write to a secondary Redis cluster during cutover. Suspect-consume lag while candidate-consume is healthy → the join uses a stale suspect set; a real brigade campaign goes unflagged. Detection: per-topic consumer-group lag (NOT aggregated). CDN purge half-failure: emitter logs-and-moves-on after Cloudflare returns 503 → permanently stale until the next emit cycle naturally overwrites; defense is the generation-id check at the CDN edge (rejects payloads older than X seconds AS WELL AS validates HMAC). Bad deploy that corrupts hysteresis state across all locales → canary-by-locale-shard catches it before 12.5% of locales drift. Correction stream straggles → monotonic event-time generation-id check refuses the stale correction. Hysteresis trade documented: a real breaking trend takes one window longer to appear. Detection: top_k_swap_rate_per_minute, emitter_redis_write_p99, cdn_purge_5xx_rate, emitter_global_p99_emit_age, suspect_consume_lag_seconds, candidate_consume_lag_seconds (split signals).
- Cache · Top-K (Redis ZSET)Redis 7 cluster mode, 6 shards
Authoritative-for-the-read-path top-K. Each (locale, window) is a Redis SORTED SET keyed trending:{locale}:{window} containing (entity_id, score) pairs; the trend-api hits ZREVRANGE 0 49 for a top-50 fetch — sub-millisecond. Also stores a versioned blob trending:{locale}:{window}:{generation_id} (JSON serialized top-500 with metadata) that the emitter atomically SETs and the trend-api fetches as a single-RTT 'give me everything for this locale-window' call. Eviction-aware: we explicitly DON'T let LRU evict trending lists — they live forever (TTL=0) and the emitter overwrites in place. Auxiliary keys: trending:hot-key-set mirrors etcd's hot-key set for the ingest-svc and Flink to read; trending:abuse-pending mirrors the human-review queue.
Why it exists. Considered serving directly from ClickHouse via materialized views. Rejected because ClickHouse single-row p99 on a top-K read at our cardinality is 30–100 ms (OLAP doesn't excel at point-query OLTP), and we need <5 ms cache reads to fit the 100 ms read budget. Considered an in-process cache in the trend-api. Rejected because cache hit ratio depends on cross-replica sharing — with 24 trend-api replicas, an in-process cache gives 1/N hit ratio. Centralizing in Redis means ONE warm copy serves all replicas; cluster failover handled by Redis Sentinel / cluster-mode without coordination from the API layer. Pinot at LinkedIn supports this query shape (250K QPS at ms latency) and we considered it — rejected because Redis is operationally simpler at this footprint (we don't need Pinot's columnar power for a 24 MB working set).
When it fails. Cluster cold restart → 30 s of trend-api degraded reads (falls through to ClickHouse); the AOF persistence keeps RPO under 5 s of emitter writes lost (acceptable — the next emit overwrites). Single shard down: the locales hashing to that shard go cold (1/6 ≈ 17% of users); mitigation is per-locale fallback to ClickHouse at trend-api (see the cache-miss alternate path). Replication lag spike on async: a follower-read could serve a 5 s-stale list — mitigation is reading from PRIMARY (we use RYW consistency, not eventual). Hot shard from US-en: split US into US-east + US-west sub-locales if read-QPS-per-shard climbs >50 K/sec (we're at ~7.5 K, plenty of headroom). Detection: cluster_state, shard_hit_rate, primary_5xx_rate, redis_aof_rewrite_lag.
- Service · Trends Read APIGo + Redis client + ClickHouse fallback driver
Hosts GET /v1/trends/place?country=US&window=5m and GET /v1/trends/for-you (for logged-in personalized). For each request: (1) resolve the locale (country code → internal locale-id via in-process LRU); (2) ZREVRANGE the (locale, window) ZSET from Redis OR fetch the versioned blob (one round-trip); (3) on cache miss / Redis 5xx, fall through to ClickHouse for the most-recent materialized top-K row; (4) for /for-you, additionally fetch the personalization re-ranker scores from the (not-modelled) ml-serving and re-rank — with a circuit-breaker fallback to the un-personalized list. Response: JSON top-50 with (entity, score_normalized, generation_id, window_end_ts). HMAC-sign the response body so the CDN can validate at edge.
Why it exists. Considered serving directly from Flink's queryable state. Rejected because Flink's QueryableState is a debug feature, not production-strength — it doesn't have the failover, the request-coalescing, or the connection limits a production API needs. Considered making the CDN's origin a Lambda that reads Redis. Rejected because cold-start tail of 200–800 ms blows the 100 ms read budget. A stateless service tier is the textbook answer; the only question is whether we even need it given Redis directly. Answer: we DO, because (a) ClickHouse fallback logic is non-trivial; (b) the personalization re-rank for /for-you needs a service tier; (c) HMAC response-signing needs a key the CDN can't have.
When it fails. Redis primary down on a shard → read falls through to ClickHouse for 1/6 of locales; ClickHouse single-row OLTP p99 is 30–100 ms vs Redis 1–3 ms, so SLO degrades but doesn't break. Detection: redis_5xx_rate, clickhouse_fallback_qps, p99 climb. ClickHouse also down → returns a stale-but-signed 'last-known-good' from an in-process 5-min LRU; UI degrades gracefully (yesterday's trends > blank UI). Personalization service down (for /for-you only) → circuit-breaker opens, returns un-personalized /trends/place response (degraded UX, not outage). HMAC signing key rotation: rolled forward via a 24 h dual-key window; if the rotation bombs, the CDN starts refusing payloads — mitigation is the canary on every rollout. Bad deploy: 5xx-rate-SLO burn-budget alert + auto-rollback via Argo Rollouts.
- Analytics DB · ClickHouseClickHouse 23.8 (6 shards × 2 replicas, ZooKeeper coordination)
Source of truth for HISTORICAL trends (per-locale, per-window, per-entity time-series) and the fallback read path when Redis is down. Materialized views ingest trending.candidates and write to two tables: trends_history (every emit recorded for offline analysis — 'what was trending in Tokyo at 14:00 UTC last Tuesday') and trends_latest (a ReplacingMergeTree keyed by (locale, window) that holds the most-recent emit for fallback reads). Also a engagement_cold MergeTree fed by the cold-archival consumer for offline reprocessing, abuse forensics, and ML training. Cold tier (>90 d) tiered to S3 via ClickHouse's storage policy.
Why it exists. Considered Druid / Pinot. Both work; Pinot's per-row low-latency reads (LinkedIn's 250K QPS at ms) are arguably a better fit for the fallback read path than ClickHouse's OLAP-tilted single-row latency. Chose ClickHouse because (a) our analyst-tooling story is built around it (Looker + Cube.dev + Grafana plugins), (b) the operator-tooling for ClickHouse Keeper / Zookeeper-based clusters at our scale is well-understood by the team, (c) MergeTree's ingest-throughput characteristics outshine Druid's segment-handoff at our 28 MB/s engagement-cold ingest rate. Considered Postgres + Citus. Rejected because at 1B events/day × 90 d × 200 B = 18 TB hot, sharded Postgres works but query-throughput on the analyst dashboard side (long time-range scans across all locales) is materially worse than ClickHouse.
When it fails. ZooKeeper / Keeper ensemble loss → ClickHouse cluster goes read-only (writes from the Kafka MV stall); detection is keeper_quorum_state P1. Single shard down: 1/6 of locales lose analyst queries; the fallback path at trend-api routes to a different shard if Redis is also down (rare cascading failure). Slow MERGE backlog: at high ingest, parts accumulate faster than they merge, query latency climbs — mitigation is right-sized merge threads + monitoring parts_to_throw_insert. Disk-full on the hot tier: alert at 70%, page at 85%, auto-tier-down to S3 at 90%. Datadog March 2023 outage applied to OLAP: storage-tier collapse cascades hard — keep at least 30% headroom in steady state.
- Coordinator · Config + Hot-Key Saltsetcd 3.5 (5-node Raft, 3 AZ)
Linearizable config store holding (a) the HOT-KEY SET — entity_ids the ingest-svc salts (and Flink reads to un-salt and merge); this is updated by a separate hot-key-detector job at ~60 s cadence; (b) the WINDOW-DEFINITION table — named windows ('5m', '1h', '24h') with their hop / size / grace parameters; rolling out a new window tier is a config push, not a code deploy; (c) the ABUSE-SUSPECT THRESHOLD table for the abuse-svc classifier; (d) the per-tenant HOME-REGION LEASE table — each tenant's authoritative region with a 30 s TTL + 10 s heartbeat (the fencing token that prevents both regions writing the same trend). The trend-emitter and trend-api check the home-region lease before doing any per-tenant work.
Why it exists. Considered using DynamoDB Global Tables (which we'd use elsewhere for cross-region dedupe). Rejected because (a) the home-region LEASE is the FENCING token — it MUST have linearizable consistency, and Global Tables's documented up-to-30 s replication lag is too long; (b) etcd / Consul give linearizable leases via Raft, matching the Cloudflare June 2022 cross-region fencing pattern. Considered ZooKeeper. Rejected because we already operate ZK for ClickHouse; running another ZK ensemble for config blurs ownership — etcd is the team-preferred coordination store and gRPC API is friendlier than ZK's. Considered Consul (similar Raft-based). Rejected: tied to HashiCorp's release cadence; etcd is the Kubernetes-native default and our SRE team owns it.
When it fails. etcd quorum lost (3 of 5 nodes failed) → writes fail-CLOSED, but READS continue serving the last-known-good cached values to all consumers (Flink, ingest-svc, trend-emitter, trend-api). The system continues running on the CACHED state; new hot-keys are not salted; new lease changes are not honored. Detection: etcd_has_leader, etcd_proposals_committed_total. Cross-region partition: each region honors its OWN lease state — risk is both regions believing they own a tenant if the cross-region sync drifts; mitigation is the dispatcher-level in-region check + a daily cross-region consistency audit job. The Cloudflare June 2022 retro showed what happens when the fencing primitive ITSELF gets misrouted — we monitor active_lease_region_per_tenant and page on multi-region overlaps.
- Metrics · Prometheus + GrafanaVictoriaMetrics (3-node TSDB cluster) + Grafana + Alertmanager
Scrapes every service for RED/USE metrics on a 15 s interval, evaluates SLO burn-rate alert rules, and serves the analytic-query backend for Grafana dashboards. Argo Rollouts polls Prometheus's /api/v1/query for the SLO-burn alert that drives auto-rollback during a bad deploy. Holds the recording-rules that pre-aggregate per-(locale, window) emit-cadence + freshness-lag SLOs so the on-call dashboards don't scan raw cardinality.
Why it exists. Considered relying on tracing alone. Rejected because (a) Argo Rollouts auto-rollback runs off a metric query, not a span query — the metrics layer is the canonical control loop primitive; (b) cardinality storms (deep-dive #2 in failure-modes) are a metrics-tier problem; tracing doesn't have the same explosion shape; (c) the on-call dashboards (Read Health, Flink Topology, Kafka Cluster, Emitter Lag, Abuse Gate) all read PromQL. Tracing complements but cannot replace metrics. Considered Datadog. Rejected on cost at our ingest cardinality. Considered direct InfluxDB. Rejected: VictoriaMetrics has a more favorable memory-per-active-series footprint at our scale.
When it fails. Self-monitoring deadlock: if the metrics tier itself is degraded, the alert that would page on-call about the metrics tier doesn't fire. Mitigation: out-of-band dead-man's-switch (a synthetic that pings VictoriaMetrics every 60 s from a separate region; if absent for >5 min, an SMS-only path pages — bypasses Alertmanager entirely). Cardinality explosion from a bad label deploy is the on-call's worst day — relabel rules MUST run at Prometheus, not at the agent, so a misbehaving scraper can't smuggle in high-cardinality data. Datadog 2023-03-08 retro is the precedent for what happens when the observability tier itself collapses. Argo Rollouts losing access to Prom → deploys can't auto-rollback; the canary-pause workflow falls through to manual SRE approval.
- Service · Schema RegistryConfluent Schema Registry (Avro + Protobuf)
Hosts the canonical Avro schemas for events.engagement and the Protobuf schemas for the internal trending.* topics. Producers (ingest-svc) and consumers (Flink, abuse-svc, trend-emitter, ClickHouse Kafka MV) fetch + cache schemas at boot; on a schema evolution, a registered backward-compatible new version goes live without redeploying consumers. Enforces compatibility-mode FULL_TRANSITIVE so a bad schema bump can't be registered without breaking the build.
Why it exists. Considered embedding schemas in service code with a CI compat-check. Rejected because (a) producers + consumers deploy on different cadences — a producer that ships a new field MUST wait for consumers to upgrade BEFORE flipping; the registry enforces this with subject-level compat history; (b) Kafka Connect / ClickHouse Kafka MV / Flink all integrate natively with the Confluent Registry — rolling our own would be more work. Considered Apicurio. Rejected: similar feature set, less ecosystem momentum at our scale.
When it fails. Registry unreachable → producers + consumers serve from cached schemas for 60 s, then hard-fail (deserialization exception → consumer lag climbs). Mitigation: long cache TTL (60 s buys the SRE a window to fix it) + soft-reject on producers (reject the request with a 503 + Retry-After rather than publishing an unparseable message). Schema-evolution-deploy gone bad (registered an incompatible schema after disabling the compat check) is the load-bearing operator error — mitigation is RBAC restricting WRITE to the registry to a small ops group + an audit log alert on any compat-check-disabled push. Detection: schema_registry_5xx_rate, schema_cache_miss_rate, deserialization_fail_rate at consumers.
- Service · Hot-Key DetectorApache Flink (separate job from the main topology)
Separate Flink job consuming events.engagement with a 1-min tumbling window. For each (locale, entity), computes per-window count + per-window-vs-previous-window growth rate. Entities that exceed a published threshold (>5× growth AND >5 K events/min OR >50 K events/min absolute) get pushed to etcd as /trending/hot-key-set/{locale}/{entity_id} with a 5 min TTL. Refreshes the set every 60 s. The ingest-svc watches this etcd key with a long-poll and applies hot-key salting on the next event for any flagged entity; the main Flink topology reads the same set to un-salt and merge.
Why it exists. Considered embedding hot-key detection in the main Flink topology. Rejected because (a) the detection logic operates on a DIFFERENT keying (per-locale, per-entity, growth-rate signal) than the main topology's (salted-entity-per-window aggregation) — running both in one job means awkward state-split; (b) detector deploy cadence is independent of the main topology; (c) blast radius isolation — a bug in the detector that emits too many hot-keys to etcd salts everything unnecessarily and inflates state 16×, but it doesn't take the main topology down. Considered a heuristic threshold in ingest-svc. Rejected: ingest-svc is stateless and needs to fast-fail on schema; growth-rate detection requires keyed state which doesn't belong there.
When it fails. Detector down → no NEW hot keys are salted; existing flagged keys persist until their 5-min TTL expires. A flash-viral key (new in the past 5 min) goes UNSALTED → main Flink topology pins one slot. Detection: hot_key_set_age_per_region_seconds (P2 >120s, P1 >300s), hot_key_detector_alive. Cross-region detector clusters: each region runs its own detector since the input stream is per-home-region; mitigation for cross-region staleness is documented in config.failureMode.
- Cache · Exact-Count Cross-Check StoreRedis 7 cluster mode, 3 shards
Holds EXACT per-(entity, locale, window) counts for entities that CMS estimates above a candidate threshold (currently 10 K events/window). Used by the trend-emitter at promote-time to verify that a CMS-promoted top-K entry has a real underlying count, defeating CMS over-estimate (deep-dive #1 — Cormode-Muthukrishnan's noted bias). Updated by the main Flink topology via a side-output every window-emit: 'these k entries crossed the threshold; here is the exact-count for cross-check.'
Why it exists. CMS NEVER under-counts but ALWAYS over-counts on hash-collision; at our 50 M-entity cardinality, the over-count bound on a long-tail entity is 20 K events/window (ε × N), which is enough to promote a never-actually-trending entity. HeavyKeeper improves over plain CMS but is still probabilistic. The production answer (referenced in the deep-dive but not previously drawn) is to keep an exact-count store for the small set of entities above the candidate threshold — a few thousand keys per locale, fits in 3 Redis shards.
When it fails. Exact-count store down → trend-emitter falls back to CMS-only promotion → over-counted long-tail items risk leaking into top-K (the failure mode this node was specifically added to prevent). Mitigation: if the cross-check store is unreachable for >5 min, the emitter SHADOW-LISTS all new-to-top-K entries until the cross-check returns (defense-in-depth with the abuse-gate which is also a promote-gate). Detection: exact_count_5xx_rate, cms_to_exact_count_drift (offline daily job — large drift = either CMS is over-counting or the threshold is mis-tuned). Single-shard down: 1/3 of entities cross-check-fails; emitter shadow-lists those for the duration.
- Tracing · OTel + TempoOpenTelemetry Collector + Grafana Tempo
Receives OTLP-encoded spans from every service (ingest-svc, trend-api, trend-emitter, abuse-svc) and the Flink JobManager. Sampled at 1% head-based for steady-state plus 100% for anything that hits an SLO-burn alert (tail-based via OTel's tail-sampler). Spans are stored in Tempo with traceID-indexing on S3-backed storage. Provides the watermark-lag trace view that on-call uses to answer 'why did the 5m window for US-en take 90 s to tick' — the load-bearing debug surface for a streaming pipeline.
Why it exists. Considered Jaeger. Rejected because Tempo's S3-backed storage scales cheaper at our cardinality without manual sharding; the operator-tooling story is Grafana-native (matches our existing dashboard infra). Considered Honeycomb. Rejected on cost at the per-event ingest rate. Considered NO tracing (relying on logs + metrics). Rejected because watermark lag is the hardest-to-debug failure mode and traces are the only tool that show the cross-service causal chain — 'event X published at T, consumed by Flink at T+5s, window emitted at T+30s' — in one view.
When it fails. Tempo backend unreachable → spans drop after the local buffer fills (64 MB ≈ 30 s of traffic); detection is otel_dropped_spans_total. Critically, tracing is NOT on the synchronous response path — user latency is decoupled from tracing availability (the lesson from 'observability outage during the incident' anti-patterns). Self-monitoring: a separate Prometheus + Loki stack monitors the tracing tier so a Tempo outage doesn't blind us to its own degradation. Span buffer overflow under load → spans drop but app continues — the explicit trade we make.
Stage by stage
The same 10 stages the workspace walks, answered.
01Clarifications
What would you ask before drawing a single box?
Typical clarifications:
- What's "trendable"? Hashtags, search queries, watched videos, viewed products — the entity model varies by surface. Affects tokenization, abuse signals, and cardinality.
- What windows? Single 5-min window vs. multi-resolution (1 m / 5 m / 1 h / 24 h). Multi-resolution adds state cost but is the production default.
- Per-locale? Global only is a toy; production is per-country/per-city (Twitter ships ~200 WOEID markets). Affects partitioning + cache size.
- Personalized? "Trending for you" sits on top of per-locale candidates. Adds a re-ranker path; the canonical here handles candidate generation, not personalization.
- Refresh cadence? 30 s for the 5 m window is product-facing; analyst dashboards can tolerate 1–5 min.
- Freshness SLO? "An event from X seconds ago must appear in the top-K." Typically 60 s for the 5 m tier.
- Abuse model? Trends drives virality — brigading is the #1 product complaint. Need a gate before promote.
- Read SLO? Sub-100 ms p99 — trends sits on app cold-start.
- At-least-once OK? Yes — a lost event biases the trend by <0.0001 %; durability isn't the prime concern. (Contrast notifications, payments.)
Assumptions stated for this canonical:
- 1 B trendable events/day (~11.6 K eps avg, 70 K sustained peak, 300 K spike).
- 500 M trends-read requests/day (≈18.5 K avg, 150 K peak — absorbed mostly at CDN).
- 200 locales (countries + cities); US-en holds ~60 % of traffic (power-law).
- Multi-resolution windows: 1 m, 5 m, 1 h, 24 h. Surfaced top-50; tracked top-500 for hysteresis.
- At-least-once ingest; lost events bias <0.1 % in CMS — within accuracy budget.
02Functional reqs
What must this system actually do?
- POST /v1/events — internal services + SDK push engagement events; idempotent batches up to 100.
- GET /v1/trends/place?country=XX&window=Wm — anonymous read of per-locale top-K.
- GET /v1/trends/for-you — authenticated personalized re-rank of per-locale top-K.
- (Internal) admin endpoints — force-flip home-region lease, manual override of abuse-suspect, hot-key salt set push.
- (Internal) replay endpoint — spin up a parallel Flink job from a chosen Kafka offset for bug-fix re-emit.
03Non-functional
What must it promise about speed, uptime and correctness?
- Read availability: 99.95 % (~4.4 h/yr error budget). Trends absent is a visible UX hit but not a data-loss event.
- Read latency: p99 < 100 ms end-to-end (CDN hit), p99 < 200 ms on cache miss to ClickHouse fallback.
- Freshness: an engagement from 60 s ago must appear in the 5 m top-K. (Watermark lag budget: 30 s; emit cadence: 30 s.)
- Refresh cadence: 30 s for 5 m window, 1 min for 1 h, 5 min for 24 h.
- Durability: at-least-once; Kafka RF=3 sync (min.ISR=2) is sufficient. Exactly-once is NOT a goal.
- Consistency: trends are per-locale-eventually-consistent across regions; explicit HOME-REGION fencing prevents cross-region double-emit.
- Security: anon read OK; ingest requires signed SDK key; admin requires user JWT + 2FA. Bot-Management at the edge; abuse-gate before promote.
- Multi-region: active-active for reads (per-locale HOME region); standby for writes (single region per tenant). RTO 20 min for region cutover (Flink savepoint restore + MM2 catch-up dominate; see multi-region section), RPO 5 min (Kafka MM2 worst-case 30 s + Flink/emitter catch-up). Datadog 2023-03-08 four-hour recovery is the realistic worst-case anchor.
04Capacity estimation
How much load and data does this have to hold?
Anchoring on a Twitter/X-scale workload (citations in research notes). Numbers are derived from the defaults set in the frontmatter; tweak the capacity sliders to see how the design moves.
Ingest pipeline.
- Events: 1 B/day → 11.6 K eps avg. Peak multiplier 6× → 70 K eps sustained, 8× spike multiplier → 560 K eps spike ceiling (Twitter Trends-tier published envelope; Super Bowl / TikTok-challenge burst).
- Avg payload 400 B → Kafka ingress 28 MB/s at peak (70 K × 400 B; the 168 MB/s figure that previously appeared in this section was double-counting the peak multiplier — flagged and corrected in critic round 1). Well below the LinkedIn-published 2 GB/s multi-cluster flip point → single Kafka cluster comfortably.
- 256 partitions on events.engagement sized for the 8× spike + salt-shard fan-out scenario (~2.2 K eps/partition at spike); at steady peak ~270 eps/partition. The load-bearing pressure is HOT KEYS, not aggregate ingress.
Stream-processing state.
- CMS: width 2048, depth 5, 4 B counters → 40 KB / sketch. 200 locales × 4 tiers × 16 salt-shards = 12,800 sketches ≈ 512 MB CMS state across the fleet.
- CMS accuracy caveat: ε = 2/w ≈ 0.001. Per 5 m window at peak, 70 K eps × 300 s = 21 M events → CMS over-count bound per entity ≈ 20 K. For a long-tail entity with true count 1, CMS can estimate up to ~20 K. CMS is therefore only a pre-filter; HeavyKeeper plus the exact-count-store cross-check are the source of truth for promotion to user-visible top-K.
- HeavyKeeper top-K: 500 entries / sketch × ~24 B = 12 KB / sketch → ~150 MB total.
- Exact-count store cross-check: ~1 K candidates above threshold per locale × 200 × 4 windows × 40 B ≈ 32 MB → 3 Redis shards.
- Sketch memory (CMS + HeavyKeeper) is only ~700 MB fleet-wide — small enough to stay in-heap; the full per-slot keyed state (sketches + hopping-window buffers + timers) targets ~2 GB / slot in RocksDB on local NVMe, well under the 50 GB spill threshold.
- Throughput rule: this topology is deeper than the published 30 K eps / slot Pinterest figure (parse + enrich-check + hot-key-lookup + 2× KeyBy + CMS-update + HeavyKeeper-update + EWMA + emit). Realistic per-slot is 10–15 K eps. At 560 K spike / 12 K = ~47 slots required; provision 64 slots = 16 TM × 4 plus 1 standby TM. Steady-state (70 K eps / 12 K eps/slot) burns 6 slots — plenty of headroom for rolling restart.
Read path.
- 500 M reads/day → 18.5 K avg / 150 K peak. CDN @ 95 % hit → 7.5 K origin reads/s → single Redis cluster, 12 nodes total.
- Top-K cache: 200 locales × 4 windows × 500 entries × 120 B ≈ 48 MB working set → trivial.
- Egress at edge: 150 K × 8 KB ≈ 1.2 GB/s — absorbed by CDN, not origin.
Storage.
- Kafka 72 h retention × 28 MB/s × RF=3 ≈ 22 TB on disk.
- ClickHouse hot tier 90 d × 200 B/row compressed ≈ 18 TB hot, 5.3 PB/yr cold (tiered to S3).
- S3 Flink checkpoints ≈ 1.3 TB (10 retained × 2 GB/slot × 64 slots).
Scale flip points.
| Threshold | Architectural change |
|---|---|
| 2 GB/s Kafka ingress | MirrorMaker 2 fan-out, multi-cluster |
| 500 GB keyed state / TM | Per-locale Flink jobs (sub-partition by locale) |
| 50 GB heap / slot | Spill from heap CMS to RocksDB-backed CMS |
| 50 K writes/s/Redis-shard | Split US-en into US-east + US-west sub-locales |
| 5 PB ClickHouse hot | Add per-region ClickHouse shards |
05API design
What does the outside world call, and what comes back?
POST /v1/events
Authorization: Bearer sk_live_xxx
Content-Type: application/json
{
"events": [
{
"ts": 1714003200,
"event_type": "search_submit",
"entity_id": "#OscarsSlap",
"entity_kind": "hashtag",
"user_id_hash": "abc...",
"session_id": "..."
},
...
]
}
202 Accepted { "received": 100, "rejected": 0 }
Ingest is at-least-once — there is no Idempotency-Key contract. A re-submitted batch (SDK retry after a 202 is lost on the wire) is re-counted; the resulting double-count is absorbed within the <0.1% CMS accuracy budget (enable.idempotence=true only dedupes producer→broker retries, not client→gateway resubmits).
GET /v1/trends/place?country=US&window=5m
Accept-Language: en-US
200 OK
Cache-Control: public, max-age=30, stale-while-revalidate=60
X-Trends-Generation-Id: 1714003200-us-5m-9bf2
{
"country": "US",
"window": "5m",
"window_end_ts": 1714003200,
"generation_id": "...",
"trends": [
{ "entity": "#OscarsSlap", "score_normalized": 1.0, "delta_rank": 0 },
...
]
}
GET /v1/trends/for-you
Authorization: Bearer <user JWT>
200 OK { "country": "US", "personalized": true, "trends": [...] }06Data model
What gets stored, and what is it looked up by?
Kafka topics (engagement spine):
| topic | partitions | retention | key | notes |
|---|---|---|---|---|
| events.engagement | 256 | 72 h | hash(entity_id) | source of truth, Kappa replay |
| trending.candidates | 64 | 24 h | (locale, window) | Flink emit |
| trending.suspect | 32 | 24 h | (locale, entity) | abuse-svc output |
| trending.correction | 32 | 24 h, compacted | (locale, window, entity) | late-event overwrite |
| engagement.cold | 64 | 24 h | hash(entity_id) | mirror for ClickHouse MV |
Redis trending cache:
trending:{locale}:{window}— Sorted Set, member=entity_id, score=normalized_count.trending:{locale}:{window}:{generation_id}— JSON blob, atomic versioned-set by emitter.trending:hot-key-set— mirror of etcd hot-key salt set.
ClickHouse trends_history (ReplicatedMergeTree):
| column | type | notes |
|---|---|---|
| locale | LowCardinality(String) | PK component |
| window | Enum8 | 1m/5m/1h/24h |
| window_end_ts | DateTime | PK component |
| entity_id | String | |
| count | UInt64 | smoothed |
| rank | UInt16 | post-hysteresis |
| generation_id | String | for replay alignment |
Partition by (locale, toYYYYMMDD(window_end_ts)); ORDER BY (locale, window, window_end_ts, rank). Tier >90 d to S3 via ClickHouse storage policy.
etcd config (selected keys):
/trending/hot-key-set/{locale}/{entity_id}— 5 min TTL, set by hot-key-detector job./trending/window-defs/{name}— hop / size / grace./trending/lease/{tenant_id}→ region_id, 30 s TTL, 10 s heartbeat.
07High-level design
Which components handle a request, and in what order?
Architecture summary (19 nodes, 37 edges — the diagram is the source of truth):
- Client (mobile / web SDK + internal readers) emits engagement events and reads trends.
- CDN (Cloudflare) absorbs ~95 % of trends reads at the edge with 30 s TTL + stale-while-revalidate + Bot-Management.
- API Gateway (Envoy + WAF) does per-tenant rate-limit + WAF + routes reads to trend-api, writes to ingest-svc.
- Auth validates SDK keys (writes) and user JWTs (personalized reads). Anonymous reads bypass.
- Engagement Ingest (ingest-svc) schema-validates, geo-enriches, tokenizes, and publishes to Kafka events.engagement with idempotent producer + acks=all.
- Schema Registry — Confluent Schema Registry holding canonical Avro/Protobuf schemas + forward-compat checks.
- Kafka (engagement spine) — the Kappa source-of-truth. 256 partitions, 72 h retention, RF=3 sync.
- Flink Trending Topology — the algorithmic core. Two-stage agg (hot-key salt → KeyBy → CMS+HeavyKeeper → un-salt merge) over hopping windows; emits to trending.candidates.
- Hot-Key Detector — separate Flink job; 1-min tumbling windows + growth-rate signal → publishes salt set to etcd.
- Exact-Count Store (Redis) — exact per-(entity, locale, window) counts above a candidate threshold; defeats CMS over-estimate (deep-dive 1).
- S3 Checkpoints — Flink RocksDB snapshots every 10 s + cold engagement archive for >72 h replay.
- Abuse Service — separate Flink job consuming the same engagement stream; XGBoost classifier publishes trending.suspect.
- Top-K Emitter (trend-emitter) — joins candidates + suspect + correction streams (separately tracked consumer-group lag per topic); hysteresis + EWMA + abuse-gate; writes versioned top-K to Redis; issues CDN surrogate-key purge. Part of the
emit-outboxtxn group. - Redis Top-K Cache (trend-cache) — ZSET per (locale, window). 6 shards, 12 nodes.
- Trends Read API (trend-api) — GET handlers; Redis ZREVRANGE; ClickHouse fallback; HMAC-sign responses.
- ClickHouse (trend-store) — history + analytics + cache fallback. 6 shards, 12 nodes.
- etcd (config) — hot-key salt set + window defs + home-region lease (linearizable fencing token).
- Metrics (VictoriaMetrics + Grafana + Alertmanager) — SLI store + Argo Rollouts source-of-truth for auto-rollback.
- OTel + Tempo (tracing) — watermark-lag debugging, tail-sampled spans.
Data flow on read (CDN hit): client → CDN → response.
Data flow on read (CDN miss): client → CDN → gw → trend-api → trend-cache → response (→ origin cache populates CDN).
Data flow on ingest: client → gw → ingest-svc → kafka. The 202 returns once Kafka acks RF=3 sync; downstream is decoupled.
Data flow on the streaming pipeline: kafka.engagement → flink (windowed CMS+HeavyKeeper, hot-key-salted) → kafka.trending.candidates. In parallel: kafka.engagement → abuse-svc → kafka.trending.suspect.
Data flow on emit: kafka.{candidates,suspect,correction} → trend-emitter (hysteresis + abuse-gate) → trend-cache (ZADD) + cdn (surrogate-key purge).
08Deep dives
Which part breaks first, and what do you do about it?
(See deepDivePrompts in the frontmatter for the prompts the canonical tries to answer. Below are the load-bearing answers.)
1. CMS over-estimate of cold items. Count-Min Sketch is biased: a hash collision between a long-tail item and a hot one inflates the cold item's estimated count and can promote it into the top-K. Fix: conservative update (Estan-Varghese 2002) reduces but doesn't eliminate; HeavyKeeper (Gong 2018) does this better by replacing count-all with exponential-decay admission — RedisBloom migrated TOPK to HeavyKeeper for exactly this reason. Defense-in-depth: keep an exact-count store (Redis SortedSet of entities with count >threshold) for the top-1 % candidate set so the final emit cross-checks CMS estimates against ground truth.
2. Hot-key meltdown. A viral hashtag takes 80 % of a Kafka partition's traffic; the Flink slot owning that partition pins one CPU. Without intervention, watermark advance stalls for ALL keys on that slot. Fix: two-stage aggregation. A separate detector job (consuming the same engagement stream) emits hot-keys to etcd; the ingest-svc and Flink job read this set every 60 s; for hot keys, the ingest-svc appends a random salt in [0..15] to the partition key, spreading the load across 16 sub-shards; a second KeyBy in Flink un-salts and merges. The trade: 60 s detection lag means a flash-viral key gets one window of slot-pinning before salting kicks in — acceptable.
3. Watermark lag and late events. A region's Kafka producer batch latency spikes; Flink's bounded-out-of-orderness watermark holds; windows don't tick. Fix: 30 s grace + late-events side output. Late events (event-time < watermark) route to a side stream that publishes trending.correction events; the trend-emitter applies corrections by overwriting the (locale, window, entity) entry IF the correction's generation_id is monotonically newer. Bound the correction window by the cache TTL (30 s for the 5 m tier) so corrections don't surface as user-visible flicker.
4. Top-K jitter. Two phrases within counter-noise swap every refresh. Fix: hysteresis — challenger must beat incumbent by ≥0.15× AND sustain for ≥2 consecutive windows. Combined with EWMA rank smoothing (α=0.3) so a one-shot spike doesn't displace a sustained trend. The trade: a real breakout takes one window longer to appear; acceptable for UI stability.
5. Abuse / brigading. Trends drives virality, so brigading is self-reinforcing. Fix: gate at PROMOTE-time, not at INGEST-time. The abuse-svc runs in parallel on the same engagement stream, scoring candidates on velocity z-score, account-age, ASN entropy, novel-account ratio. First-time-trending entities are SHADOW-listed until the classifier OR a human reviewer signs off. Twitter's published Q1-2018 spam-suspension numbers (>142 K applications) are the precedent; the gate is an arms race, not a solved problem.
6. Kappa replay. A bug mis-tokenized 6 h of events; the top-K is wrong. Fix: spin up a parallel Flink job (new consumer-group id) from the corresponding earliest offset; let it catch up; promote its emit topic to the trend-emitter; drain the old job. The Kappa promise (Kreps 2014) is exactly this — a replayable log IS your batch layer. No batch-vs-stream code split, no double-pipeline maintenance.
09Trade-offs
What did this design cost, and what breaks at 10×?
| We chose | Alternative we rejected | Why |
|---|---|---|
| Kafka + Flink (Kappa) | Lambda (batch + stream) | One pipeline; reprocess by offset replay (Kreps 2014). |
| CMS + HeavyKeeper | Exact-count everywhere | 50M-cardinality × 200 locales × 4 windows in exact RAM is ~40 GB+; sketches collapse to ~500 MB. |
| Two-stage agg (salt+rekey) | Single-stage agg | One viral key would pin one slot; salt-fan-out + un-salt merge defeats it. |
| Redis ZSET per (locale, window) | Pinot / Druid | At a 48 MB working set, Pinot's columnar power is overkill; Redis is operationally simpler. |
| Cloudflare edge cache 30 s | No CDN, origin-served | 150 K read/s peak → collapses to 7.5 K origin; CDN is mandatory at this scale. |
| ClickHouse as fallback | Postgres + Citus fallback | Long time-range scans for analyst dashboards perform 10× better in ClickHouse. |
| etcd for hot-key set + lease | DynamoDB Global Tables | Lease is the fencing token; needs linearizable, not 30 s replication lag. |
| OTel + Tempo tracing | Honeycomb / no tracing | Watermark-lag debugging needs cross-service span traces; the cost is bounded. |
| Outbox txn on ingest → Kafka | 2PC across ingest+downstream | 2PC across Kafka + ClickHouse is a non-starter; outbox via the Kafka idempotent producer is the documented pattern. |
| At-least-once | Exactly-once | Exactly-once across Kafka + downstream stores is illusory; we dedupe at the emitter via generation-id. |
| Per-locale HOME region (active-active reads, active-standby writes) | Active-active writes | Avoids dueling-emitter writes when cross-region partitions occur (Cloudflare Jun 2022 retro). |
Primary sources
- Cormode & Muthukrishnan — An Improved Data Stream Summary (CMS, J.Alg. 2005)
- Metwally, Agrawal, El Abbadi — Space-Saving top-K (ICDT 2005)
- Gong et al. — HeavyKeeper (USENIX ATC 2018)
- Kulkarni et al. — Twitter Heron: Stream Processing at Scale (SIGMOD 2015)
- Apache Flink — Stateful Stream Processing + Watermarks docs
- Confluent — Windowing in Kafka Streams + KIP-633 grace periods
- Confluent — KIP-794 Strictly Uniform Sticky Partitioner
- Confluent — KIP-429 Cooperative-Sticky Rebalance
- Jay Kreps — Questioning the Lambda Architecture (Kappa)
- Boykin, Ritchie, Singhal — Summingbird / Algebird (Twitter)
- Apache DataSketches — Theta sketch (Yahoo)
- LinkedIn — Open-sourcing Apache Pinot
- Discord — How Discord Stores Trillions of Messages
- Datadog — March 8 2023 Multi-region outage retro
- Cloudflare — June 21 2022 cross-region routing retro
- Twitter — How we're fighting spam and malicious automation (2018)
- Pinterest — TransAct real-time user actions for Homefeed
- Beyer et al. — SRE Workbook (Managing Load + Cascading Failures)
Now defend it
Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.
Work Trending Topics yourselfMore in Feeds, Timelines, Counters & Ranking
What to show and in what order: fanout on write versus read, hot/top/new scoring, approximate counters, trending, and recommendation.
- Twitter / X TimelinePush or pull? Both. The canonical fanout problem.
- Instagram News FeedRanked feed with cursor pagination. No `OFFSET`.
- Reddit / Hacker NewsVote-driven ranking with hot/top/new at scale.
- Like Button at ScaleEventual consistency, but the liker sees their own write. Counts are approximate by design; hot keys are the real enemy.
- View Count on a Video/PostDedup, bot-filter, batched aggregation.