System designHard4 min readruns in the simulator

Netflix System Design (High-Scale)

Open Connect CDN for playback, a queue-driven transcode pipeline: flush the cache and watch the origin take the hit.

  • cdn
  • streaming
  • storage
  • microservices
  • queue
Live model
throughput
6,000 rps
play p99 latency
134 ms
cost
$10,417/mo
S3 Backfill MissUpload masterSmart TV App — client · tvSmart TV Appclient · tvMobile App — client · mobileMobile Appclient · mobileZuul API Gateway — gateway · ×30%Zuul API Gatewaygateway · ×3Open Connect CDN — cdn0%Open Connect CDNcdnAuth Service — api · ×30%Auth Serviceapi · ×3Billing Service — api · ×20%Billing Serviceapi · ×2Geo-Licensing Service — api · ×20%Geo-Licensing Serviceapi · ×2Recommendation Engine — server · ×40%Recommendation Engineserver · ×4Transcode Scheduler — orchestrator0%Transcode SchedulerorchestratorTranscode Task Queue — queue0 msgsTranscode Task QueuequeueGPU Transcoder Node — worker · ×100%GPU Transcoder Nodeworker · ×10Video Block Store (S3) — store0%Video Block Store (S3)storeBilling DB (MySQL) — db0%Billing DB (MySQL)dbUser Profiles DB (Cassandra) — db · cassandra0%User Profiles DB (Cassandra)db · cassandraSession & Geo Cache — cache · redishit 0%Session & Geo Cachecache · redisTelemetry Bus (Kafka) — bus · kafka0 msgsTelemetry Bus (Kafka)bus · kafkaMaxMind Geo-IP Service — cloud0%MaxMind Geo-IP ServicecloudStudio Ingest Client — client · studioStudio Ingest Clientclient · studio
The model you’ll run. Open it to see requests move, then break a part.

Requirements

Functional

  • Users can browse a personalized home feed of recommended titles.

Non-functional

  • Playback starts within 500 ms p99 for content already on the CDN edge.play p99 latency < 500 ms
  • Support 200M subscribers, up to 50M concurrent streams at peak (modeled here at a scaled-down 6,000 req/s so the simulation runs at a workable size).play successful requests ≥ 4000 rps

The requirements with a measurable target are checked live in the simulator: break something and watch them fail.

High-level architecture

Smart TV Appclient · tv

Mobile Appclient · mobile

Zuul API GatewayapiGateway

Open Connect CDNcdn

Auth Serviceserver · api

Small-scale compute cluster

Billing Serviceserver · api

Small-scale compute cluster

Geo-Licensing Serviceserver · api

Small-scale compute cluster

Recommendation Engineserver

Transcode Schedulerorchestrator

Transcode Task Queuequeue

GPU Transcoder Nodeworker

Video Block Store (S3)objectStore

Billing DB (MySQL)db

User Profiles DB (Cassandra)db · cassandra

Distributed NoSQL database with 3 replicas

Session & Geo Cachecache · redis

Single-shard in-memory cache

Telemetry Bus (Kafka)messageBus · kafka

Distributed message broker

MaxMind Geo-IP Servicecloud

Studio Ingest Clientclient · studio

Request paths

Browse the home feed

Zuul API GatewayRecommendation EngineUser Profiles DB (Cassandra)

Start playback

Open Connect CDNon a missVideo Block Store (S3)

Start playback via the CDN

  1. Press play.

    70% of viewer requests are playback. The TV and mobile apps fetch video segments straight from the Open Connect edge; the API gateway isn't on this path.

  2. Served from the edge.

    92% of segment requests are edge hits, served in about 10 ms.

  3. A miss goes to S3.

    The rest are read through from the S3 block store, about 60 ms more, and kept at the edge for the next viewer. A cold edge takes around 45 s to warm up.

  4. The target.

    Playback must start within 500 ms p99 for content already at the edge, and hits keep it far inside that.

Studio upload & encode

Video Block Store (S3)Transcode SchedulerTranscode Task Queue

Video transcoding pipeline

  1. Upload the master.

    A studio uploads the master file to the S3 block store.

  2. Schedule the encode.

    The upload triggers the transcode scheduler, which splits the work into encoding tasks.

  3. Queue the tasks.

    Tasks wait in the queue, up to 5,000 deep, so a burst of uploads never swamps the encoders.

  4. Encode on GPUs.

    10 GPU workers pull tasks, about 5 a second each, and write the encoded blocks back to S3, ready for the CDN.

What happens when it breaks

The edge goes dark

Open Connect fails entirely: every playback start errors out, but browsing the catalog keeps working — the edge is a single point of failure for playback only.

ReadingHealthyDuring the failure
p99 latency252 ms251 ms
error rate0%70%
throughput6,000 rps1,800 rps
cost$10,417/mo$10,417/mo
  • Playback starts within 500 ms p99 for content already on the CDN edge. (fails)
  • Support 200M subscribers, up to 50M concurrent streams at peak (modeled here at a scaled-down 6,000 req/s so the simulation runs at a workable size). (fails)
Run this failure yourself, then pick the fix.Break it yourself

Trade-offs

Fixes and what they cost

FixWhat it doesTrade-off
Coalesce CDN missesConcurrent misses for the same title share one origin fetch, which prevents a thundering herd on the object store.Viewers asking for the same title wait for the first origin fetch.
Add origin read capacityMore origin read throughput, at a real cost per month.Costs more per month for origin capacity that sits idle most of the time.
Pre-warm the CDN edgeA flushed or restarted edge loads 90% of the popular titles before it serves, so the origin never sees a cold-cache stampede. It takes effect on the next flush or restart.A restarted edge takes longer to come back while it loads popular titles.
Circuit breaker on gateway to recommendationsWhen recommendations fail, the gateway stops waiting on them and serves a generic row. Latency drops, and the struggling service gets room to recover.Viewers see a generic row instead of personal recommendations while it's open.

Sources

  1. Netflix Open Connect

Run this design. A live model of this system. Watch requests flow, then break it.

Open in the simulator