Command Palette

Search for a command to run...

Hectal
PHASE 16Advanced ~8 min· topic 3 of 7

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.

0/7 · 0%

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

  1. 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.

  2. 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).

  3. 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.

  4. 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.

  5. 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

analytics.mermaiddiagram
Rendering diagram…

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

01

Why land raw events in a data lake as well as aggregating them?

Practice

01

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.