YouTube / Netflix Streaming
Worked solution

YouTube / Netflix Streaming — a worked solution

ABR, CDN, transcoding, hot/cold — the egress, fan-out, and popularity Pareto.

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 YouTube / Netflix Streaming workspace

The problem

Build a video streaming platform like YouTube or Netflix. Users press play and expect a 4 K stream to start within two seconds; creators upload 50 GB mezzanines and expect the full ABR ladder published within minutes. Almost every byte the system serves is read; almost every byte it stores was produced by transcoding. The architectural pressure is not QPS — it is egress (you are 15% of the US internet at peak), fan-out (one source mezzanine becomes 250 K encode jobs across rungs × codecs × shots), and the popularity Pareto (the top 1% of catalogue serves 50% of bytes; the bottom 50% serves <0.1%).

The reference design below is what an SRE staff team at Netflix-class scale would actually run a five-year incident retro against — not the 45-minute whiteboard sketch. Every node has load-bearing internals; every edge declares wire semantics; every chaos card is grounded in a cited real-world incident. Focus is on the four pillars: ABR, CDN, transcoding, and hot/cold tiering. Comments / search / recommendations / ads are explicitly out of scope.

The reference architecture

Reference architecture for YouTube / Netflix Streaming: 17 components — Client, CDN · Open Connect Edge, CDN · Regional Shield, API Gateway, Service · Manifest Builder, Auth / DRM License Server, Cache · Catalogue (EVCache), SQL DB · Metadata (Vitess), Object Store · Hot (S3 Standard), Object Store · Cold Archive, Service · Upload Orchestrator, Stream · Transcode Bus, Worker · Encoder Farm, Worker · Packager + DRM, Worker · Popularity / Tier Mover, Queue · Encoder DLQ, Monitoring — connected by 28 flows.ClientclientCDN · Open Connect Ed…Netflix Open Connect (F…CDN · Regional ShieldOpen Connect storage ti…API GatewayEnvoy + Cloudflare WAF …Service · Manifest Bu…Go / Java; HLS m3u8 + D…Auth / DRM License Se…Widevine (Google) + Fai…Cache · Catalogue (EV…Netflix EVCache (Memcac…SQL DB · Metadata (Vi…Vitess (sharded MySQL) …Object Store · Hot (S…AWS S3 Standard / GCS /…Object Store · Cold A…S3 Glacier Deep Archive…Service · Upload Orch…Cosmos ingest control-p…Stream · Transcode BusKafka 3.7 (3 topics: en…Worker · Encoder FarmFFmpeg + x264/x265/libv…Worker · Packager + D…Bento4 / Shaka Packager…Worker · Popularity /…Apache Flink stream pro…Queue · Encoder DLQKafka tombstone topic +…MonitoringNetflix Atlas (in-memor…
17 components, 28 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Client

Three concrete callers share this node: (a) the playback player (Shaka / hls.js / ExoPlayer / AVPlayer / dash.js) that fetches manifests + segments + DRM licences; (b) the upload UI used by creators to push mezzanine sources; (c) telemetry on the same player emitting playback events back to the popularity pipeline. The player runs an ABR algorithm (BBA, BOLA, or Pensieve-class) that picks rungs from the manifest based on buffer level and throughput estimate.

Why it exists. Considered modelling the upload UI as a separate edge node because its failure modes (chunked-resumable upload, large body, retryable PUT) differ from playback. Rejected because both run on the same SDK and ship from the same release pipeline — splitting them only matters when we model authn (OAuth + device binding for playback, signed-upload-key for creators), and that lives at the gateway downstream. Player IS the surface where the ABR regression that puts the fleet on 480p actually originates.

When it fails. Bad client deploy regresses ABR — within minutes, fleet drops to lowest rung; egress goes DOWN, error rate stays zero, p99 latency looks fine. Detection: average_delivered_bitrate, rebuffer_ratio, rung_switch_rate. Mitigation: 1% / 10% / 50% staged rollout gated on bitrate-delta SLO. SDK quickstarts ship the canonical ABR + retry helpers so naive clients don't regenerate manifest URLs on retry.

CDN · Open Connect EdgeNetflix Open Connect (FreeBSD + NGINX) / CloudFront / Akamai

Serves manifests + segments to nearby viewers. One appliance is 96 GB RAM + 16 TB NVMe + dual 100 Gb NIC, sustains ~90 Gb/s TLS at ~55% CPU under Netflix's kernel-TLS + sendfile + RACK/BBR build. Holds the ~5-10% of catalogue popular in this country (per-region popularity feed). Serves manifests with short TTL (60-300 s VOD, 1-5 s live), serves segments with long TTL (hours-to-days, immutable URL).

Why it exists. Considered using only a public CDN (Cloudfront / Akamai) — that's how Disney+, Hotstar, and most others run. Rejected for Netflix-class scale because at >5% of any ISP's traffic, settlement-free peering economics flip in favour of embedding appliances inside the ISP's rack: no transit costs, predictable latency, ISP gets a free peering agreement. Open Connect carries ~100% of Netflix bytes; the public CDN model is the right answer for catalogues that haven't crossed the embed threshold.

When it fails. Most-watched failure: BGP withdrawal / mis-route at the CDN provider (Fastly 2021-06-08 pattern) makes the whole edge tier unreachable from the public internet — internal health checks pass because internal paths are fine. Detection: external synthetic monitor 5 s polling against a non-CDN'd path, SYN-ACK rate by ASN. Mitigation: multi-CDN steering with real-user-measurement (Conviva / NPAW). Hot-key on a viral release saturates one OCA's outbound link; mitigation is per-segment sharding + edge prefill of top-N titles.

CDN · Regional ShieldOpen Connect storage tier / CloudFront Origin Shield / Akamai tiered cache

Sits between the edge POPs and the origin. When the edge OCA misses, the request falls through here before touching the object store. Sized at ~10x edge capacity per region; holds the long-tail catalogue that doesn't fit in edge NVMe. Akamai calls this 'tiered cache', CloudFront calls it 'Origin Shield', Open Connect calls it the 'storage tier'.

Why it exists. Considered going edge → origin direct (no shield). Rejected because the cache miss explosion at origin is the main cost driver — without the shield, origin egress is ~1/cacheHitRatio of total egress; with a shield at 95% hit, origin sees ~0.05 × 0.05 = 0.25% of traffic instead of 5%. The shield exists to protect the origin's bandwidth bill and the metadata path, not to lower viewer latency.

When it fails. Shield outage pushes 100% of edge misses to origin — that's a ~20x origin spike (5% miss rate × 1/0.05 amplification). Detection: shield_hit_ratio, origin_5xx_rate, origin_egress_gbps. Mitigation: edge falls through directly to origin with a stale-while-revalidate budget on the popular tier; long-tail viewers see latency tail. Cross-region replication of the shield's catalogue prefill so a regional shield outage doesn't take a full region's mid-tier offline.

API GatewayEnvoy + Cloudflare WAF + TLS terminator

L7 ingress for every request that misses the edge: manifest builds, license fetches, uploads, telemetry, control-plane. Terminates TLS, runs the WAF (token-auth signature check, country geo-block, rate-limit per device), copies the entitlement claim from the JWT onto upstream headers signed with HMAC, mints short-lived signed URLs the manifest service injects into the response.

Why it exists. Considered baking auth + signed-URL minting into the manifest service. Rejected because (1) license requests + upload requests need the same edge concerns, so duplicating them is bug-prone, (2) the WAF rule set (geo-blocks, abuse signatures) belongs at one layer that fronts everything, (3) signed-URL minting holds an HMAC secret that should NOT live next to product code. Centralizing auth means upstream services trust the signed envelope without round-tripping.

When it fails. Gateway pool exhaustion during a live-event join surge masks as 'manifest fetch slow' but is actually upstream backpressure from license / manifest. Detection: gateway_active_streams, downstream_5xx_rate. Mitigation: drain mode + per-zone failover; circuit-break a misbehaving device fingerprint before 503'ing the world. The Disney+ 2019-11-12 launch lesson: test at 5x expected peak, not at average — premiere nights hide a 20x minute-burst inside a 5x hour-average.

Service · Manifest BuilderGo / Java; HLS m3u8 + DASH mpd + CMAF

Generates the master + media manifests on demand. Reads the asset's encoded variants (rungs × codecs × DRM systems) from the catalogue cache; if missed, falls through to the metadata DB. Injects signed URLs (minted by the gateway) for every segment, applies the per-region geo-rule (rung trimming for poor networks, codec gating per device class), and writes the rendered manifest back into the edge with a 60-300 s TTL.

Why it exists. Considered serving manifests directly from object storage as static files. Rejected because (1) signed-URL injection requires a per-request HMAC signature that must NOT live in S3, (2) per-device rung trimming (don't expose AV1 to a 2019 iPad) is dynamic, (3) the manifest re-renders every TTL anyway. Considered building manifests at the edge — same problem with HMAC custody. Centralized rendering is the simplest design that keeps signing keys away from caches.

When it fails. Catalogue-cache cold start: hit ratio drops from 95% to 60% during a regional cache restart, manifest p99 climbs from 30 ms to 250 ms because every miss touches the DB. Detection: catalogue_cache_hit_ratio, manifest_build_p99. Mitigation: warm the cache from the popularity feed BEFORE the restart; bounded-staleness manifests (serve last-known-good manifest with a longer TTL during cache warmup). The metadata-DB primary failure mode is what the trace simulator's manifest-build-cold-db trace exists to model.

Auth / DRM License ServerWidevine (Google) + FairPlay (Apple) + PlayReady (Microsoft)

Two concerns under one roof: (1) entitlement — does this account have rights to play this title in this country today? (2) DRM license — return a content-key wrapped to the device so the player can decrypt the stream. Validates the signed envelope from the gateway, checks entitlement against metadata DB (with a hot Redis cache), and returns a Widevine / FairPlay / PlayReady license depending on what the device asked for. Widevine ≈ 60-65% of devices, FairPlay ≈ 25-30%, PlayReady fills the rest.

Why it exists. Considered splitting entitlement and license minting into two services. Rejected because (1) every license fetch needs the same entitlement check (skip it and you've broken DRM contracts), (2) the latency budget is tight (p99 ≤ 200 ms from the player) and an extra hop breaks it, (3) the failure modes are coupled: entitlement DB down → license fetch fails. We do split the license server pool by DRM system internally (different binary deploys for Widevine vs FairPlay so a Widevine bug can't tank FairPlay), but they share the same entitlement code path.

When it fails. Premiere-night overload: every player on Earth requests a license in the same 5-min window; pool saturates, exponential back-off retries make it worse. Disney+ Mandalorian 2019-11-12 was exactly this. Detection: license_5xx_rate, license_fetch_p99, gateway_active_streams. Mitigation: shadow-load test at 3x peak, license-server pre-warm, regional sharding, pre-authenticated tokens. The 'license server overload' chaos card models this directly.

Cache · Catalogue (EVCache)Netflix EVCache (Memcached fork, Ketama consistent hashing)

Stores asset:{titleId}:variants (rungs × codecs × DRM × signed-URL templates), entitlement:{accountId}:{titleId} (rights-cache for license server), popularity:{titleId} (per-region rank). Cache-aside on read: manifest service tries here first (p99 ~5 ms), falls through to DB on miss; DRM service tries here first too. Populated lazily on miss; INVALIDATED on every catalogue write so a buggy writer can't poison the cache.

Why it exists. Considered Redis cluster (we already use it elsewhere). Picked EVCache (Memcached + Ketama) because the read pattern is 99.9% point-lookup with no need for set/list/sorted-set ops, and EVCache's heterogeneous-cluster allocation means losing one shard takes 1/N of the keyspace offline (graceful degradation), not a quorum failure. Considered making the cache the source of truth for popularity scores. Rejected because the popularity worker recomputes them from the playback stream — making the cache authoritative would mean every popularity update has to invalidate it anyway.

When it fails. Cache cold restart sends all manifest builds to the DB → manifest p99 climbs from 30 ms to 250 ms; DB connection pool saturates; manifest builds start 503'ing. Detection: cache_hit_ratio, manifest_p99, db_pool_active. Mitigation: warm from the popularity feed BEFORE restart; single-flight on the manifest service so 1000 concurrent misses for the same key don't hammer the DB 1000 times. Hot shard on a viral release: kill-shard chaos card models this directly.

SQL DB · Metadata (Vitess)Vitess (sharded MySQL) — YouTube pattern; alt: Spanner

Source-of-truth for the catalogue: title rows (id, region rights, takedown state), asset rows (variants × codecs × DRM × URLs), entitlement rows (accountId × titleId × geo × expiry), popularity scores. Schema is sharded by titleId (Vitess vindex on title_id_v) for a 64-shard layout. Reads served by semi-sync followers within the AZ; writes go to the shard primary; cross-shard joins are forbidden — denormalize at write time.

Why it exists. Considered DynamoDB (single-key access, zero ops). Rejected because (1) takedown queries — 'show me everything in country X marked disabled this week' — need indexes that don't fit DynamoDB's GSI model cleanly, (2) the auditor escape-hatch matters: regulators ask for 'show me all titles licensed by Disney with rights in Brazil' and you need ad-hoc SQL. Considered Spanner — it's the right answer for global serializable, but at our QPS the cost is ~3-5x Vitess for the marginal consistency we don't need (catalogue writes are not contended).

When it fails. Primary failover during prime time: 30 s of write outage on one shard; reads from followers continue at bounded staleness. Manifest builds that hit catalogue cache (95% path) survive cleanly; manifest builds that miss BREAK for the failover gap. New title launches that race the gap surface as user-visible publish delays. Detection: shard_leader_health, oldest_unreplicated_lsn, manifest_build_cold_p99. Mitigation: catalogue cache absorbs 95% of reads through the gap; new launches gated on synthetic 'first-watch' probe.

Object Store · Hot (S3 Standard)AWS S3 Standard / GCS / Colossus, multi-region replicated

Stores everything that can be read by the CDN: source mezzanines (ProRes / IMF, ~400 GB/hr), the full encoded ladder (10-50x source, ~7-10 rungs × 3 codecs × 2-3 audio tracks per asset), packaged CMAF + DRM-encrypted segments, manifests-on-disk for offline rebuild. Sized for the active catalogue (single-digit PB at Netflix-class, ~2 EB at YouTube-class). Replicated cross-region asynchronously (RPO minutes-to-hours).

Why it exists. Considered self-hosting on hardware (Open Connect appliances ARE Netflix-built hardware, but for serving — not for the origin tier). Rejected because the origin needs 11×9s durability across multiple AZs and the operational burden of running tape-as-storage at this scale is enormous. S3 + cross-region replication gives us the durability story; we layer caching (OCAs + regional shield) on top to dodge the egress bill.

When it fails. us-east-1 partial outage (2017-02-28, 2021-12-07): uploads fail, encoder writes fail, cold-tail viewers fail, hot bytes keep playing because OCAs are the source of truth for those. Detection: PutObject_5xx_rate, transcode_publish_lag, oldest_unreplicated_master_age. Mitigation: warm-warm cross-region; route encoder writes to secondary region during the outage; do NOT flip read traffic (OCAs already serve). The 's3-origin-region-outage' chaos card models this.

Object Store · Cold ArchiveS3 Glacier Deep Archive / GCS Coldline / tape

Holds the ~95% of mezzanine masters and encoded variants that aren't actively served — long-tail catalogue, retired-but-licensed titles, the 50-year retention horizon on every master. Retrieval is batched (hours) and cheap; storage is ~$1/TB/month for Deep Archive vs ~$23/TB/month for Standard. Tier-mover (the popularity worker) demotes here when popularity drops; promotes back to hot when a back-catalogue title trends.

Why it exists. Considered keeping everything in S3 Standard. Rejected on cost: at YouTube-class catalogue (~2 EB), the standard-tier cost differential is in the hundreds of millions per year, and the hit ratio on the cold archive is by definition <0.1%/asset/month. Considered tape-only (the way old broadcasters do it). Rejected because retrieval latency (days) is too long for popularity surges — we need 'hours, not days' to move a trending title back to hot.

When it fails. Glacier retrieval rate-limit: a popularity surge that needs 1000+ titles restored simultaneously hits the per-account rate limit; restore queue grows to days. Detection: glacier_restore_queue_age, popularity_promote_lag. Mitigation: pre-emptive restore of trend-candidates 24 h ahead via the popularity worker's prediction model; tier-promote scheduling that batches restore requests.

Service · Upload OrchestratorCosmos ingest control-plane; S3 multipart-upload pre-signed URLs

Three control-plane endpoints — bytes do NOT flow through this service. (1) POST /v1/uploads/init: validates creator ownership against metadata-db, validates declared content-type / size, creates an S3 multipart-upload session, returns the upload-id + a set of pre-signed PUT URLs (one per part). (2) Client uploads each part DIRECTLY to S3 against those URLs. (3) POST /v1/uploads/:id/complete: receives the per-part ETags from the client, calls S3 CompleteMultipartUpload (a metadata-only API), validates the resulting object's checksum + size, then writes the asset row to metadata-db and publishes the encode-job to Kafka in the same outbox txn. The orchestrator stays out of the byte path entirely.

Why it exists. Considered routing bytes through the service (the URL-shortener writesvc shape). Rejected because at YouTube-class ingest of ~3.3 GB/s sustained, pass-through doubles the bandwidth bill (client→service + service→S3 = ~6.6 GB/s of provisioned egress) AND the service tier becomes the bandwidth bottleneck — a 30-replica pool simply cannot saturate 6.6 GB/s without becoming an absurd network appliance. Considered fronting S3 with a transparent proxy that streams chunks. Rejected because the proxy still sits in the byte path with the same economics — direct-to-S3 with signed URLs is what every reference design (Cosmos / MezzFS / Disney+ ingest / Hotstar) actually does. The orchestrator's job is auth + lifecycle + transactional publish, not byte forwarding.

When it fails. Orphan multipart uploads: client crashes between init and complete; bytes accumulate in S3 with no asset row referencing them. Detection: incomplete_multipart_uploads_age, S3 storage-class growth rate. Mitigation: S3 lifecycle policy auto-aborts after 7 days; nightly sweeper reconciles open multipart-uploads against issued init records. Signed-URL leak: a creator's URL is exfiltrated and used by an attacker to upload garbage to the slot. Detection: per-uploadId byte-hash + per-init session-token. Mitigation: 15 min URL TTL + per-IP rate limit on init. Orchestrator down: new uploads can't START, but in-flight uploads continue (clients have signed URLs and S3 is up) — graceful degradation. Detection: upload_init_5xx_rate, signed_url_mint_failures.

Stream · Transcode BusKafka 3.7 (3 topics: encode-jobs, chunk-ready, playback-events)

Three logical topics on one cluster: (1) encode-jobs — upload publishes, encoder consumes; ~250 K work items per 30-min episode (rungs × codecs × ~900 shots) at full per-shot fan-out. (2) chunk-ready — encoder publishes when an encode chunk is durable, packager consumes for stitching + DRM-encrypt. (3) playback-events — CDN edge emits view events, popularity worker consumes for tier scoring. Single cluster keeps cross-topic ordering possible when needed.

Why it exists. Considered separate clusters per topic (better blast-radius isolation). Rejected because the operational overhead of running 3 Kafka clusters per region (3 × N regions = 9-15 clusters globally) didn't pay for the marginal isolation — a single cluster with per-topic ACLs, per-topic compaction, per-topic retention is the cheaper-and-equally-isolated choice. Considered AWS Kinesis. Rejected because per-shot encoding fan-out (250 K work items / episode) needs >10 K partitions for parallel consumer scaling — Kinesis has hard shard limits that bind us into rebalance hell.

When it fails. Tentpole release lands all 8 episodes simultaneously; encoder queue depth blows to 2 M jobs; consumer lag goes from seconds to hours; bottom rungs ship on time but AV1-4K rungs ship 6 h late. Detection: kafka_consumer_lag_records, oldest_unconsumed_msg_age_seconds. Mitigation: priority partitions (tentpoles get reserved partitions that don't share with long-tail re-encodes); spot autoscale on the encoder consumer group; DLQ for poison chunks routed to encoder-dlq so one un-encodable chunk doesn't pin a worker into retry-loop. SECOND failure mode: outbox CDC relay lag — the asset-publish saga writes a row to metadata-db AND publishes a chunk-ready / cache-invalidation event via Kafka; if the CDC relay stalls (poison row, partition leader churn, schema-registry wedge), packager publishes never reach the catalogue cache invalidation consumer, and freshly-packaged titles stay invisible at the edge until the relay drains. Detection: oldest_unpublished_outbox_age, debezium_slot_lag_bytes. Mitigation: bound the WAL retention so the slot can't grow unbounded; alert at age > 30 s; switch to a fresh slot + replay from WAL archive if catastrophic.

Worker · Encoder FarmFFmpeg + x264/x265/libvpx-VP9/libaom-AV1/SVT-AV1; YouTube Argos VCU ASIC

Pulls encode-jobs from the transcode bus, downloads the source chunk from object-hot via MezzFS-style streaming mount (no full-download), runs the codec at the requested rung × codec, writes the encoded chunk back to object-hot, then publishes a chunk-ready event so the packager can stitch. Software encoders run on spot instances; the AV1 hot path runs on YouTube-class Argos VCU ASICs (10x more compute-efficient than CPU encoders).

Why it exists. Considered making the upload service kick off encodes synchronously. Rejected because (1) per-shot fan-out for a 30-min episode is ~250 K work items — the upload SLA can't wait, (2) encoding is CPU-bound and bursty (tentpoles) so the right tool is async with elastic capacity, (3) encoder failure modes (poison chunks, codec crashes) need a DLQ pattern that's natural in async but ugly in sync. Considered ASIC-only. Rejected because not all codec combinations are supported on Argos hardware yet (libaom-AV1 quality tuning still needs CPU); we run the hybrid.

When it fails. Tentpole queue backlog (Squid Game S2 lands all 8 episodes at midnight) — bottom rungs ship on time, top rungs ship 6 h late. Detection: encoder_queue_oldest_msg_age, sla_breach_count. Mitigation: priority queues (tentpoles get reserved encoder pool), DLQ for poison chunks. Codec regression: a bad libaom-AV1 build corrupts a fresh batch of titles before the SSIM/PSNR audit catches it. Detection: encode_psnr_p99, encode_ssim_p99. Mitigation: ringed encoder rollouts (1% / 10% / 50%) gated on quality SLOs.

Worker · Packager + DRMBento4 / Shaka Packager + CENC + Widevine/FairPlay/PlayReady

Consumes chunk-ready events from the transcode bus, fetches the encoded chunks for an asset from object-hot, stitches them into CMAF segments + master/media manifests, encrypts with CENC (CTR for Widevine/PlayReady, CBCS for FairPlay) using a content key, writes the packaged bundle back to object-hot, and writes the asset's published=true row to metadata-db (with the variants fingerprint + DRM key URLs). Once metadata-db acks, the title is live: catalogue cache invalidation propagates, edge fills can pick it up.

Why it exists. Considered doing packaging inside the encoder. Rejected because (1) packagers stitch ACROSS shot boundaries (encoder is per-shot), so a packager has to wait for all shots' chunks before finalizing a manifest, (2) DRM custody is sensitive — the content key wrapping and license-server handshake belong in a separate, hardened service, (3) packager bugs (the chunked-encode fence-post bug Netflix wrote about) produce A/V drift, and isolating that failure mode from encoder churn matters operationally. Considered shipping the packager output to the gateway directly. Rejected because publish must be transactional with metadata-db (asset visible iff packaged).

When it fails. Chunked-encode fence-post bug: chunk boundaries misaligned, stitching produces A/V drift > 40 ms — viewers report 'dialog out of sync'. Detection: pre-publish A/V sync probe (SSIM + audio waveform alignment), post-publish error reports. Mitigation: overlap-and-validate chunk boundaries at the encoder; A/V sync probe gates publish. DRM key rotation desync between encoder + license server: encoder ships ciphertext with key K1, license server still serving K0 → decryption fails on every player. Mitigation: atomic key bundles, ringed deploys, license caching on-device with a 24 h TTL so a brief desync doesn't affect existing sessions.

Worker · Popularity / Tier MoverApache Flink stream processing + tier-move scheduler

Consumes playback-events from the transcode bus (the player's session-start + bitrate-switch telemetry, ingested via the gateway). Computes per-title popularity scores per region with a sliding decay window (1 hr / 24 hr / 7 day / 30 day buckets). Writes scores to metadata-db. Drives two control loops: (1) every 6 h, demote bottom-decile titles from object-hot → cold-archive (storage lifecycle action). (2) when a back-catalogue title's 24 h score crosses a promote threshold, prefetch from cold-archive → object-hot so the next OCA fill window picks it up.

Why it exists. Considered keeping popularity scoring inside the catalogue cache directly. Rejected because the input stream (playback events at ~20 M events/sec) doesn't fit a cache's write pattern, and the windowed-aggregation logic needs a stream processor. Considered Spark batch (run nightly). Rejected because tier-promote latency matters for trending content — by the time a TikTok meme-trend hits 'titanic-rerelease' levels, you want the title hot within hours, not the next morning. Flink streaming gives us 5-minute granularity at acceptable cost.

When it fails. Glacier restore rate-limit: a TikTok-trending sports replay needs 500 historical titles restored simultaneously; restore queue grows to days. Detection: glacier_restore_queue_age, popularity_promote_lag. Mitigation: pre-emptive restore 24 h ahead via predictive trend signals; budget per-region per-day so a single trend can't pin the restore queue. Promote storm: every region promotes the same 50 titles within an hour after a global event — encoder farm gets re-encode burst. Mitigation: stagger promotes by region; reuse already-encoded variants whenever possible (no re-encode needed if we already have the right rungs).

Queue · Encoder DLQKafka tombstone topic + S3-archived sample chunks

Receives encode jobs that have failed > N retries (default 5) — corrupt mezzanine, codec crash, missing audio track, malformed timecode. Each entry persists the failing chunk's S3 key, codec parameters, error stack trace, and the originating titleId so a human can triage. Drains slowly via a manual ops review queue; auto-resurrects entries when the encoder ships a fix that resolves a class of failures (e.g. libaom-AV1 1.4 fixes a crash bug → all DLQ entries with that error replay).

Why it exists. Considered letting the consumer simply skip + log a poison chunk. Rejected because (1) silently dropping a chunk produces an A/V-misaligned final stream that ships to viewers, (2) Kafka's at-least-once semantics + idempotent encoder means retry is the expected path; with no DLQ, a single corrupt chunk pins one consumer in a retry-loop forever and hides capacity from healthy work. Considered using a separate Kafka cluster for DLQ. Rejected — same cluster, separate topic, separate consumer group is enough.

When it fails. DLQ fills past 100 K entries: indicates a systemic encoder regression (bad codec build, mezzanine validator bug). Detection: dlq_depth_total, dlq_growth_rate. Mitigation: pause the encoder fleet, roll back to last-known-good build, replay DLQ once the regression is confirmed fixed. Without the DLQ, the same poison chunks would silently retry forever and consume encoder capacity that healthy titles need.

MonitoringNetflix Atlas (in-memory dimensional TSDB) + Mantis stream processing + Conviva / NPAW for client-side QoE

Two streams of signals land here: (a) service-side metrics scraped from every named SLO in the canonical (manifest_p99, license_5xx_rate, gateway_active_streams, kafka_consumer_lag_records, glacier_restore_queue_age, encode_psnr_p99, oldest_unreplicated_lsn, dlq_depth_total) — Atlas serves the dashboards + the alert rules; (b) client-side QoE via Conviva / NPAW: average_delivered_bitrate, rebuffer_ratio, rung_switch_rate, concurrent_session_derivative. SPS (Started Plays per Second) is the load-bearing top-line.

Why it exists. Considered relying on cloud-provider monitoring (CloudWatch / Stackdriver). Rejected because (1) the dimensional cardinality of streaming KPIs (per-region × per-codec × per-rung × per-DRM × per-device-class) blows up cloud-monitoring quota costs at YouTube scale; (2) the JOIN-derivative + per-key QPS alerts need a streaming compute layer (Mantis), which CloudWatch doesn't provide; (3) the client-side QoE signal can't be derived from server-side metrics alone — Conviva / NPAW listen on the player's CDN-fetch results, not the origin's. Without this node, every failureMode claim in the canonical has nowhere to live.

When it fails. Monitor down → on-call goes blind for the duration. Detection: synthetic health check from a pager-bot every 30 s. Mitigation: warm DR Atlas in a second region, alert routing failover via PagerDuty's redundant ingestion endpoints. Gray failure: a metric pipeline regression silently drops dimensions — the dashboard renders but the alert never fires. Mitigation: synthetic 'inject + assert-fired' daily test on 5 critical alert rules.

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:

  • Catalogue scale? Netflix-class (~17 K titles, multi-PB encoded) or YouTube-class (billions of videos, ~2 EB stored)?
  • Live or VOD? Live changes everything — manifest TTL drops from 60-300 s to 1-5 s, segment durations drop from 6 s to 200-400 ms (LL-HLS partial segments), and encoder fall-behind becomes a user-visible latency.
  • Geo coverage? Single-region (Hulu-class) or global (Netflix-class)? At >5% of any ISP's traffic, embedded ISP CDN (Open Connect) flips on.
  • Codec mix? H.264 only (legacy, every device), HEVC + VP9 (most), AV1 (top-watched titles only)?
  • DRM? Free + ad-supported (unencrypted), or premium (Widevine/FairPlay/PlayReady)?
  • Latency targets? P99 startup ≤ 2 s VOD, ≤ 500 ms live; rebuffer ratio < 0.5%; manifest p99 < 100 ms; license p99 < 200 ms.
  • Upload fidelity? ProRes-class mezzanine (400 GB/hr) or compressed source (50 GB/hr)?
  • Per-title vs fixed ladder? Per-title encoding (Netflix 2015) saves ~25% bitrate at constant VMAF; per-shot (Netflix 2018) adds another ~10%.

Assumptions to state (Netflix-class):

  • 90 M concurrent at prime-time peak; 65-120 M for tentpole live events.
  • 5 Mbps average delivered bitrate (ABR-determined; ranges 235 kbps mobile-cellular to 16 Mbps 4 K HDR).
  • 6 s segments VOD / 2 s LL-HLS partial segments live.
  • 7-10 ABR rungs × 3 codecs × 2-3 audio tracks per asset (~250 K work items per 30-min episode at full per-shot fan-out).
  • 17 K+ titles in active catalogue; ~2 EB total stored at YouTube-class.
  • Top 1% of catalogue serves 50% of bytes; top 10% serves 90%.
  • 17 K+ Open Connect Appliances embedded in ISPs globally; 600+ regional shields.
  • 50-year master retention; encoded variants on a 24-month decay (re-derivable).

02Functional reqs

What must this system actually do?

  • Playback — player fetches manifest, selects rungs via ABR, fetches segments, fetches DRM license, decrypts + decodes + renders.
  • Upload — creator pushes mezzanine (chunked, resumable); system acks once master is durable; transcode runs async.
  • Transcode — fan-out per-shot encoding × N rungs × M codecs; CMAF + CENC packaging; DRM-encrypted bundles published to object store.
  • Catalogue — title metadata + asset variants + entitlement; reads dominate (manifest builds, license entitlement checks).
  • Hot/cold tiering — popularity worker scores titles per-region; promotes from cold archive when trending, demotes when popularity drops.
  • Live ingest (out of scope here, sketched in tradeoffs) — RTMP / SRT push → live transcoder → LL-HLS/DASH packager → CDN.

03Non-functional

What must it promise about speed, uptime and correctness?

  • Availability: 99.99% on the read path (playback). Reads are the product. 99.9% on creator upload (a 30-second hiccup is recoverable).
  • Latency: P99 startup < 2 s VOD / < 500 ms live; manifest fetch < 100 ms; license fetch < 200 ms; bitrate switch < 1 s.
  • Durability: Mezzanines are 11×9s durable across regions (50-year retention). Encoded variants can be re-derived from masters.
  • Consistency: Eventual is fine for catalogue reads (60 s TTL on the cache). Strong on entitlement check (revoked subs must stop within 30 s). Read-after-write within a region for fresh uploads (manifest references segments that exist).
  • Scalability: Horizontal at every layer. The bottleneck is egress at the ISP edge, which is solved by embedded OCAs (Open Connect) or multi-CDN (everyone else).
  • Security: Signed URLs prevent leech; DRM gates protected content; CMAF + CENC means one ciphertext serves all DRM systems (~66% storage savings vs per-DRM packaging).

04Capacity estimation

How much load and data does this have to hold?

Inputs: 90 M peak concurrent, 5 Mbps avg bitrate, 6 s segments, 500 hr/min upload (YouTube-class), 400 GB/hr mezzanine, 30x encoded ladder, 97% edge cache hit.

  • Segment GETs / sec: 90 M / 6 s = 15 M GET/s globally (peak). At LL-HLS 2 s segments: 45 M GET/s.
  • Peak egress: 90 M × 5 Mbps = 450 Tbps. This is roughly Netflix's published US-prime-time number (~15% of US internet) plus a global multiplier.
  • Origin offload: with 97% edge cache hit + 80% regional shield hit, origin sees (1 - 0.97) × (1 - 0.80) = 0.6% of total egress = ~2.7 Tbps at the origin. Without the shield, origin would see ~13 Tbps.
  • Mezzanine ingest: 500 hr/min × 525 600 min/yr × 400 GB/hr = 105 EB/yr at full ProRes. (YouTube uses lower-fidelity sources, so real ingest is ~7-10 EB/yr.)
  • Encoded ladder storage: ~30x ingest = ~3,150 EB/yr at full ProRes (30 × 105 EB), or ~210-300 EB/yr on YouTube's realistic lower-fidelity sources (30 × ~7-10 EB). Most variants decay to cold archive within 24 months, so the active hot tier holds ~2 EB rather than a year's full-ladder output.
  • Transcode jobs / sec: 500 hr/min upload × 60 min/hr = 30 K source-min/min = 500 source-min/sec ingested. Per-shot fan-out: 250 K work items per 30-min episode → ~8 333 work items per source-minute. Total: 500 × 8 333 ≈ 4.16 M jobs/sec at YouTube-class peak ingest. Solved by elastic spot autoscale + Argos VCU ASICs for AV1 hot path.
  • Live event surge: Tyson-Paul Nov 2024: ~65 M concurrent vs ~30 M sized → ~10x over-peak. Hotstar IPL 2023: 25 M concurrent in 30 s → ~30x join-rate over baseline. Sized for 5-10x sustained, 25x burst on join derivative.
  • DRM license fetches: 1 per session start + 0-5 renewals → at 90 M sessions starting per day, ~3 K/sec sustained, ~50 K/sec burst at premiere.

05API design

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

GET /v1/manifest/:titleId/:profile
Accept: application/vnd.apple.mpegurl, application/dash+xml
Authorization: Bearer <session-token>

200 OK
Content-Type: application/vnd.apple.mpegurl
Cache-Control: public, max-age=120

#EXTM3U
#EXT-X-VERSION:6
#EXT-X-STREAM-INF:BANDWIDTH=5800000,RESOLUTION=1920x1080,CODECS="avc1.640028"
https://edge.cdn.example/title/strangerthings/s05e08/avc1/1080p/playlist.m3u8?Expires=...&Signature=...
...
POST /v1/license/:titleId
Authorization: Bearer <session-token>
Content-Type: application/octet-stream    # Widevine/FairPlay/PlayReady CDM challenge (request body)

200 OK
Content-Type: application/octet-stream    # license response: content-key wrapped to the device

# entitlement failure — the endpoint's primary authz outcome (sub lapsed / geo-blocked / title pulled):
403 Forbidden   { "error": "not_entitled" }
# Two-step upload — the 50 GB of bytes go DIRECT to S3, never through the app tier.
# 1) init: app validates ownership + quota and mints presigned multipart PUT URLs.
POST /v1/uploads/init
Authorization: Bearer <creator-token>

{ "filename": "...", "contentType": "video/mp4", "sizeBytes": 53687091200 }   # 50 GB

200 OK
{
  "uploadId": "up_abc123",
  "partSizeBytes": 67108864,                  # 64 MB
  "parts": [
    { "partNumber": 1, "url": "https://s3.../?partNumber=1&uploadId=...&X-Amz-Signature=...", "expiresAt": "..." },
    { "partNumber": 2, "url": "https://s3.../?partNumber=2&uploadId=...&X-Amz-Signature=...", "expiresAt": "..." }
    // … one short-lived presigned PUT per part
  ]
}
# 2) client PUTs each part DIRECTLY to S3 (not to the app). S3 returns an ETag per part.
PUT <parts[i].url>
Content-Length: 67108864
<raw bytes>

200 OK
ETag: "9b2cf5..."
# 3) complete: client hands back the per-part ETags; app finalizes the S3 multipart
#    object and kicks off transcode/packaging.
POST /v1/uploads/up_abc123/complete
Authorization: Bearer <creator-token>

{ "parts": [ { "partNumber": 1, "etag": "9b2cf5..." }, { "partNumber": 2, "etag": "..." } ] }

201 Created
{ "assetId": "asset_77", "status": "encoding" }
# Play event (fire-and-forget, async telemetry)
POST /v1/events/play
{ "titleId": "...", "sessionId": "...", "rung": "1080p_avc1", "ts": "..." }

204 No Content

06Data model

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

Title (sharded by title_id via Vitess vindex):

fieldtypenotes
title_idbigintPK; sharding key
nametext
publishertext
disabledbooltakedown
rights_geojsonbper-country license window
created_attimestamp

Asset (sharded by title_id):

fieldtypenotes
asset_idbigintPK
title_idbigintFK; covering index
variantsjsonbrungs × codecs × DRM × signed-URL templates
statusenumencoding / packaging / published / takedown
published_attimestamp?when packager committed

Entitlement (sharded by account_id):

fieldtypenotes
account_idbigintPK; sharding key
title_idbigint
geochar(2)
valid_fromtimestamp
valid_totimestamp?

Popularity (sharded by region + title_id):

fieldtypenotes
regionchar(2)composite PK
title_idbigintcomposite PK
score_5mfloat5-min sliding window
score_1hfloat
score_24hfloat
score_7dfloat
tierenumedge / regional / hot / cold

Database choice — recommended:

Vitess (sharded MySQL) for the catalogue / asset / entitlement metadata. YouTube uses this in production. Why: 64 hash shards on title_id give linear scale; SERIALIZABLE single-shard transactions cover the writes (insert-asset, update-popularity, update-entitlement); Vitess vindex gives us pluggable sharding without app rewrites; the auditor escape-hatch (ad-hoc SQL for licensing queries) matters. Semi-sync replicas in-region give RPO ≤ 1 s; async cross-region follower for DR.

Spanner / CockroachDB is the right answer if you need globally-serializable writes (e.g. account creation must be unique across regions instantly). Costs ~3-5x Vitess at our QPS for marginal consistency gain.

DynamoDB is fine for KV access patterns but the takedown / compliance queries don't fit GSIs cleanly. Cassandra (Netflix's choice for some metadata) is a fine alternative — RF=3 LOCAL_QUORUM, but the takedown queries again push us toward SQL.

07High-level design

Which components handle a request, and in what order?

Architecture summary:

  1. Player (Web / Mobile / TV) runs an ABR algorithm (BBA / BOLA) that picks rungs by buffer level + throughput estimate. Fetches manifest, fetches segments, fetches DRM license, decrypts + renders.
  2. CDN Edge POP — Netflix Open Connect Appliance (FreeBSD + NGINX, 90 Gb/s TLS, 16 TB NVMe) embedded in ISPs; or CloudFront/Akamai/Fastly for non-Netflix-class scale. Holds the ~5-10% of catalogue popular in this country.
  3. Regional Shield — mid-tier cache (Akamai tiered cache / CloudFront Origin Shield / Open Connect storage tier). Sized 10x edge; absorbs edge misses; protects origin from cache-miss explosion.
  4. API Gateway (Envoy + WAF) — handles everything that misses the edge: manifest builds, license fetches, uploads, telemetry. mTLS internally; signed URLs minted here.
  5. Manifest Service — stateless; reads catalogue cache (95% hit) or metadata DB (cold path) to render HLS/DASH/CMAF manifests. Injects signed URLs; per-region rung trimming; per-device codec gating.
  6. DRM License Server (Widevine / FairPlay / PlayReady) — entitlement check + license minting in one box (latency budget too tight for an extra hop). Pre-authentication tokens for live-event surges.
  7. Catalogue Cache (EVCache — Memcached + Ketama) — hot title metadata + entitlement cache. Cache-aside; invalidated-not-written-through.
  8. Metadata DB (Vitess sharded MySQL) — source-of-truth catalogue. 64 hash shards × (1 primary + 2 semi-sync followers + 1 async cross-region follower).
  9. Origin Object Store (S3 / GCS / Colossus) — mezzanines + encoded ladder + packaged bundles. Cross-region asynchronously replicated.
  10. Cold Archive (Glacier Deep Archive / tape) — long-tail masters + retired variants. Lifecycle-promoted by the popularity worker.
  11. Upload Orchestrator (control plane only — Cosmos ingest pattern) — issues S3 multipart-upload pre-signed URLs; client uploads bytes DIRECTLY to S3; orchestrator handles init / complete and publishes encode-job (outbox).
  12. Transcode Bus (Kafka 3.7) — three logical topics: encode-jobs (2 048 partitions), chunk-ready (1 024 partitions), playback-events (1 024 partitions). ~12 K partition replicas total — within metadata-budget for one cluster, but tentpole growth may force per-topic clusters.
  13. Encoder Farm (FFmpeg + AV1/HEVC/H.264 + Argos VCU ASIC for AV1 hot path) — per-shot encoding × rungs × codecs.
  14. Packager (Bento4 / Shaka + CMAF + CENC) — stitches chunks, encrypts (one ciphertext serves all DRMs), publishes asset (outbox).
  15. Popularity / Tier Mover (Apache Flink) — scores titles per-region from playback events; drives hot/cold transitions.
  16. Encoder DLQ (Kafka tombstone topic) — dead-letter queue for poison chunks; bounded auto-resurrect when encoder ships a fix.
  17. Monitoring (Atlas + Mantis + Conviva/NPAW) — service-side metrics + client-side QoE; the SLOs every failureMode references live here.

Data flow on hot read (95% case): Client → CDN Edge → (NVMe hit) → 200 OK + segment bytes. ~25 ms p99.

Data flow on cold read (long-tail): Client → CDN Edge → Regional Shield → Object Store (S3) → 200 OK + segment bytes. ~600 ms p99.

Data flow on manifest build (cache hit): Client → CDN Edge → API Gateway → Manifest Service → Catalogue Cache → 200 OK + .m3u8/.mpd. ~120 ms p99.

Data flow on upload + transcode (direct-to-S3 control plane):

  1. Creator → API Gateway → Upload Orchestrator: POST /v1/uploads/init (validates ownership, mints signed multipart URLs).
  2. Creator → S3 directly via signed URLs: PUT each part. Bytes never touch the orchestrator.
  3. Creator → API Gateway → Upload Orchestrator: POST /v1/uploads/:id/complete (orchestrator calls S3 CompleteMultipartUpload, INSERTs asset row, PUBLISHes encode-job — outbox group).
  4. Encoder × N (per-shot fan-out) → S3 (encoded chunks) + Kafka (chunk-ready) → Packager → S3 (packaged bundle) + Metadata DB (asset published). ~30 min SLA short / hours long.

Data flow on hot/cold promote: Edge → Kafka (playback-events) → Flink (popularity) → Metadata DB (score) + S3 lifecycle (cold-archive PROMOTE). ~hours SLO.

Trace catalogue

The simulator has eight authored journeys for this canonical. Each is a deterministic walk through the architecture; click between them to compare hot vs cold paths, read vs write, and sync vs async lanes:

TraceJourneyBudget
youtube-streaming:edge-hitHot segment served at the OCA / CDN edge25 ms
youtube-streaming:cold-segment-origin-pullEdge miss → regional miss → origin S3 (the latency tail)600 ms
youtube-streaming:manifest-build-cache-hitManifest builder hits catalogue cache (95% path)120 ms
youtube-streaming:manifest-build-cold-dbCatalogue cache miss → Vitess / Spanner (cold path)280 ms
youtube-streaming:drm-license-fetchPlayer asks for Widevine / FairPlay license250 ms
youtube-streaming:creator-uploadFull upload (multi-segment): orchestration + direct-to-S3 bytes, rendered as parallel paths30 min
youtube-streaming:transcode-fanoutCosmos fan-out: encoder consumes → encodes → object30 min
youtube-streaming:cold-tier-promotePopularity worker promotes long-tail title from cold to hot1 h

> The creator-upload trace uses the simulator's Phase-10 multi-segment feature: one trace, two parallel BFS paths rendered as separate cards in the stepper. Segment 1 (Orchestration) walks client → gateway → upload → transcode-stream — init mints signed URLs, complete commits the asset row + outbox publish. Segment 2 (Bytes) walks client → object-hot directly via signed URL — the byte PUT that the orchestrator stays out of. One click, both paths visible. The encode pipeline that the outbox publish triggers lives in transcode-fanout.

Failure scenarios we model

The simulator's chaos panel ships eight scenarios for this canonical, each grounded in a cited real-world incident:

ScenarioCategoryReal-world precedent
cdn-edge-bgp-blackholeinfraFastly 2021-06-08 — BGP withdrawal blackholes ~85% of network for 50 minutes; BBC, Reddit, GitHub, Hulu offline.
s3-origin-region-outageinfraAWS S3 us-east-1 2017-02-28 + 2021-12-07 — multi-hour outage; Netflix survived via OCA cache, Hulu did not.
metadata-db-leader-failoverdataGitHub 2018-10-21 — 43 s WAN partition produced 24 h of follow-on remediation; teaches RTO/RPO trade-offs.
viral-release-hot-shardtrafficNetflix Open Connect — top 100 titles serve >50% of bytes; the canonical hot-key story.
drm-license-launch-overloaddepsDisney+ 2019-11-12 — Mandalorian launch-night DRM/login overload.
encoder-farm-tentpole-backlogprocessNetflix Cosmos / Squid Game S2 pattern — 250 K work items per 30-min episode × 8 episodes = queue blowout.
abr-client-regression-fleet-480pprocessNetflix Falcor / playback-team retros — bad client deploys regress ABR; egress drops, error rate stays zero, bitrate-delta is the only signal.
live-event-thundering-herdtrafficHotstar IPL 2019/2023 + Netflix Tyson-Paul 2024-11-15 — 25-65 M concurrent in a 30 s window.

08Deep dives

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

1. Per-title encoding economics. Per-title (Netflix 2015) and per-shot (Netflix 2018) encoding save 20-40% bitrate at constant VMAF. The encoder cost is real (a 30-min episode → ~250 K work items at full per-shot fan-out). Pays for itself once bytes_saved × egress_cost > extra_compute_cost × (encode_jobs - fixed_ladder_jobs), which Netflix flipped at ~5 K-title catalogue scale. YouTube's mix is different: H.264 ships first (minutes after upload), VP9 added once viewership picks up, AV1 reserved for top-watched titles where compute amortizes against billion-view egress savings.

2. Hot key on a viral release. A tentpole release can take >50% of global egress for a window. Naive consistent-hashing on titleId pins all the bytes to one OCA shard. Mitigations: per-segment sharding ((titleId, segmentId)) so a 2 hr movie is ~1200 keys, replicate top-1% to every edge, edge prefill of top-N titles, ABR drops players to a lower rung as the OCA outbound link saturates (rebuffer stays low at the cost of bitrate). Detection: per-key QPS (NOT global QPS), OCA outbound link by ASN.

3. ABR algorithm choice. BBA (buffer-based, Netflix 2014) is robust and dominant in production: pick the highest rung the buffer comfortably sustains. BOLA (Lyapunov-optimal, 2016) is provably optimal in QoE for some metrics but harder to tune. Pensieve / Puffer (RL-based, MIT 2017) wins on benchmark traces but needs to be retrained per network. The failure mode of ALL three: the high-water mark gets misread, the player oscillates → the fleet drops to lowest rung. Mitigation: bitrate-delta as the load-bearing client KPI (alert on it directly, not on error rate).

4. Multi-CDN steering. Hotstar runs three CDNs; Netflix runs only Open Connect for hot bytes. The second CDN is 1.5-3x the per-byte cost of the primary; pays for itself if a 50-minute Fastly-class outage costs more than the persistent overhead. Steering decision uses real-user-measurement (Conviva / NPAW) — which CDN has the lowest rebuffer in this device's ASN this minute?

5. Live event surge. The load-bearing metric is the JOIN derivative (sessions/sec joining), NOT absolute count. By the time count spikes, you've already saturated. Mitigations: pre-authenticated tokens (T-30 min mint short-lived tokens for known subscribers; T-0 they bypass auth), aggressive license caching on-device, manifest TTL drops to 1-2 s during the event (player stays fresh, cache-fronted within the window), dedicated reserved-capacity gateway pool for the event, Hotstar's "Connect-Stream-Connect" client protocol that explicitly tears down and reconnects rather than silently degrading.

6. DRM and key rotation. CMAF + CENC means one ciphertext serves Widevine (CTR), FairPlay (CBCS), PlayReady — saves ~66% storage vs per-DRM packaging. Key rotation is the failure mode: encoder ships ciphertext with key K1 before license server learns about it → every player decrypts gibberish. Mitigations: atomic key bundles (deploy encoder + license-server config together, never piecewise); ringed encoder rollouts gated on decryption-success SLO; license caching on-device (24 h TTL) so a brief desync doesn't affect existing sessions.

7. Hot/cold tier promote. When a back-catalogue title trends, the popularity worker has hours-not-days to promote it from Glacier Deep Archive to S3 Standard so the next OCA fill picks it up. Glacier expedited restore is 5 minutes / 10x cost; batch is 12 hours. The worker uses a 24 h trend signal to pre-emptively restore predicted candidates so the promote latency is hidden. Promote only re-encodes when the popularity tier requires a codec we didn't ship (e.g. an old back-catalogue title that's H.264-only suddenly trending → no re-encode unless AV1 is needed).

8. Regional origin failure. S3 us-east-1 has had two multi-hour outages (2017, 2021). What keeps playing: hot bytes (OCA NVMe is the source of truth). What dies: fresh uploads, encoder writes, cold-tail viewers, new launches. Mitigations: warm-warm cross-region (replicate masters + variants); route encoder writes to secondary region during the outage; do NOT flip read traffic (OCAs already serve); accept that fresh launches land in the secondary region for a few hours.

09Trade-offs

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

What breaks at 10x scale (Netflix-class → YouTube-class, ~500 hr/min upload, ~2 EB stored):

  • Encoder farm can't keep up on CPU-only — the Argos VCU ASIC flip becomes mandatory (~10x compute efficiency on H.264).
  • Per-shot encoding moves from "nice optimization" to "necessary" — fixed ladders waste too many bytes at tail catalogue scale.
  • Public CDN stops being economical past ~5% of any ISP's traffic — embedded ISP appliances (Open Connect) flip on.
  • Per-region popularity gets too coarse — sub-region (state / metro / ASN) popularity feeds dominate edge prefill.
  • Cold archive must move from S3 Glacier to true tape storage — Deep Archive cost is still ~$1/TB/mo, tape is ~$0.10/TB/mo at YouTube's scale.

Per-component failure stories:

  • CDN edge BGP blackhole (Fastly 2021) — multi-CDN steering with RUM signals saves you. Multi-CDN is expensive insurance; we underwrite it.
  • Origin region outage (S3 2017/2021) — OCAs absorb hot reads; fresh uploads + cold viewers + new launches break. Cross-region warm-warm.
  • Metadata DB primary failover — 30 s RTO with semi-sync replicas; manifest builds that hit cache survive, cold builds break for the gap. Cache hit ratio is the load-bearing number.
  • Catalogue cache hot shard (Squid Game S2) — per-segment sharding + edge prefill of top-N titles + manifest TTL randomization.
  • DRM license overload (Disney+ 2019) — pre-warm + regional sharding + pre-authenticated tokens + license caching.
  • Encoder backlog (tentpole drop) — priority queues + spot autoscale + DLQ for poison chunks.
  • ABR client regression — bitrate-delta SLO, staged rollout, server-side rung trim as belt-and-braces.
  • Live event surge (Tyson-Paul) — JOIN derivative monitoring; pre-auth tokens; Connect-Stream-Connect; reserved capacity.

Out of scope here (would extend the canonical):

  • Live ingest path (RTMP/SRT push → live transcoder → LL-HLS packager). Adds: ingest service, live transcoder pool with SLA against real-time, partial-segment producer for LL-HLS, DVR window manager.
  • Recommendations / search / ads. Separate systems with their own architectures.
  • Subtitles / dubbing / accessibility tracks. Multiplies asset count but doesn't change the architecture shape.
  • Content ID / copyright detection. YouTube-specific; runs as a separate fingerprint pipeline against the encoded ladder.

Primary sources

  • Netflix Tech Blog — Per-Title Encode Optimization (2015)
  • Netflix Tech Blog — Optimized Shot-Based Encodes (2018)
  • Netflix Tech Blog — Rebuilding Video Pipeline with Microservices (Cosmos, 2023)
  • Netflix Tech Blog — Serving 100 Gbps from a single Open Connect Appliance
  • Netflix Open Connect — Overview whitepaper + Fill Patterns docs
  • Google / YouTube — Reimagining YouTube Video Infrastructure (Argos VCU, 2021)
  • Apple — Enabling Low-Latency HLS
  • Hotstar Engineering — How we scaled to 25M concurrent for IPL
  • AWS — CloudFront Origin Shield in Multi-CDN Deployments
  • BBA paper — Stanford SIGCOMM 2014 (Buffer-Based ABR)
  • Pensieve — MIT SIGCOMM 2017 (RL-based ABR)

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 YouTube / Netflix Streaming yourself