Topic 16.3
Designing Analytics Pipelines: Netflix Events and YouTube Analytics
In one line
High-volume analytics pipelines ingest client events through a gateway into Kafka, aggregate them with stream processors in windows, and land raw and aggregated data in a lake or warehouse. Key decisions: partition keys for aggregation, compression, retention, backpressure handling, and idempotent or exactly-once aggregation.
Think of it like this
Counting footfall in a huge shopping mall. Sensors at every door stream counts to a control room (Kafka); staff aggregate per shop per hour (stream processing); the full raw log is archived for later analysis (data lake).
Key ideas
- 01
Ingest: clients batch events to an edge gateway (HTTP), which validates, enriches (geo, device) and produces to Kafka with large batches and zstd compression. Clients buffer and retry offline.
- 02
Keys: by user or session ID for session analytics (order per user), by video ID for per-video aggregation (beware viral videos: use two-stage aggregation), or null for pure archival streams (best spread).
- 03
Processing: Kafka Streams or Flink compute windowed aggregates (views per video per minute, watch time), with event time and grace periods for mobile clients that upload late. Exactly-once or idempotent upserts into the serving store.
- 04
Storage: raw events to object storage via a sink connector (Parquet, partitioned by date/hour) for replay and batch analysis; aggregates to an OLAP store (ClickHouse, Druid, Pinot) or warehouse for dashboards.
- 05
Retention: Kafka keeps days (enough to replay processing failures); the lake keeps years. Backpressure: Kafka absorbs spikes as lag; scale processors on lag; drop or sample low-value events under extreme load.
Code & diagrams
Interview problem
The problem
Netflix-style event pipeline at millions of events/sec
Design the pipeline for play, pause, stop, seek, rating and recommendation-interaction events at millions of events per second: clients → gateway → Kafka → stream processing → lake and warehouse. Discuss partitioning, compression, retention, replay and backpressure. Then adapt it for YouTube view counts (billions/day, exactly-once or idempotent counts).
You're given
- 3M events/sec peak
- ~1 KB per event
- Dashboards within 1 minute
- Raw data kept for years
Explain it without notes
Why land raw events in a data lake as well as aggregating them?
Practice
Compute broker count for 600 MB/s compressed ingest, RF 3, 3 consumer groups, 3 days retention, 12 TB usable per broker.
Trade-offs
- ↔
Longer Kafka retention simplifies replays but multiplies storage; lakes are cheaper for history.
Done when you can
I can design a high-volume analytics pipeline with sizing, keys, windows, storage and replay.