Live Viewer Count (YouTube/Twitch)
Worked solution

Live Viewer Count (YouTube/Twitch) — a worked solution

Millions of viewers on one entity. Approximate by design.

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 Live Viewer Count (YouTube/Twitch) workspace

The problem

The "1.4M watching now" counter on a live stream is the canonical approximate counting problem at internet scale. The number doesn't have to be exact — it has to be believable, monotone-feeling, fast-enough to feel live, and robust to coordinated viewbot attacks. Production teams at YouTube, Twitch, Meta, JioHotstar, and Discord have converged on a remarkably consistent shape: heartbeat ingest at MMevents/sec, two-stage HyperLogLog aggregation against hot keys, per-region sketch-merge into a designated primary region, monotone-display smoothing at the publish boundary, and hierarchical WebSocket fanout to push the displayed number back to every subscribed player. Every component has a real production failure mode; every knob has a story.

The reference architecture

Reference architecture for Live Viewer Count (YouTube/Twitch): 23 components — Player (Web / Mobile / Smart-TV), CDN · Edge Count Cache, API Gateway · Public Edge, Auth · Stream-Session-Token Service, WebSocket Hub · Realtime Count Fanout, Service · Heartbeat Ingest, Stream · Heartbeat Spine + Emit Topics, Service · Flink Viewer-Count Topology, Cache · HLL Sketch Bag (per-region registers), Service · View-Fraud Detector, Service · Counter Promotion + WS Publisher, Cache · Counter Read Path (Redis), Service · Counter Read API, Analytics DB · Counter History (ClickHouse), Object · Flink Checkpoints + Cold HB Archive, Coordinator · Config + Hot-Stream Set, Tracing · OTel Collector + Jaeger, Metrics · VictoriaMetrics + Argo Rollouts SoT, Service · Hot-Stream Detector, Service · Schema Registry, Service · Stream Lifecycle Controller, Service · WS Fanout Hub (tier-2), Service · Panic-Mode Controller — connected by 43 flows.Player (Web / Mobile …clientCDN · Edge Count CacheCloudflare + Akamai + F…API Gateway · Public …Envoy + WAF + TLS termi…Auth · Stream-Session…Go + SPIFFE service-mes…WebSocket Hub · Realt…Go + nhooyr/websocket +…Service · Heartbeat I…Go + librdkafka idempot…Stream · Heartbeat Sp…Apache Kafka 3.7 (KRaft…Service · Flink Viewe…Apache Flink 1.18 + Roc…Cache · HLL Sketch Ba…Redis 7 cluster mode, 1…Service · View-Fraud …Apache Flink (separate …Service · Counter Pro…Go + Redis pipeline + K…Cache · Counter Read …Redis 7 cluster mode, 8…Service · Counter Rea…Go + Redis client + Cli…Analytics DB · Counte…ClickHouse 23.8 (8 shar…Object · Flink Checkp…S3 (cross-region replic…Coordinator · Config …etcd 3.5 (5-node Raft, …Tracing · OTel Collec…OpenTelemetry Collector…Metrics · VictoriaMet…VictoriaMetrics (3-node…Service · Hot-Stream …Flink (lightweight side…Service · Schema Regi…Confluent Schema Regist…Service · Stream Life…Go + Kafka producer + a…Service · WS Fanout H…Go + gRPC-streaming + p…Service · Panic-Mode …Go + PromQL client (met…
23 components, 43 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Player (Web / Mobile / Smart-TV)

Three coupled responsibilities share one client: (a) the VIDEO player — terminating HLS / LL-HLS / DASH chunk requests at the CDN; (b) the HEARTBEAT producer — emitting a /v1/heartbeat POST every 30 s in steady state, 5 s during scheduled ad-breaks or DRM-quality-shift windows, with the session-bound stream-token in the body; (c) the COUNT subscriber — opening a long-lived WSS connection to ws-hub at stream-start (subscribe topic viewers:{stream_id}) AND issuing a one-shot HTTP GET /v1/streams/:id/viewers on cold-load to grab the current number before the first WS push arrives. Display logic is local: receives WS push or HTTP fetch, applies the platform's monotone-smoothing + milestone-pinning rules in-player, paints the count.

Why it exists. Considered splitting heartbeat-producers, count-readers, and count-subscribers into distinct services / endpoints. Rejected because the SAME player session is producing AND consuming — same auth context, same session token, same in-process state. Considered a server-rendered count painted into the player HTML (no client logic). Rejected because the count changes every 30 s and 30 s is too slow if the user is staring at it — WS push is what moves the platform from 'count refresh' to 'count tick'. YouTube's per-stream concurrent ticks updates faster than the page would on its own; Twitch's player-embedded count likewise comes via PubSub frames, not a re-render.

When it fails. Player-side retry storm is the canonical client failure — a brief gateway 503 → SDK without backoff → 100M players retry at the same second = the AWS Kinesis 2020-11-25 pattern. Mitigation: hard client retry budget + decorrelated jitter + the SDK MUST be most resilient on this path (SRE Workbook ch. 'Addressing Cascading Failures'). Player-side cache poisoning of the count payload (over-aggressive HTTP cache, never refreshes) → mitigation is short Cache-Control (5 s) AND the WS push is authoritative; HTTP poll is a fallback only. Stale stream-token after long-pause-in-tab: SDK refreshes the token on player-resume; HB without a valid token returns 401 + retry-once-with-fresh-token (NOT retry-forever).

CDN · Edge Count CacheCloudflare + Akamai + Fastly (multi-CDN with sub-250 ms steering)

Caches the public count payload at the edge keyed by (stream_id, window=current) with a 5 s TTL and stale-while-revalidate. Absorbs the cold-load HTTP fetch — every player that joins a stream issues exactly one GET /v1/streams/:id/viewers before the WS push arrives, so the edge collapses 100M cold-load fetches into ~75K origin fetches per stream-start spike (assuming 1 % cache-miss). Also runs Bot-Management — viewbot infrastructure scrapes the count endpoint to validate their inflation worked, so this is a frequent abuse target. Multi-CDN steering (Akamai + CloudFront + Cloudflare) gives sub-250 ms failover when one provider degrades (the JioHotstar playbook).

Why it exists. Considered serving the count directly from count-cache (Redis) via the count-read-api with no edge layer. Rejected because the count is identical per (stream, window) across all viewers in a region, so caching at the EDGE collapses ~3M cold-fetch QPS into ~30K origin reads, AND because the count endpoint is the #1 viewbot/scraper target — without edge Bot-Management we'd be running OWASP rules in the origin. Considered making the count a payload of the HLS manifest (segment metadata) — saves the round-trip. Rejected because the manifest cache TTL is 4–6 s (matches segment cadence) and changes the cache-invalidation story dramatically; the dedicated count endpoint with surrogate-key purge gives us the explicit invalidate primitive emitters need.

When it fails. CDN serves a stale count after stream-end is the load-bearing UX failure: stream ended at t=0, player joined at t=4 s, gets a cached '1.4M watching' for the 5 s TTL, then suddenly sees 'stream ended' — brand crisis for the creator. Mitigation: surrogate-key purge wired to the emitter so stream-end fires an immediate invalidate; the player /live endpoint is authoritative, count payload contains an ended:true flag once detected. POP-level outage handled by multi-CDN steering (sub-250 ms). Cache poisoning of count payload: signed response body (HMAC over (stream_id, count, generation_id, ts)) + CDN rejects unsigned / older-than-X payloads at the edge. Cloudflare 2022-06-21 cross-region routing retro reminds us anycast misconfig can route a country to a stale POP — per-POP count divergence alert at >5σ.

API Gateway · Public EdgeEnvoy + WAF + TLS terminator + per-IP token-bucket

L7 ingress for POST /v1/heartbeat (heartbeat ingest), GET /v1/streams/:id/viewers (count read), and the WS upgrade for wss:// to ws-hub. Terminates TLS, applies WAF (OWASP top-10 + scraper rules + JA3 fingerprint matching against the viewbot block-list), enforces per-IP + per-ASN + per-session token-buckets BEFORE routing so a retry storm or coordinated viewbot campaign cannot poison the ingest pipeline, validates the session-stream-token for /v1/heartbeat by calling out to auth-svc (cached 60 s), routes heartbeats to the ingest-svc pool, routes reads to count-read-api, routes WS upgrades to ws-hub. WATCHES etcd's panic-mode/{cohort} keyspace so the heartbeat-response header carries the panic-mode flag to clients within one heartbeat-interval. Routes admin / stream-lifecycle endpoints (RTMP-publish, stream-end) to a separate admin pool not modelled here.

Why it exists. Considered inline rate-limit + auth at the ingest-svc and ws-hub. Rejected because (a) a misbehaving tenant or a viewbot campaign must be circuit-broken WITHOUT 503'ing legitimate traffic — that needs to happen at the EDGE, not at services whose only response is more 503s; (b) WS upgrade-time auth at ws-hub means every WS connect carries an auth round-trip even for legitimate traffic — the gateway caches the auth result and rejects the WS upgrade BEFORE it touches ws-hub; (c) the WAF + Bot-Management surface is a separate competency from the count pipeline. Twitch's 'Breaking the Monolith' post showed what happens when rate limits live in the wrong layer — one channel's raid 503'd everyone.

When it fails. Edge pool exhaustion under a coordinated viewbot campaign IS the load-bearing failure — the heartbeat endpoint is the #1 abuse target after the count endpoint. Detection: gateway_active_streams, downstream_5xx_rate, per-IP/ASN rate-limit-hit, JA3 entropy collapse per stream. Mitigation: aggressive Bot-Management challenges, per-IP circuit-breaker, drain-mode + per-zone failover. Cert expiry is the OTHER classical failure (Microsoft Teams 2020-02-03) — cert-manager + alert at T-30 days. Auth-svc 503 → gateway falls back to a LRU snapshot of recently-valid tokens (300 s grace); after that, fail-CLOSED for new sessions. WS upgrade auth: every reconnect storm puts pressure on auth — mitigation is upgrade-token-cache + connection rate-limit per session-id.

Auth · Stream-Session-Token ServiceGo + SPIFFE service-mesh JWT + scoped session-tokens + Redis token cache

The trust boundary for the heartbeat protocol. Three flavors: (1) /play-start issues a session-stream-token bound to (viewer_session_id, stream_id, ip_prefix, ja3_hash, attestation_blob), short-lived (15 min), signed Ed25519. The PLAYER receives it and includes it in every subsequent /v1/heartbeat AND in the wss:// subscribe upgrade. (2) Per-HB validation: gateway calls /validate with the token; auth-svc verifies signature, expiry, ja3-match-against-stream-context, ASN reputation; returns a signed envelope (viewer_id_anonymized, stream_id, tier, weight) cached at gateway for 60 s. (3) Service-mesh SPIFFE JWTs for internal RPC. Anonymous count reads of /v1/streams/:id/viewers do NOT pass 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 the token mint + validate into the gateway. Rejected because (a) the validation logic (JA3 match + attestation check + per-ASN reputation) is non-trivial AND mutates over time as we learn new viewbot fingerprints — coupling it to the gateway makes deploy of an abuse-rule require a gateway restart, which is unsafe at this fan-out; (b) the token mint requires private-key custody, which we want sealed inside an HSM-fronted service, not on every gateway pod. Considered token-less heartbeats (just session cookie). Rejected because that's exactly the path the Twitch 2025 viewbot crackdown closed — without session-bound stream-tokens, a bot can forge arbitrary (viewer, stream) pairs at the protocol layer. Twitch's 2025-08-21 deploy dropped public concurrent viewership 24% — the load-bearing nature of this defense is empirically confirmed.

When it fails. Auth-svc down → heartbeat ingest fails-CLOSED 503 (better than fail-OPEN which would let viewbots in unfettered); /play-start fails over to a degraded token (anonymous, narrow-scope) so cold viewers can still join (no count contribution; degraded count, not outage). Cached scopes survive 60 s then 401. Detection: auth_5xx_rate, token_validation_latency_p99, ja3_match_failure_rate, revocation_propagation_lag. Microsoft Teams 2020-02-03 cert-expiry (4h global) is the cautionary tale — cert-manager + alert at T-30 days. Key-rotation bombing the validation: dual-key window 24 h, canary every rollout. Adversarial: a sophisticated bot mints fresh tokens at scale — defense is per-IP rate-limit on /play-start AND JA3 entropy collapse triggers manual review per (ASN, stream).

WebSocket Hub · Realtime Count FanoutGo + nhooyr/websocket + per-stream topic fanout tree (Twitch-PubSub-style)

The realtime fanout fabric, internally a THREE-TIER tree (Twitch-PubSub topology — confirmed by Twitch Engineering 2022 'Breaking the Monolith' and Discord's 2017 Manifold post). Three concerns split across the tiers: Tier-1 publish-front (~12 pods) receives count-update events from the emitter and fans them out to Tier-2; Tier-2 mid-fanout (~100 pods) holds the (stream_id → Tier-3 subset) routing table and replicates each event to ~30 Tier-3 pods per stream; Tier-3 edge (3000 pods × 75K conns = 225M conn capacity) terminates wss:// connections from players, hosts per-(stream_id) topics, and pushes the new count to every subscriber of the affected stream topic. The replicas: 3000 field models the Tier-3 edge count (the load-bearing capacity dimension); Tier-1 and Tier-2 counts are documented under keyChoices/internals so the diagram doesn't fragment three tightly-coupled processes into separate components. Per-connection state minimal (subscribe-list + last-push-ts + per-second send-quota); zero persistent state on pod loss (player reconnects, re-subscribes, gets a snapshot).

Why it exists. Considered HTTP long-poll instead of WS. Rejected because (a) at 200M concurrent viewers polling every 30 s = 6.7M poll-QPS on the gateway — order of magnitude more expensive than maintaining a TCP socket; (b) push-latency on long-poll is dominated by the poll interval, defeating the 30 s freshness budget. Considered server-sent events (SSE). Rejected because mobile networks aggressively close idle long-lived HTTP connections (carrier NATs) — WS pings keep the connection warm in a way SSE doesn't. Considered HTTP/2 push. Rejected because the player needs to control subscribe/unsubscribe topology, and HTTP/2 push has no concept of topic subscription. The Twitch PubSub (≤10 conns/IP, ≤50 topics/conn) + Discord Manifold (batched send/2 per remote-NODE) + Slack PS/GS (presence + delivery split) + JioHotstar (50K–100K conns/pod, 400–600M scorecard updates in 800 ms) production proof points all converge.

When it fails. Slow-client cascade is the load-bearing failure: a mobile player on a degraded link can't drain its send buffer, the WS pod's memory climbs, OOM-kill takes 75K connections with it. Detection: ws_conn_buffer_high_water + per-pod memory + per-pod conn-drop-rate. Mitigation: per-conn buffer cap + aggressive disconnect on backpressure (drop the slow client; the next reconnect gets a snapshot). Pod loss → 75K players reconnect within 30 s — anti-thundering-herd: client uses decorrelated-jitter exp backoff on reconnect AND the gateway rate-limits WS upgrades per session-id. WS reconnect storm × auth degraded cascade: a regional flap closes millions of sockets; every reconnecting client triggers the WS upgrade auth-call (e25) — if auth-svc is at 50% error the upgrade pool saturates and the reconnect storm self-amplifies. Defense (the SRE Workbook ch. 'Addressing Cascading Failures' play): gateway-validated upgrade-token-cache (60 s TTL) so the ws-hub→auth call (e25) is hit only on COLD subscribe; ws-hub-side upgrade rate-limit per session-id (≤1 upgrade per 5 s) caps the reconnect-driven auth load; on auth-svc circuit-breaker open, ws-hub falls back to gateway-signed envelopes for a 60 s grace window before failing closed. Cross-pod fanout-tree partition: emitter publishes to Tier-2 hubs; if a Tier-2 hub goes down, ~150 Tier-3 pods (its subtree) lose pushes — mitigation is Tier-2 redundancy (each Tier-3 subscribes to 2 Tier-2 hubs) and clients seeing stale-count alert fires within 60 s. Discord's 2017 'how we scaled to 5M' post documents the BEAM-VM equivalent of these failure modes; the fix (Manifold per-NODE batching, FastGlobal hot-path lookups) is the playbook.

Service · Heartbeat IngestGo + librdkafka idempotent producer + MaxMind geo + tier-based sampler

Hosts POST /v1/heartbeat. For each HB: (1) validate the session-stream-token envelope from gateway (already validated, just unpack); (2) tier-based sampler — reads stream-tier from etcd; for tier-3 (<1K viewers) keep 100%, for tier-2 (<100K) keep 10%, for tier-1 (>100K) keep 1% (rationale: variance on the count for a low-count stream is dominated by sample noise so we keep all HBs, but for a 60M-viewer stream 1% gives ±0.81% HLL accuracy on a 600K-sample HLL); (3) hot-stream-salt: if stream_id is in the etcd hot-stream-set, append a random salt[0..15] to the partition key — spreads viral keys across 16 sub-shards; (4) enrich (geo, ASN, CDN POP, tls_ja3 from the gateway envelope); (5) publish to Kafka with the idempotent producer (acks=all, enable.idempotence=true) keyed by hash(stream_id ‖ salt) mod 1024. Returns 202 in <30 ms; durability is via Kafka's RF=3 sync replication, NOT via any synchronous downstream write — heartbeats are deliberately lossy (RPO=60 s).

Why it exists. Considered fronting Kafka directly with a thin proxy. Rejected because (a) the tier-based sampler, hot-key salt, and per-stream enrichment are stateful concerns the stream-processing tier shouldn't own — every downstream consumer would have to repeat them; (b) the partition-key choice (with salt for hot keys) is set HERE — if a producer chose its own key, two-stage aggregation downstream becomes impossible to enforce. Considered Flink's source connector handling sampling. Rejected because Flink restarts (savepoint restore) replay sampling for everything since the savepoint, including HBs that shouldn't have been sampled in the first place — doing the sampling at INGEST captures the decision at event-time and locks it. Considered direct edge-pre-aggregation (compute partial HLLs at the edge pop). Rejected for v1 because it requires Cloudflare-Workers / Lambda@Edge tier code which is operationally heavier; revisit at the 10× scale tier flip.

When it fails. Kafka producer unreachable → fail-CLOSED 503 + SDK retries with backoff; lost-HB count bounded by the SDK 300 s buffer. Schema-registry outage → cached schemas survive 60 s, after that reject with a soft code. Geo-enrichment DB stale → events misattributed to wrong region; daily auto-refresh + canary against held-out test set. The path NEVER blocks on a downstream service — Kafka durability is the only sync guarantee. Detection: kafka_produce_ack_p99 (alert <200 ms P2), schema_registry_5xx, geo_enrich_p99, ingest_per_replica_qps_skew (alert >2.5× mean — indicates partition imbalance OR hot-stream-salt not propagating). Hot-stream-salt-lag = the load-bearing failure: a viral stream goes live, the hot-stream-set hasn't been published to etcd yet → 60 s of UNSALTED HBs slamming one partition → the partition's broker melts. Mitigation: per-partition lag alert + an aggressive hot-key detector consuming the same Kafka stream that publishes to etcd within 30 s. Sampler-tier mismatch (stream is tier-1 but ingest still has stale tier-3 config) → over-counts because we sampled 100% when we should have sampled 1%; mitigation is etcd-watch and a guard on per-stream HB-rate that detects this within 60 s.

Stream · Heartbeat Spine + Emit TopicsApache Kafka 3.7 (KRaft, 36 brokers, 3 AZ)

Single Kafka cluster carries the entire pipeline. Inbound topic: heartbeat.raw (post-ingest, post-sampling, partitioned 1024 ways by hash(stream_id‖salt)). Emit topics: viewers.candidates (Flink window-tick emit per (stream, window)), viewers.suspect (abuse-svc fraud signals), viewers.correction (late-event corrections, log-compacted), viewers.rollup (1 s rollups for ClickHouse cold archive), viewers.lifecycle (stream-start / stream-end signals). Retention: heartbeat.raw 24 h (for Kappa replay if pipeline needs reprocess); viewers.* topics 24 h; viewers.rollup 72 h (ClickHouse loader has time to catch up). NOT log-compacted on heartbeat.raw (we want the full log for replay); viewers.correction IS log-compacted (latest correction per (stream, hour) wins).

Why it exists. Considered SQS / RabbitMQ / Kinesis. Rejected because (a) Kappa replay needs offset-addressable retention which SQS doesn't offer; (b) multi-consumer fanout (Flink + abuse-svc + cold-archive-loader all reading the same heartbeat stream) wants Kafka's pub-sub semantics; (c) at 6.67M HB/s × 380 B = 2.5 GB/s peak ingest (post-sampling: 570 MB/s) we're well within single-Kafka-cluster comfort (LinkedIn flips multi-cluster at ~2 GB/s × RF=3 = 6 GB/s broker-write). Considered Kinesis — actually Twitch DOES use Kinesis for QoUX, but their 'Spade' data lake (S3 + analytics) is the primary; we follow the Kafka + S3 archive pattern because operational tooling (MirrorMaker 2, Kafka Connect, Schema Registry) is mature in our SRE org. AWS Kinesis 2020-11-25 retro showed front-end thread-limit cascading failures — Kafka's per-broker concurrency model decouples this.

When it fails. Exactly-once-effective via idempotent producer (enable.idempotence=true, acks=all) plus emitter-side generation-id dedup — NOT Kafka KIP-98 transactions (which we do not use; their write-amplification and consumer-isolation cost is too high for a 6.67M HB/s spine). Hot partition during stream-start IS the load-bearing failure: viral go-live, all HBs hash to one Kafka partition before the hot-stream-salt-set propagates → that partition's broker pins one CPU. Detection: per-partition Kafka lag skew >5× mean, per-broker CPU asymmetry. Mitigation: aggressive hot-key detector job (separate Flink job consuming heartbeat.raw) emits to etcd within 30 s; partitioning by hash(stream_id ‖ salt) is the structural fix; cell-isolated dedicated cluster for top-100 (tier-1) streams pre-warmed on stream-start. Single-broker loss → controller re-elects partition leaders, ~5 s gap on affected partitions; min.ISR=2 means single-broker loss never blocks producers. Consumer-group rebalance during deploys: KIP-429 cooperative-sticky + static membership (group.instance.id) → seconds-of-pause per partition not minutes. Offset-reset on cluster failover: auto.offset.reset=none MANDATORY; manual operator advance via runbook. Confluent MM2 docs warn about cluster failover offset translation.

Service · Flink Viewer-Count TopologyApache Flink 1.18 + RocksDB state + DataSketches (HLL + Theta) + per-region windowed agg

The algorithmic core. Consumes heartbeat.raw and produces viewers.candidates + viewers.rollup. Topology: (1) Map: parse HB, extract (stream_id, salt, viewer_id_anonymized, region, ts, weight); (2) KeyBy(stream_id, salt, region, window): partition by the FULL salted key so per-region state lands on a Flink slot whose RocksDB key — AND the corresponding sketch-bag Redis HLL key viewers:hll:{stream}:{salt}:{region}:{window} — both spread across 16 sub-shards on a hot stream; (3) Sliding windows: 1 s tumbling for the live count + 60 s sliding for the smoothed published number, grace=10 s for late events; (4) HLL update: per (stream, salt, region, window) HyperLogLog sketch in RocksDB, PFADD viewer_id_anonymized per event; (5) Per-region per-salt emit: every 1 s a ProcessFunction emits the per-salt HLL to viewers.candidates keyed by (stream, region, salt). The salt-merge (PFMERGE across the 16 salts) is performed at the EMITTER, not in Flink, so the salted keyspace stays sharded all the way to the Redis HLL layer and never collapses onto a single hot key. Watermark: event-time minus 10 s bounded out-of-orderness; late events route to side-output → viewers.correction.

Why it exists. Considered Kafka Streams. Rejected because (a) Kafka Streams' state store rebalance during deploys is non-trivial at our 1024 keyed partitions × 4 window tiers → 4096 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 Spark Streaming. Rejected: micro-batch latency tail of 100–500 ms vs Flink's per-record processing; for a 1 s window cadence this matters. Considered counting directly in Redis (PFADD per heartbeat from ingest-svc). Rejected because at 6.67M PFADD/s across hot streams Redis single-key hits the 50K-ops/s ceiling AND because the per-region accuracy logic (sampler weighting, watermark, late-arrival) is correctness-critical to express as a Redis script. Two-stage aggregation (salt+unsalt+rekey) is the canonical Flink answer for hot keys; without it ONE viral stream pins ONE slot and stalls watermark advance for everyone 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 → count goes stale (or under-counts if we drop late events). Mitigation: bounded out-of-orderness watermark + late-events side-output → viewers.correction stream that the emitter overwrites with. Checkpoint storm to S3 → S3 SlowDown 503 → checkpoint duration climbs → back-pressure cascades → count freezes for minutes. Detection: checkpoint_duration_seconds_p99 (alert >60 s P1), state-size growth, consumer lag. Mitigation: incremental RocksDB checkpoints + prefix-jitter S3 keys + exp-backoff in Flink S3 client + bound checkpoint duration to 90 s then abandon. Hot-key meltdown without salting: per-slot CPU pin + per-key state-size skew; hot-key detector (separate Flink job) emits to etcd within 30 s; the topology re-salts on next event for that key. Viewer-facing degraded mode during the 30–60 s pre-salt window of a viral go-live: the affected stream's emit may stall on the one Flink slot pinned by the unsalted hot key; the displayed count is frozen at last-good and is non-decreasing per the monotone-display rule — viewers see a number that is 1–60 s stale, not wrong. The count resumes a normal tick once salt-fanout activates and the 16 sub-shards redistribute the load. This is the documented graceful-degrade SLO; we accept 1–60 s of count staleness in exchange for not pre-salting every stream (which would inflate state 16× for the 99.99% of streams that never go viral). JVM GC pause cascade: ZGC + off-heap RocksDB; full-GC >10 s triggers Kafka rebalance — keep them <1 s. DLQ poison message — malformed HB poisoning the job, restarting every 10 min (the Twitch 2019 NA outage shape): schema registry + reject-on-deserialize → DLQ; never block the hot path on a single bad record.

Cache · HLL Sketch Bag (per-region registers)Redis 7 cluster mode, 12 shards (HLL register storage)

Authoritative-for-the-emit-path HLL register storage. Each (stream, region, window) is a Redis HLL key viewers:hll:{stream}:{region}:{window} populated by Flink via PFADD. The emitter reads these registers via PFMERGE across (region, window) tiers to produce the global count per stream. Two distinct workloads share one cluster: (a) FLINK writes — high-frequency PFADDs across millions of stream-region keys (the count workload); (b) EMITTER reads — periodic PFMERGE across 5 regional sketches per (stream, window-tier) to compute the published count. Auxiliary keys: viewers:hot-streams mirrors etcd's hot-stream-set for ingest-svc and flink; viewers:weighting mirrors the per-tier sampler weights so the emitter can compute weighted-cardinality.

Why it exists. Considered keeping HLL state ONLY inside Flink's RocksDB (no Redis). Rejected because (a) the emitter is a separate service from Flink — coupling the emitter to Flink's QueryableState ties their deploy cycles and savepoint compatibility; (b) the cross-region merge happens at the EMITTER which is in a designated primary region — sending it 5 regional Flink jobs' QueryableState introduces a fan-out that Redis's PFMERGE primitive obviates. Considered Druid Theta sketches. Rejected because Druid is excellent for OLAP query-time merge but heavy for the per-second-per-stream PFMERGE we need (Druid query latency is 100ms+, not the 5 ms we have). Considered Cassandra with counter columns. Rejected: counter columns are notoriously bad on hot keys; Redis HLL is the production-proven sketch primitive (Redis Labs benchmarks). LinkedIn moved their distinct-counters off counter columns onto sketches.

When it fails. Write-side SLOs (this is a write-heavy register store, NOT a cache-hit-ratio workload): PFADD p99 < 10 ms, PFMERGE p99 < 50 ms, AOF-rewrite-lag < 30 s. Async replication means HLL register writes are SEQUENTIAL per-shard (NOT linearizable cluster-wide); PFADD is CRDT-mergeable (same viewer_id → no-op) so the 5 s RPO is rebuilt by the next 10 s of PFADDs in the common case — see the stream-start undercount caveat below. Cluster cold restart → 30 s of count-staleness (emitter falls through to last-known-good); AOF keeps RPO under 5 s of PFADDs lost. BUT HLL undercount during stream-start is PERMANENT for the new-distinct-viewers in the RPO window (bounded by new_distinct_viewers_in_5s — at viral go-live with 1M viewers/sec joining that's up to 5M undercount on the affected stream's first-window count; the HLL does NOT 'rebuild' those distincts because the same viewer IDs heartbeat again with the same hash → already counted). Mitigation for tier-1 streams: opt-in sync replication during the first 60 s of stream-start (cost: ~3× write latency, acceptable on the small tier-1 set). Single shard down: the (stream, region) tuples hashing to that shard go cold → 1/12 of streams' count-merge starves → mitigation is per-stream fallback to the most-recent ClickHouse rollup at the emitter. Hot shard from US-en (one Mr Beast stream): split by region-sub-bucket if read-QPS-per-shard climbs >50K/sec. PFMERGE-fan-out: a global stream merges across 5 regions × 4 window tiers = 20 PFMERGEs per emit cycle; if any region's shard is slow, the WHOLE emit is slow — mitigation is timeout-per-region + exclude-on-timeout (better to publish a count missing one region than block the publish). Detection: cluster_state, per-shard PFADD/PFMERGE latency, primary_5xx_rate, AOF rewrite lag. Per-shard cluster-rebalance during slot migration: AOF + replication mean no data loss, but write-availability dips for ~5 s — schedule rebalances during low-tier-1 hours.

Service · View-Fraud DetectorApache Flink (separate job) + Python UDFs + XGBoost classifier + TLS-JA3 ML model

Separate Flink job consuming the SAME heartbeat.raw topic (different consumer-group). For each (stream, 60 s window) candidate that the trending topology produces, computes anomaly features: TLS-JA3 entropy (a coordinated bot fleet uses the same fingerprint → entropy collapses); ASN/IP-prefix concentration (bot farms cluster on cheap-cloud ASNs); heartbeat cadence z-score (real players' cadence has a natural jitter, bots are too clock-perfect); attestation-blob validity rate (mobile devices have Play Integrity / DeviceCheck attestation; web players use a custom challenge); device-renewal rate (real viewers' sessions persist; bots burn through sessions). Scores via XGBoost; high-suspect (stream, window) tuples are published to viewers.suspect. The emitter consumes this stream and DOWNWEIGHTS the suspect signal at promote-time (a high-suspect stream's HLL cardinality is multiplied by 1 - suspect_score) — bots are NOT removed from raw HBs (we keep them for forensics + product analytics).

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 (stream, window); at 6M (stream, window)/min that's 30K 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 it to the count topology forces correlated outages; (c) the suspect store is QUERIED by other surfaces (creator-fraud system, ad-billing-clawback, ban-evasion) — it earns its own service boundary. Twitch's 2025-08-21 viewbot crackdown that dropped public concurrents 24% confirms the load-bearing nature; running this inline would have made that deploy a 'reboot the count pipeline' incident.

When it fails. Abuse-svc down → new (stream, window) tuples block at the emitter pending review; existing streams' counts keep flowing. Detection: trending_first_time_promotion_lag_p99 (P2 >10 min). Classifier drift → false-positive rate climbs, real streams' counts artificially deflated; mitigation is weekly canary against held-out validation set + manual-override admin endpoint. Adversarial: sophisticated actors pre-warm accounts to defeat account-age features → mitigation is ensemble (no single feature is load-bearing); the Twitch 2025 deploy used a multi-feature ensemble. Detection: false_positive_rate_24h, classifier_skew, suspect_rate_per_creator. Twitter's published 'state of platform' integrity reports + Twitch's 2025 deploy confirm this is an arms race; the gate is a tax not a solution.

Service · Counter Promotion + WS PublisherGo + Redis pipeline + Kafka cooperative-sticky + WS hub publish protocol

Consumes three Kafka topics in lock-step: viewers.candidates (per-(stream, region, salt) HLLs from Flink), viewers.suspect (per-(stream, window) abuse downweighting), viewers.correction (late-event corrections). For each (stream, window) emit cycle (every 1 s for live count, 30 s for displayed): (1) PFMERGE across regions → one global HLL → cardinality; (2) apply weighting (sample-rate inverse) → weighted-cardinality; (3) apply suspect downweighting (cardinality × (1 - suspect_score)); (4) apply monotone-display smoothing — displayed = max(last_displayed, new) with cap-rate-of-decrease 5%/min; (5) apply milestone-pinning — if last cross was within 120 s, the displayed floor is the milestone; (6) write versioned count payload to count-cache keyed viewers:{stream}:{generation_id} (atomic SET); (7) issue CDN surrogate-key purge for viewers:{stream}; (8) publish a count-update event to ws-hub for fanout to subscribed players; (9) publish a 1 s rollup to viewers.rollup for ClickHouse archive.

Why it exists. Considered making count-read-api do the smoothing/milestone-pinning on every read. Rejected because (a) the smoothing state (last_displayed, last_milestone_crossing_ts) is per-(stream) and would have to live in a separate state store the read-api reads on every request — blowing the 80 ms read SLO; (b) milestone-pinning state is OWNED by the emit-time logic; the read path is a dumb cache lookup. 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. The Twitch QoUX architecture (Lambda transform+aggregate after Kinesis) is the analogue — the emitter is the 'transform + smooth + publish' service after the sketch stage.

When it fails. FALLING-BEHIND-GLOBALLY is the load-bearing failure: emitter is the SOLE writer to the count-cache + WS push path, so global lag = global staleness. Detection: emitter_global_p99_emit_age (P1 >90 s — beats per-stream alerts that lag a global outage by minutes). Cache write failure (count-cache down on a shard) → emitter fail-CLOSED on primary-down (better to serve stale than wrong); dual-write to secondary cache during cutover. CDN purge half-failure → emitter logs-and-moves-on after Cloudflare returns 503 → permanently stale until next emit cycle naturally overwrites; defense is generation-id check at CDN edge (rejects payloads older than X seconds AS WELL AS HMAC). WS hub publish fail → count-cache still has the value, but subscribed players don't get the push → fallback: player polls HTTP every 60 s if WS-stale-alert fires; player-side staleness threshold 90 s. Suspect-consume lag while candidate-consume healthy → join uses stale suspect set, real viewbot campaign goes unflagged; detection: per-topic consumer-group lag (split signals, NOT aggregated). Bad deploy that corrupts smoothing state for all streams → canary-by-stream-shard catches before 8% of streams drift.

Cache · Counter Read Path (Redis)Redis 7 cluster mode, 8 shards (versioned counter blob storage)

Authoritative-for-the-read-path counter. Each stream has a single key viewers:{stream_id} holding the most-recent count payload (count, generation_id, last_emit_ts, suspect_score, ended:bool, signed_hmac). The count-read-api hits GET — sub-millisecond response. Also stores viewers:milestone:{stream_id} (most-recent milestone + crossing-ts) for the emitter's milestone-pinning logic. Versioned key pattern: emitter atomically SETs viewers:{stream}:{generation_id} first, then SETs the unversioned viewers:{stream} pointer — so a slow emitter's stale write LOSES the SET race (the unversioned key always points at the highest generation_id seen). Eviction: TTL-keyed (5 min for inactive streams; live streams refresh TTL on every emit).

Why it exists. Considered serving directly from ClickHouse via materialized views. Rejected because ClickHouse single-row p99 on a counter read is 30–100 ms (OLAP doesn't excel at point-query OLTP), and we need <5 ms cache reads to fit the 80 ms read budget. Considered an in-process cache in count-read-api. Rejected because cache hit ratio depends on cross-replica sharing — with 48 read-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. JioHotstar's >95% cache-hit SLO on the live-counter (published in Kovvuru's Medium write-up) is the production-proven analogue.

When it fails. Cluster cold restart → 30 s of count-read-api degraded reads (falls through to ClickHouse for the most-recent rollup); AOF keeps RPO under 5 s of emitter writes lost. Single shard down → streams hashing to that shard go cold (1/8 ≈ 12.5% of streams) → mitigation is per-stream fallback to ClickHouse at count-read-api. Replication lag spike on async: a follower-read could serve a 5 s-stale count → mitigation is RYW from primary (we use RYW consistency, not eventual). Async replication means writes are SEQUENTIAL per-shard (NOT linearizable cluster-wide) with up-to-5 s RPO on failover — the correctness primitive we actually need is the monotonic-generation-id Lua SET script that REJECTS any write with generation < current. Hot shard from one Mr Beast stream — 5M cold-fetch reads/sec hitting one shard: detection is per-shard QPS asymmetry. Mitigation: edge CDN absorbs >99% of cold-fetch QPS; if origin egress climbs we cell-isolate tier-1 streams onto dedicated shards.

Service · Counter Read APIGo + Redis client + ClickHouse fallback driver + HMAC response signer

Hosts GET /v1/streams/:id/viewers. For each request: (1) resolve home-region for the stream (in-process LRU; etcd-backed); (2) GET viewers:{stream} from count-cache Redis; (3) on cache miss / Redis 5xx, fall through to ClickHouse for the most-recent 1 s rollup (degraded; signed as stale:true); (4) HMAC-sign the response body so the CDN can validate at edge (defense against cache poisoning); (5) return JSON: { count, generation_id, ts, ended, signed_hmac }. Cold-load path: player issues this exactly once on stream-join before subscribing to the WS push channel — so this is the 'cold reader' origin, NOT the steady-state path (steady-state is WS push).

Why it exists. Considered serving directly from emitter's in-process state. Rejected: emitter is per-(stream, region)-shard, doesn't have global counts in-memory, would need to fan-out reads. Considered making CDN's origin a Lambda that reads Redis. Rejected because cold-start tail of 200–800 ms blows the 80 ms read budget. A stateless service tier is the textbook answer; the only question is whether we even need it given CDN+Redis. Answer: we DO because (a) ClickHouse fallback logic is non-trivial; (b) HMAC response-signing needs a key the CDN can't have; (c) suspect-downweighting display logic (showing 'verifying viewers' UI flag during high-suspect windows) needs a service tier.

When it fails. Redis primary down on a shard → read falls through to ClickHouse for 1/8 of streams; ClickHouse single-row OLTP p99 is 30–100 ms vs Redis 1–3 ms — SLO degrades but doesn't break. Detection: redis_5xx, clickhouse_fallback_qps, p99 climb. ClickHouse also down → returns a stale-but-signed 'last-known-good' from in-process 5-min LRU; UI degrades gracefully (yesterday's count > blank). HMAC signing key rotation: dual-key window 24 h; if rotation bombs, CDN starts refusing payloads — canary on every rollout. Bad deploy: 5xx-rate-SLO burn-budget alert + auto-rollback via Argo Rollouts.

Analytics DB · Counter History (ClickHouse)ClickHouse 23.8 (8 shards × 2 replicas, tiered storage, Kafka engine ingest)

Three roles: (1) READ-FALLBACK — count-read-api falls through here on Redis miss; serves the most-recent 1 s rollup (signed stale:true). (2) CREATOR-ANALYTICS — backs the Creator Studio 'concurrent viewers over time' chart with sub-second response on per-stream queries. (3) BILLING + LEGAL — 5-year retention of 1-hour rollups for creator-payout audit (ad-CPM is per-concurrent-viewer-minute) and DMCA / law-enforcement requests. Ingest via the ClickHouse Kafka engine consuming viewers.rollup (the emitter publishes 1 s rollups). Tiered storage: 7-day hot (NVMe), 90-day warm (SSD), 5-y cold (S3 backed by ClickHouse-cloud or ZNS). Materialized views collapse 1 s → 1 m → 1 h on ingest.

Why it exists. Considered serving creator-analytics from Druid (Twitch's published preference). Rejected because (a) the read-fallback workload + analytics workload + long-haul billing workload all want the same data; running three stores is operationally heavy; (b) ClickHouse's sub-second response on the per-stream time-series query is a stronger fit than Druid's broader query surface. Considered InfluxDB / TimescaleDB. Rejected: ClickHouse scales horizontally; InfluxDB at the storage tier we need (PB+ cold) is awkward. Considered S3 + Parquet + Athena for cold tier. Rejected because the latency of cold-tier billing audit (creator disputes a payout) is unacceptable on Athena — ClickHouse's tiered storage gives us seconds, not minutes. Disney+ Hotstar's published simplification (move OFF complex stack ONTO ClickHouse-centric) is the precedent.

When it fails. Kafka-engine consumer-lag is the canonical failure: a deploy of the materialized view DDL stalls the consumer, lag climbs, the fallback-read serves data many seconds old. Detection: kafka_engine_lag_seconds (alert >30 s P2). Mitigation: dual-MV (old + new) during DDL changes, swap-on-validate. ZooKeeper quorum loss → ReplicatedMergeTree writes stall until quorum recovers; mitigation is ZK on dedicated nodes + KRaft migration roadmap. Shard loss → 1/8 of streams' history goes cold; per-shard recovery from object-storage snapshot in <60 min. Cardinality explosion (a bad stream_id tag) → ClickHouse handles cardinality far better than TSDBs but the MV refresh slows; mitigation is per-stream cardinality SLO + drop-on-violation. Cold-tier read latency for old billing audits: 10–30 s acceptable for that workload.

Object · Flink Checkpoints + Cold HB ArchiveS3 (cross-region replicated to us-west-2; 11×9s durability)

Two concerns under one bucket: (1) Flink's RocksDB checkpoint store — each TaskManager async-uploads incremental SST files every 10 s + a manifest. (2) Long-haul cold archive of heartbeat.raw Kafka log dumped via the heartbeat.cold consumer to Parquet for >24h-old replay (Kappa horizon extends beyond Kafka retention via S3 — billing-correction replays can reach back 5 years). Bucket structure: s3://viewer-count-state/checkpoints/{job-id}/chk-{n}/ + s3://viewer-count-state/heartbeat-cold/year=Y/month=M/day=D/hour=H/. Cross-region replicated so a region loss doesn't lose Flink savepoints OR billing-audit raw data.

Why it exists. Considered local-disk-only checkpoints. Rejected: a TaskManager host loss + RocksDB state loss means re-bootstrapping from Kafka offset — fine for replay but 24 h of state to rebuild = ~4 h catch-up at our throughput. Considered EFS/EBS-backed durable. Rejected: cross-region durability is what we need for region-loss DR; EBS doesn't replicate cross-region; S3 has 11-nines durability cross-region. Considered HDFS. Rejected: another distributed system to operate. Pinterest Flink retros, Lyft Flink-at-scale talks, AWS Flink docs all converge on S3 — with prefix-jitter for the request-rate ceiling.

When it fails. S3 SlowDown / 503 throttling during a checkpoint burst is the canonical failure: detection is checkpoint_duration_seconds >60 s + s3_5xx climbing; mitigation is incremental checkpoints + prefix-jitter + exp backoff in 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 (a CIDR change locks Flink out): IAM-policy CI tests against held-out 'Flink-reader' role. AWS S3 us-east-1 has had multiple regional events (2017, 2021); cross-region replication is the only defense.

Coordinator · Config + Hot-Stream Setetcd 3.5 (5-node Raft, 3 AZ)

Tiny but load-bearing coordinator. Holds: (1) hot-streams/{stream} — the set of streams flagged for salt-fanout (TTL'd, rewritten by the hot-key detector job every 60 s); (2) stream-tier/{stream} — the sampling-tier (1/2/3) per stream (set on stream-start, mutable via admin); (3) home-region/{stream} — the canonical region for ordering / merge-primary election; (4) panic-mode/{cohort} — the server-side panic-mode flags read by the gateway and SDK (set during overload to trigger client-side degradation); (5) tier1-cells/{cell} — capacity-pinning records for tier-1 (cell-isolated) streams; (6) revoked-tokens — short-lived denylist of session-stream-tokens. Watched by ingest-svc, flink, emitter, gateway, ws-hub.

Why it exists. Considered Redis for these primitives (already deployed). Rejected because (a) Raft-replicated linearizable writes are the correctness primitive we need for the hot-stream-set — Redis async replication risks two ingest pods seeing different salt decisions for the same key (split state on the same stream → unrecoverable HLL miscount); (b) etcd's watch semantics are stronger than Redis pub-sub (guaranteed delivery, ordering, snapshot-on-reconnect) — critical for the panic-mode flag which MUST propagate within seconds. Considered ZooKeeper. Rejected: etcd is more operationally tame at this footprint AND already deployed for Kubernetes' control plane (we co-locate). Considered Consul. Rejected: etcd's gRPC API is a better fit for our Go services.

When it fails. etcd Raft loss-of-quorum (split-brain on 2 of 5 nodes) → reads and writes STALL; services run on cached state (~60 s tolerable). Detection: etcd_quorum_state, watch_disconnect_rate, raft_propose_failures. Mitigation: 5-node quorum tolerates 2 losses; AZ-spread; never co-locate >2 etcd nodes per AZ. Blast-radius callout: a single etcd outage cascades into (a) auth (revoked-tokens stale → potentially-leaked tokens valid for cache-TTL), (b) panic-mode (load-shed cannot fire), (c) hot-key salt (new viral streams melt one partition), (d) sampler-tier (defaults to tier-3 → 100% retention until refresh) — six concerns on one Raft is acknowledged debt; tradeoff with the operational simplicity vs splitting into etcd-config + revoked-tokens-redis is in the tradeoffs section. Compaction-storm: forgotten retention → etcd disk fills → write-block; mitigation: compaction every 60 s + alert on db-size >2 GB. Adversarial: a runaway hot-key detector dumping 1M hot-stream entries per minute → etcd write-rate exceeded → mitigation is rate-limit at the detector side AND etcd-side --max-txn-ops. Watch-storm during a deploy (10K pods reconnecting) → throttled with --max-concurrent-streams; pods exp-backoff.

Tracing · OTel Collector + JaegerOpenTelemetry Collector (gRPC + tail-sampling) + Grafana Tempo backend

Receives OTLP spans from ingest-svc, flink, emitter, ws-hub, count-read-api via in-process OTel SDK + sidecar collector. Tail-based sampling at the collector (NOT head-based) — every trace's spans buffered for 30 s then a decision is made (errors + slow traces + 1% baseline always kept). Persisted to Grafana Tempo for 7 days. Correlation IDs propagated via W3C TraceContext through Kafka headers (so a heartbeat that triggers a count emit can be traced end-to-end across the async boundary). NOT on the synchronous response path — ingest-svc / emitter / count-read-api emit spans fire-and-forget; the collector buffers locally for 30 s before dropping.

Why it exists. Considered no tracing (just logs + metrics). Rejected because the cross-async-boundary semantics (a slow count-emit caused by a slow Flink slot caused by a stuck S3 checkpoint) are impossible to debug from logs alone — Dapper-style trace IDs are load-bearing for the on-call's first-30-min triage. Considered head-based sampling at SDK. Rejected because errors are rare and head-sampling at 1% misses 99% of them; tail-based at the collector keeps every error + every slow trace + 1% baseline (the cost is the 30-s buffer at the collector). Considered Honeycomb (managed) vs Tempo (self-hosted). Tempo because at our span volume (~30M spans/min peak) managed pricing is prohibitive.

When it fails. Tracing backend unreachable → spans buffer locally for 30 s at the SDK, then dropped (back-pressure on application traffic deliberately disabled — tracing MUST NOT take down the count). Detection: collector_5xx, sdk_export_drop_rate, tempo_storage_lag. Mitigation: collector pre-aggregates + back-pressures span emission; SDK never blocks application threads on export. Self-monitoring: a separate non-OTel metrics system (Prometheus) alerts on tracing-down so we don't lose visibility AND debug capacity simultaneously. Head-based sampling fallback (when collector is overloaded): drop to 0.1% baseline; keep errors + p99-band trips.

Metrics · VictoriaMetrics + Argo Rollouts SoTVictoriaMetrics (3-node TSDB cluster) + Grafana + Alertmanager + Argo Rollouts AnalysisTemplate

Scrapes every service for RED/USE metrics on a 15 s interval (5 s on the load-bearing SLO surfaces: emitter_global_p99_emit_age, count_monotone_violation_rate, kafka_produce_ack_p99, ws_concurrent_connections_per_pod). Evaluates SLO burn-rate alert rules every 30 s and pages on 2% / 1 h fast AND 10% / 6 h slow shapes. Argo Rollouts polls VictoriaMetrics's PromQL endpoint as the AnalysisTemplate source-of-truth — every canary stage queries (success-rate, p99-latency, burn-rate) and auto-rollbacks if any breaches its threshold. Holds the recording rules that pre-aggregate per-stream-tier rollups so on-call dashboards (Heartbeat health, Flink topology, Emitter, WS hub, Count integrity) don't scan raw cardinality at query time. Cardinality control: per-stream_id labels are FORBIDDEN at the ingest relabel step (we bucket into stream_tier=tier1|hot|warm|cold instead); per-IP labels likewise — without this rule a 200M-viewer / 2M-stream system would explode the head-block memory.

Why it exists. Considered Datadog. Rejected on cost at our cardinality (per-pod × per-region × per-tier series fan-out at the 200M-viewer scale is well into seven-figure / month territory). Considered Prometheus single-node. Rejected — at 24 h retention × 5 K active series / pod × ~6 K pods we burn the head-block memory ceiling on one node; VictoriaMetrics's vmselect / vmstorage split scales horizontally to our cardinality. Considered tracing-as-metrics (RED from Tempo span counts). Rejected because Argo Rollouts queries PromQL natively, NOT TraceQL; the canary auto-rollback control loop runs against metrics, end of story. Tracing complements but cannot replace.

When it fails. Cardinality explosion (a new label leaking — e.g. someone instruments a per-stream_id metric on the hot path) is the load-bearing failure: head-block memory pressure climbs, ingest stalls, queries time out, the Argo Rollouts canary auto-rollback control loop goes SILENT — bad deploys ride to 100% because the rollback query can't evaluate. Detection: vm_rows_inserted_total rate-drop, vm_cardinality_limits_exceeded, alert_evaluation_duration_seconds_p99. Mitigation: per-tenant series-limits hard-rejected at relabel ingest, recording-rule downsampling for high-cardinality series, recording-rule eviction in panic mode. Alertmanager down → dead-man's-switch alert via a SECONDARY paging path (PagerDuty integration on a separate route AND a cross-region paired Alertmanager). Datadog March 2023 multi-region outage is the cautionary tale for the metrics tier of an OLAP-class system — keep 30% headroom in steady state and DO drill failover monthly.

Service · Hot-Stream DetectorFlink (lightweight side-job) + Count-Min Sketch over heartbeat.raw + etcd writer

Separate Flink side-job (own consumer-group on heartbeat.raw) that maintains a per-(stream, 60 s sliding) Count-Min Sketch keyed by stream_id and emits the top-N streams (typically top-1000) into etcd's hot-streams/ keyspace with a 5-min TTL. ingest-svc and the main Flink topology watch this keyspace; when a stream appears, ingest-svc starts salting its partition key and Flink keeps that salt in its KeyBy — the 16 sub-shards are merged downstream at the emitter (PFMERGE), not un-salted inside Flink. This is the load-bearing structural defense against the 'viral go-live → one Kafka partition melts' failure mode (failure #2 in the research catalog) — without it the salt-fanout path is dead code.

Why it exists. Considered baking hot-stream detection into the main Flink topology. Rejected because (a) hot-stream detection's update cadence (every 30 s) is different from the count topology's emit cadence (every 1 s) — coupling them forces compromises; (b) hot-stream detection state (the CMS over stream_ids) is independent of the per-stream HLL state — sharing a Flink job means correlated savepoint failures take both down; (c) the abuse-svc precedent (separate Flink job for the same heartbeat.raw stream) is the playbook. Considered Kafka topic with consumer-group lag scraped from Prometheus + alerting rule. Rejected: too slow (60 s alert lag), and the salt set needs to be authoritative in etcd, not derived from a metric scrape.

When it fails. Detector down → ingest-svc + Flink fall back to whatever salt-set is cached (last 60 s of state); a NEW viral stream after the outage is unsalted until the detector recovers → that partition melts within 30 s. Detection: hot_key_detector_last_emit_age (P1 >120 s). Mitigation: 1 active + 1 standby; auto-failover via Flink JobManager HA; chaos-test detector outage monthly. Stale set (the cached set says stream X is hot but it's no longer) → tier-1 stream consumes 16× the state for no benefit; mitigation is the 5-min lease TTL. Detector backlog under viral spike (rare but possible) → at >10K viral streams the CMS overflows or top-N computation slows → mitigation is bound top-N at 1000 (sufficient for the load shape — at 200M viewers / 2M streams the top-1000 captures >99% of skew).

Service · Schema RegistryConfluent Schema Registry 7.5 (HA pair, Avro / Protobuf)

Holds the authoritative Avro/Protobuf schemas for every Kafka topic (heartbeat.raw, viewers.candidates, viewers.suspect, viewers.correction, viewers.rollup, viewers.lifecycle). Producers fetch the latest schema-id at startup (cached in-process for 60 s); consumers fetch by schema-id encoded in each Kafka message. Enforces backward-compatibility on schema evolution (a producer-rolled forward schema must be readable by every consumer on the older schema) so an ingest-svc deploy cannot poison the flink topology. DLQ-rule: a message whose schema-id is unknown OR fails deserialize is rejected at the consumer and routed to kafka.heartbeat.dlq for ops review — the structural defense against the 'one poison message restarts the Flink job every 10 min' failure mode (Twitch 2019 NA outage shape).

Why it exists. Considered inline JSON parsing with schemaless contracts. Rejected because (a) at 6.67M HB/s a deserialize-error in 0.001% of messages = 67 errors/sec = a Flink-restart loop without DLQ routing; (b) backward-compat is a deploy-safety primitive — without it, ingest-svc and flink deploys are coordinated, which kills the per-service deploy cadence; (c) DLQ routing of poison messages requires a stable schema-id-or-reject contract. Considered Apache Avro without a registry (schema-in-message). Rejected: 200 B Avro-schema-per-message overhead at 6.67M HB/s = 1.3 GB/s of pure schema metadata, dwarfs the actual heartbeat payload. Confluent Schema Registry pattern is the production-default; Confluent docs document the DLQ + schema-id-reject contract.

When it fails. Schema-registry down → producers fail-CLOSED on schema-id mint (new schemas can't be registered, but cached schema-ids stay valid for 60 s); consumers fail-CLOSED on unknown schema-id (rejects to DLQ instead of crashing). Cached schemas survive 60 s then the path stalls — by which time the on-call is paged. Detection: schema_registry_5xx_rate, schema_lookup_p99, dlq_message_rate (a spike here means a schema mismatch — investigate the most recent deploy). Mitigation: HA pair across AZ; Kafka _schemas topic RF=3 means the schemas survive a broker loss. Backward-incompatible schema accidentally registered (a producer rolls forward an incompatible Avro change) → registry REJECTS the registration at write time per the BACKWARD policy; a deploy-time CI test against the registry should catch this BEFORE the deploy hits production. Disney+ Hotstar and Twitch both publish the schema-registry-as-deploy-gate pattern.

Service · Stream Lifecycle ControllerGo + Kafka producer + admin-RPC + RTMP-ingress integration

The producer of kafka.viewers.lifecycle. Two signal sources: (1) the platform's RTMP / WHIP encoder-ingress fires stream-start when the first encoder packet arrives AND stream-end when the encoder disconnects or admin terminates; (2) the Creator Studio admin console fires manual start/end/pause events. Each event published with (stream_id, event_type, ts, encoder_session_id, home_region, expected_viewers_tier) into viewers.lifecycle (32 partitions, 7-day retention). The emitter, ingest-svc, and ws-hub all consume this topic — stream-start warms the count to 0, registers home-region in etcd, and pre-allocates HLL keys in sketch-bag; stream-end triggers the count-drain protocol (cap-rate-of-decrease 5%/min for 60 s then hard-zero, see deep-dive 7).

Why it exists. Considered making ingest-svc emit lifecycle events when it sees the first HB for a stream. Rejected because (a) the encoder side is the authoritative source — a stream might have ZERO viewers (just the broadcaster encoding) and still need lifecycle tracked; (b) waiting for first-HB delays the count-zero-init (the emitter would have no viewers:{stream} key when the first cold reader arrives → fallback to ClickHouse for a stream with literally no history → blank UI); (c) the stream-end signal MUST come from the encoder side, not from 'no HBs in 60 s' — the latter is indistinguishable from a transient network blip. Considered embedding into the API gateway admin path. Rejected: lifecycle has its own retention, schema, and consumers — earns its own service boundary. JioHotstar's published architecture has a dedicated stream-state-controller for exactly this reason.

When it fails. Service down → no new stream-start signals → new streams' counts stay zeroed (degraded UX: count shows '0 watching' even with viewers actively heartbeating); the count-drain on stream-end can't fire so a CDN-cached '1.4M watching' stays stale for the 5-s TTL window (then naturally clears). Detection: lifecycle_publish_5xx_rate, encoder_disconnect_to_lifecycle_publish_p99, stream-start-detected-but-not-emitted-count. Mitigation: 6 replicas across 3 regions, encoder-side retry budget (3 attempts × 200 ms = 600 ms), fallback to ingest-svc's first-HB heuristic after 5 min of missed lifecycle events. RTMP-ingress integration brittleness (a new encoder protocol version mis-fires events) → schema-registry + canary on the encoder integration path.

Service · WS Fanout Hub (tier-2)Go + gRPC-streaming + per-stream topic registry + consistent-hash routing

The MIDDLE tier of the Twitch-PubSub-style hierarchical fanout tree. The emitter publishes count-update events into ~100 tier-2 hubs (each owning a stream-id range by consistent-hash); each tier-2 hub fans out to its subtree of ~30 ws-edge pods (which terminate the actual client TCP sockets). This collapses the emitter→ws-edge fan-out from 12-emitters × 3000-edges = 36000 connections to 12 × 100 + 100 × 30 = 4200 connections — an order of magnitude reduction in connection-state at the publisher side. The tier-2 hub also performs LAST-MILE message dedup (a slow emitter retry doesn't re-publish to the subtree) and per-subtree backpressure (a wedged edge pod doesn't slow the whole subtree).

Why it exists. Considered flat fanout (emitter directly publishes to all 3000 ws-edge pods). Rejected because (a) at 12 emitters × 3000 edges = 36000 long-lived gRPC streams per emit cycle — each emitter pod holds 3000 streams in-process, an operational nightmare on TCP-keepalive, certificate rotation, and graceful shutdown; (b) a deploy of 100 edge pods causes 100 × 12 = 1200 gRPC-stream reconnect events at the emitter (each emitter is the hot connection-target); a tier-2 layer absorbs that into 100 × 1 = 100 reconnects at the tier-2 layer per emitter. Twitch's PubSub three-tier hierarchy (Edge → PubSub-core → backend publishers) is the production-proven shape; Discord's Manifold's 'send per remote-NODE not per-PID' is the BEAM equivalent. Without this tier, the slow-client cascade authored on ws-hub's failureMode is mis-attributed — the tier-2 layer is where node-to-node fanout-tree partitions are detected and contained, not at the edge.

When it fails. Tier-2 hub down → its subtree of ~30 ws-edge pods (carrying up to 2.25M client connections) loses pushes for ~5 s during reroute; consistent-hash flips half of the affected streams to a neighboring tier-2. Detection: tier2_hub_heartbeat_age (P1 >10 s), per-tier2 publish_p99. Mitigation: stateless tier-2 with consistent-hash ring stored in etcd; auto-reroute on heartbeat miss; clients reconnect-and-resnapshot from ws-edge cache. Fanout-tree partition (network split between tier-1 emitter and tier-2 hub on one region) → that region's tier-2 layer goes dark; the regional ws-edge subtree falls back to HTTP-poll fallback (60 s cadence — the panic-mode degraded path). The tier-2 layer is the CONNECTION-INTERMEDIATE layer described in the Discord Manifold post; without it, the 1B-concurrent tier-flip in 'tradeoffs' is infeasible.

Service · Panic-Mode ControllerGo + PromQL client (metrics-svc) + etcd writer + Raft-elected single-writer

The AUTONOMOUS writer of etcd.panic-mode/{cohort}. Every 10 s polls the metrics tier for SLO-burn signals: emitter_global_p99_emit_age > 60 s, kafka_engine_lag_seconds > 30 s, gateway_5xx_rate > 1 %, ws_pod_conn_buffer_high_water > 80 %, hb_ingest_qps > 1.5× provisioned. On sustained-30-s violation across the rule set, writes panic-mode/{cohort} = on to etcd with a 5-min lease (so it auto-clears if forgotten). Cohort scoping: per-region (panic-mode/region/us-east), per-tier (panic-mode/tier/tier1), or global (panic-mode/global). The flag propagates via etcd watches in <500 ms to the gateway (which embeds it in the next /heartbeat response header) and to the SDK on the client (which drops non-essential surfaces per the deep-dive 8 priority order: reactions → recommendations → comments → HB cadence 30 s→60 s → WS push paused → recommended sidebar disabled).

Why it exists. Without an autonomous trip, panic-mode is human-in-the-loop at 3 a.m. — the difference between graceful degradation in 30 s and a 10-min Twitter-trending outage while the on-call joins the bridge. The whole point of the panic-mode protocol is that it's autonomous; an SRE manually flipping a feature flag is NOT panic-mode, it's just feature-flag-flipping. JioHotstar / Pragmatic Engineer published that automated server-signaled degradation is the only viable path at the 25M+ concurrent tier — they hit 821M during the ICC T20 World Cup 2026 final on the same playbook. Considered embedding the trip logic into the metrics-svc itself (the AlertManager already evaluates these rules for paging). Rejected because metrics-svc's job is to ALERT humans; this is a side-effecting CONTROL plane action and deserves a separate audit log + canary path. The lease+hysteresis logic (30 s sustained to trip, 5 min clean to clear) is also non-trivial.

When it fails. Controller down (all 3 replicas) → panic-mode flags freeze at last-known state; existing tripped flags persist for their 5-min lease (acceptable), but new overloads can't auto-trip → falls back to manual SRE flip via the admin CLI. Detection: controller_leader_alive, controller_decisions_per_minute, panic_mode_lease_age_seconds. False-positive trip (panic-mode tripped while system actually healthy) is the operator's worst day — mitigation: 30 s sustained-violation rule + 5 min hysteresis + canary on every threshold-update path (the rules-config in etcd is itself canary-deployed). Lease-loss flapping → the 3 replicas' lease-election bounces every few seconds → mitigation: etcd lease TTL >> max-pause (we set 30 s lease, 10 s renew). Metrics-svc down → controller cannot evaluate → fails-CLOSED (does NOT trip a new panic-mode in absence of signal), but existing flags remain — better than spurious trips.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Typical clarifications to surface:

  • What does "viewer" mean? Authenticated session? Unique device-id? Heartbeat in the last 60 s? Chunk-request in the last segment? Twitch counts lurkers (muted/backgrounded players) but not chatbots; YouTube discards "low-quality playbacks." The definition gates everything downstream.
  • Accuracy budget? ±2% on hot streams, ±5% on long-tail. Twitch publicly admits the long-tail figure; HLL at p=14 gives ±0.81% intrinsic accuracy.
  • Freshness? 30–90 s industry norm. YouTube and Twitch explicitly slow / freeze the displayed count for verification — the freeze is intentional, not a bug.
  • Hottest single entity? 8.06M (YouTube Chandrayaan-3, world record), 59M (JioCinema CWC 2023 Ind-Pak final), 65M+ (JioHotstar T20 finals). Design for 60M on the hottest entity, 200M platform-wide.
  • What's broken if it's wrong? Public count = creator monetization signal + ad billing + social-proof flywheel. A monotone-violating display ("1.4M → 1.2M → 1.5M") becomes a brand crisis within minutes.
  • What's broken if it's late? A 60 s lag is fine; a 5 min lag is "the count is broken" on social.
  • Multi-region? Yes — the count must merge across regions; the cross-region presence-divergence failure (eu sees 5M, us sees 3M) is a published Cloudflare / BGP-leak-class concern.
  • Bots? Yes — viewbotting is an entire economy. Twitch's 2025-08-21 deploy dropped public concurrents ~24% in days; the defense is structural (session-bound stream-tokens), not aspirational.

Assumptions stated:

  • 200M concurrent viewers (peak), 60M on hottest entity, 2M concurrent active streams.
  • 30 s heartbeat cadence (matches HLS segment); 5 s during ad-break or quality-shift.
  • 5 active regions.
  • 5-year retention for the per-hour rollup (creator billing + DMCA + legal).

02Functional reqs

What must this system actually do?

  • Heartbeat ingest — POST /v1/heartbeat from the player every 30 s with the session-bound stream-token.
  • Count read (cold) — GET /v1/streams/:id/viewers, served by CDN ≥ 99% of the time.
  • Count subscribe (warm) — wss:// subscribe to viewers:{stream}; server pushes new count every 30 s.
  • Count promote — internal: emit per-stream count to count-cache + WS push + CDN purge.
  • Stream lifecycle — admin endpoints (publish, end) propagate to the count pipeline; end triggers count-drain.
  • Creator analytics — Creator Studio queries the per-second timeseries from ClickHouse.
  • Billing rollup — per-hour rollups retained 5 y for creator-payout audit.
  • Viewbot scrub — abuse signals downweight suspect counts before publish.
  • Panic mode — server-side flag tells SDKs to degrade non-essential surfaces.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Availability — 99.99% on heartbeat ingest (gates billing); 99.9% on the displayed number (one bad minute tolerable). Slack S&L target.
  • Latency — heartbeat ingest p99 ≤ 5 s (fire-and-forget); count-read p99 ≤ 80 ms origin (CDN cache absorbs >99%); WS-push delivery p99 ≤ 200 ms after emit; promote-to-public freshness 30 s p50, 90 s p99.
  • Accuracy — ±2% on tier-1 (top-100) streams, ±5% on long tail. HLL p=14 gives intrinsic ±0.81%; the rest is sampling + viewbot-scrub uncertainty.
  • Durability — Heartbeats are LOSSY (RPO 60 s); per-second rollups MUST persist (RPO 0 in ClickHouse) for billing.
  • Consistency — Display: monotone, milestone-pinned. Internal: per-(stream, region) HLL is RYW from primary; cross-region merge is eventual within 16 s.
  • Scalability — Horizontal on every tier. Hot-stream salt-fanout is the structural answer to the hot-key problem at the 60M-on-one-entity scale.
  • Security — Session-bound stream-tokens at the heartbeat protocol; HMAC-signed count payloads; per-IP/ASN rate-limits; OWASP at the edge; JA3 fingerprint matching for viewbot detection.

04Capacity estimation

How much load and data does this have to hold?

Assumptions: 200M concurrent viewers peak, 60M on hottest entity, 30 s heartbeat cadence, 5 s ad-break cadence, HLL p=14.

  • Heartbeat QPS (steady) = 200M / 30 s = 6.67M HB/s.
  • Heartbeat QPS (ad-spike) = 200M / 5 s = 40M HB/s (provision for this — ad-breaks are scheduled).
  • Hot-entity HB/s = 60M / 30 s = 2M HB/s on one stream — this is the load-bearing hot-key problem; mitigation is salt-fanout (16 sub-shards) + tier-based sampling (1% at tier-1 → 20K HB/s actually published per shard).
  • Heartbeat ingress bytes = 6.67M × 380 B = 2.5 GB/s raw; with tier sampling → 570 MB/s to Kafka.
  • HLL bytes/sketch = 2^14 × 6 bits = 12 KB per sketch; total state = 12 KB × 2M streams × 5 regions × 4 tiers ≈ 480 GB — fits across 128 Flink slots (3.7 GB/slot) AND across 12 Redis sketch-bag shards (40 GB/shard).
  • Cross-region merge = 5 regions × PFMERGE 50µs ≈ 250 µs per stream emit; PFMERGE 2M streams in 10 s window = 500 µs × stream-shard = parallelizable.
  • WebSocket push QPS = 200M / 30 s = 6.67M push/s; with 75K conns/pod = ~2700 WS edge pods needed (we provision 3000).
  • Cold storage — Per-second per-stream sample ≈ 64 B compressed; 2M streams × 86400 × 7 d × 64 B ≈ 77 TB hot tier. Per-hour rollups × 5 y × 2M streams ≈ 7 TB cold tier — affordable.
  • Kafka on-disk = 570 MB/s × 86400 × 24 h × RF=3 ≈ 148 TB across 36 brokers ≈ 4.1 TB/broker.

The hot-key tier flip:

  • Pure Redis INCR on one key works to ~10K viewers / stream.
  • Local pre-aggregation (per-edge counters flushing every 1 s) to ~100K.
  • HLL with per-region sketch-merge gets to ~10M.
  • Tier-1 sampling at 1% + salt-fanout gets to ~60M (our target).
  • Beyond 60M, the displayed number is a statistical projection — and that's correct, because at ±1.6M nobody notices.

05API design

What does the outside world call, and what comes back?

/play-start mints the session-stream-token that bootstraps everything below — the heartbeat Authorization: Bearer and the WS Sec-WebSocket-Protocol token are both this token. Token-mint goes through auth-svc (Ed25519, HSM-fronted key custody); the count read does NOT.

POST /v1/play-start
Content-Type: application/json

{
  "stream_id": "abc123",
  "viewer_session": "vs-9f3...",
  "ja3_hash": "a0e9f5...",
  "attestation": "<opaque>"
}

200 OK
{
  "session_stream_token": "sst.<ed25519-signed>",
  "expires_in": 900,
  "next_heartbeat_ms": 30000
}

Token is bound to the tuple (viewer_session, stream_id, ip_prefix, ja3_hash, attestation) and signed Ed25519 with a 15-min (900 s) TTL — short enough to bound a leaked token, long enough that a typical watch never refreshes mid-session. On auth-svc failure, /play-start fails over to a degraded token (anonymous, narrow-scope): the cold viewer still joins and watches, but contributes no count — degraded count, not outage. A HB or WS upgrade presenting an expired token gets 401; the SDK refreshes once via /play-start, not retry-forever.

POST /v1/heartbeat
Authorization: Bearer <session-stream-token>
Content-Type: application/json

{
  "stream_id": "abc123",
  "viewer_session": "vs-9f3...",
  "ts": 1709251200000,
  "cdn_pop": "AMS-3",
  "bitrate_rung": 1080,
  "buffering_ms": 0,
  "abr_switches": 1,
  "ja3_hash": "a0e9f5...",
  "attestation": "<opaque>"
}

202 Accepted
{ "next_heartbeat_ms": 30000, "panic_mode": false }
GET /v1/streams/:id/viewers
# Served by CDN > 99% of the time.

200 OK
Cache-Control: public, max-age=5, stale-while-revalidate=10
{
  "stream_id": "abc123",
  "count": 1402356,
  "generation_id": 1709251230,
  "ts": 1709251230000,
  "ended": false,
  "signed_hmac": "AbCd1234..."
}
wss://hub.example.com/v1/stream/{stream_id}
Sec-WebSocket-Protocol: token.<session-stream-token>

→ SUBSCRIBE { topics: ["viewers:abc123"] }
← SNAPSHOT { topic: "viewers:abc123", count: 1402356, generation_id: 1709251230 }
← UPDATE   { topic: "viewers:abc123", count: 1404102, generation_id: 1709251260 }
← ENDED    { topic: "viewers:abc123", count: 0, ended: true }

06Data model

What gets stored, and what is it looked up by?

heartbeat.raw Kafka topic (partitioned 1024 ways by hash(stream_id‖salt)):

fieldtypenotes
stream_idstringsalted upstream for hot keys
viewer_sessionuuidanonymized at SDK
tsint64event-time ms
regionstringderived from POP
cdn_popstringPA, AMS, …
bitrate_rungint
buffering_msint
ja3_hashstring
attestationbytesopaque, validated upstream
weightfloat1 / sample-rate (tier-based)

viewers.candidates Kafka topic (1 s window, per (stream, region, salt)):

fieldtype
stream_idstring
saltint (0–15 salt-shard for hot streams, 0 for cold)
regionstring
window_end_tsint64
hll_sketchbytes (12 KB)
weighted_cardinalityint64
sample_ratefloat
generation_idint64 (monotonic, event-time-based)

viewers.rollup Kafka topic (1 s rollup for ClickHouse):

fieldtype
stream_idstring
tsint64
countint64
suspect_scorefloat
generation_idint64

ClickHouse viewers_1s table (sharded by cityHash64(stream_id), ordered by (stream_id, ts), TTL 7 d → 1 m MV → 90 d → 1 h MV → 5 y):

CREATE TABLE viewers_1s (
  stream_id String,
  ts DateTime64(3),
  count UInt64,
  suspect_score Float32,
  generation_id Int64,
  region LowCardinality(String)
) ENGINE = ReplicatedMergeTree
ORDER BY (stream_id, ts)
TTL ts + INTERVAL 7 DAY TO VOLUME 'warm',
    ts + INTERVAL 90 DAY TO VOLUME 'cold';

07High-level design

Which components handle a request, and in what order?

Architecture summary:

  1. CDN (Cloudflare + multi-CDN steering) caches the count JSON at the edge with a 5 s TTL + surrogate-key purge. Absorbs >99% of cold-fetch reads.
  2. API Gateway (Envoy + WAF + per-IP rate-limit) terminates TLS, validates session-stream-tokens via auth-svc, routes heartbeats to ingest-svc / reads to count-read-api / WS upgrades to ws-hub.
  3. Auth · Stream-Session-Token Service mints session-bound stream-tokens at /play-start and validates per-HB; the load-bearing viewbot defense.
  4. Heartbeat Ingest Service validates the envelope, applies tier-based sampling (100% / 10% / 1%), applies hot-key salt (etcd-driven), enriches with geo/ASN, publishes to Kafka with idempotent producer.
  5. Kafka is the spine — heartbeat.raw (1024 partitions) plus emit topics for candidates / suspect / correction / rollup / lifecycle.
  6. Flink Viewer-Count Topology consumes heartbeat.raw, per-(stream, region) HLL state in RocksDB, 1 s window emit to viewers.candidates + 1 s rollup to viewers.rollup.
  7. Sketch-Bag (Redis HLL) stores the per-(stream, region) HLL registers for the emitter's PFMERGE step.
  8. Abuse / View-Fraud Detector (separate Flink job) consumes the same heartbeat.raw and emits viewers.suspect with downweighting scores.
  9. Counter Promotion / Emitter consumes candidates + suspect + correction, PFMERGEs across regions, applies monotone-display + milestone-pinning + suspect-downweighting, writes to count-cache, issues CDN purge, publishes to ws-hub.
  10. Count-Cache (Redis) is the read-path source-of-truth for count-read-api on CDN miss.
  11. Count-Read API serves GET /viewers (origin behind CDN) with fallback to ClickHouse on Redis miss.
  12. WebSocket Hub holds 75K conns / pod, fans count-updates to subscribed players via Twitch-PubSub-style hierarchical tree.
  13. Count-History (ClickHouse) stores 1 s rollups (7 d hot), 1 m (90 d warm), 1 h (5 y cold) for analytics + billing + DMCA.
  14. Checkpoints (S3) holds Flink RocksDB checkpoints (cross-region replicated) + cold raw-HB Parquet archive for Kappa replay beyond Kafka retention.
  15. etcd holds the hot-stream-set, stream-tier table, panic-mode flags, revoked-tokens; watched by ingest-svc / flink / emitter / gateway.
  16. Tracing (OTel + Tempo) with tail-based sampling at the collector.

Data flow on heartbeat: Player → CDN-bypass → Gateway → Auth (validate token) → Ingest-svc (sample + salt + enrich) → Kafka (publish, 202 returns) → Flink (consume + HLL update) → Sketch-bag (PFADD) → Flink (1 s window emit) → Kafka.candidates → Emitter (PFMERGE + smooth + suspect-downweight) → Count-cache (versioned SET) → CDN (purge) + WS-hub (push to subscribers) → Player (display).

Data flow on count-read (cold): Player → CDN (hit > 99% of the time) → return cached JSON. On miss: → Gateway → Count-Read API → Count-Cache (hit > 99% in steady state) → CDN-populate → Player. On Redis miss: → ClickHouse fallback (signed stale:true).

08Deep dives

Which part breaks first, and what do you do about it?

See the deepDivePrompts in the frontmatter (1: hot-key salting; 2: monotone-display + milestone-pinning; 3: viewbot defense in depth; 4: WS fanout tree; 5: ad-break spike; 6: multi-region merge; 7: stream-end drain; 8: panic-mode; 9: Kappa replay).

09Trade-offs

What did this design cost, and what breaks at 10×?

Production choices and rejected alternatives:

  • Sampled heartbeats vs every-event ingest — chose tier-based sampling (100/10/1%) because at tier-1 the 0.81% HLL accuracy floor dominates and we'd burn compute counting bot-noise. Twitch and YouTube freeze the published number; we sample.
  • HLL (Redis HLL primitive) vs Theta-sketch (DataSketches) vs literal SET — chose HLL because it's mergeable, bounded (12 KB at p=14), and is a first-class Redis primitive. SET is 5+ GB per stream at 100M cardinality. Theta-sketch is overkill for cardinality-only (it gives union/intersection/difference — we only need union for region-merge).
  • WS push vs HTTP poll — chose WS because at 200M concurrent the poll-storm at 30 s cadence = 6.7M poll/s. WS holds the socket; only the publisher sends.
  • Single Kafka cluster vs per-tier clusters — chose single (with cell-isolated tier-1 dedicated cluster). LinkedIn flips multi-cluster at ~6 GB/s broker-write; we're at 1.7 GB/s aggregate.
  • Emitter writes both count-cache AND ws-hub vs WS hub reads from cache — chose dual-publish because cache-read on every WS-push is N×M fan-out; the emitter's dual-write is N+M.
  • etcd vs Redis for hot-stream-set — chose etcd because linearizable replicated writes prevent two ingest pods seeing different salt decisions; Redis async replication risks split-state.
  • 30 s freshness vs 5 s — chose 30 s because YouTube/Twitch publish on 30 s and product wants what users expect. 5 s freshness costs 6× the WS push QPS and provides no measurable UX gain.

What breaks at 5× scale (1B concurrents, 300M on hottest entity):

  • Sketch-bag Redis cluster ceiling — at 480 GB total HLL state × 5 = 2.4 TB; need to re-shard to 60 shards (already cluster-mode-friendly) OR move sketches into Flink-only QueryableState.
  • WebSocket fanout — 1B / 75K conns = 13,300 pods; manageable but expensive. Move to MQTT-based persistent broker (Meta Live's pattern) or push via the HLS manifest itself (avoid the WS layer entirely).
  • Kafka spine — 6 GB/s broker-write × 5 = 30 GB/s; multi-cluster with MirrorMaker 2 federation (LinkedIn's pattern).
  • CDN egress — at 1B concurrent reads/s × 200 B = 200 GB/s edge egress; multi-CDN steering becomes mandatory.

Primary sources

  • Flajolet, Fusy, Gandouet, Meunier — HyperLogLog (DMTCS 2007)
  • Heule, Nunkesser, Hall — HyperLogLog in Practice (EDBT 2013)
  • Cormode, Muthukrishnan — Count-Min Sketch (J.Alg. 2005)
  • Apache DataSketches — Theta sketch (Yahoo)
  • Engineering at Meta — Under the hood: Broadcasting live video to millions (2015)
  • Engineering at Meta — Scaling Live streaming for millions of viewers (2020)
  • Twitch Engineering — State of Engineering 2023 (Spade, PubSub, Kinesis)
  • Twitch Engineering — Breaking the Monolith at Twitch (2022)
  • Twitch Engineering — The QoUX Journey (2025)
  • Twitch Engineering — How Twitch Uses PostgreSQL (2016)
  • Twitch Developers — PubSub API (≤10 conns/IP, ≤50 topics/conn)
  • Discord — How Discord Scaled Elixir to 5,000,000 Concurrent Users (Manifold + Semaphore + FastGlobal)
  • Discord — Real-time Communication at Scale with Elixir (2020)
  • Slack Engineering — Real-time Messaging (Presence Servers + Gateway Servers)
  • Slack Engineering — Migrating Millions of Concurrent WebSockets to Envoy
  • HasGeek Rootconf — Scaling hotstar.com for 25M concurrent viewers (2019)
  • Pragmatic Engineer — Live streaming at world-record scale with Ashutosh Agrawal (JioHotstar)
  • Last9 — Cricket Scale Series #1 (IPL 30M concurrent)
  • ByteByteGo — How Disney+ Hotstar / JioHotstar scales (NAT-per-subnet, multi-CDN)
  • Cloudflare — June 21 2022 cross-region routing retro
  • AWS — Kinesis Data Streams Nov 25 2020 retro (thread-limit cascading failure)
  • Apache Flink — Stateful Stream Processing + Watermarks + RocksDB checkpoints
  • Confluent — KIP-429 Cooperative-Sticky Rebalance
  • Confluent — KIP-794 Strictly Uniform Sticky Partitioner
  • Kreps — Questioning the Lambda Architecture (Kappa, 2014)
  • Beyer et al. — SRE Workbook (Managing Load, Addressing Cascading Failures)
  • YouTube Help — How engagement metrics are counted (the 'we freeze on purpose' rule)

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 Live Viewer Count (YouTube/Twitch) yourself