System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

YouTube: streaming, cost optimisation and failure handling


YouTube: scope, scale and the upload path built the upload path, which is where the engineering is. The playback path is where the money is: 90 PB a day of egress, 150 bytes out for every byte in.

This lesson designs that path, then the optimisations that make the whole platform affordable, then what breaks and the questions the problem always attracts.

The streaming path

Adaptive bitrate, one segment at a timeFetch themanifest4 ssegment at 1080pBufferdrains,link slowsNextsegment at 480pRecovers,steps back upThe client picks the rendition; the server only serves files.
Chunking exists so the quality decision can be remade every four seconds instead of once per video.

Why video is chunked at all

A video is not delivered as one file. It is delivered as a manifest plus a sequence of short segments, typically 2 to 6 seconds each.

Three reasons, and none of them is optional:

  • Seeking. Jumping to minute seven means fetching the segments around minute seven, not downloading six minutes of video first.
  • Quality switching. Bitrate can change at any segment boundary, which is what makes adaptive streaming possible.
  • Cacheability. A segment is a small, immutable, individually addressable object — exactly the shape a content delivery network caches well. One large file with range requests caches far less predictably.

Adaptive bitrate streaming

Putting the decision on the client matters architecturally: the server serves static files. There is no session, no per-viewer state, and no server-side logic in the playback path — which is precisely what makes it fully cacheable at the edge.

The CDN does nearly all the work

Playback traffic is 7 Tbps average. Serving that from origin infrastructure means building a network with a peak capacity of ~15 Tbps and enough points of presence worldwide to keep start-up latency under two seconds. That is the business of a content delivery network.

Viewership is heavily concentrated: on most platforms, a small fraction of videos accounts for the large majority of views. Suppose 1% of videos generate 90% of watch time. Then a cache holding that 1% serves 90% of traffic, and the origin sees only the remaining 10% — a 10× reduction in origin bandwidth from caching alone.

Push that further and the origin's real job becomes: serve cache misses, absorb the long tail, and hold the authoritative copy.

The cost model, stated plainly

Three costs, in the order they matter:

CostRough shareDriver
Egress bandwidthDominantWatch-hours × bitrate
StorageSubstantial, and it only growsUploads × renditions × retention
Transcoding computeReal but boundedUploads × ladder size

Two consequences shape the design:

  1. Every percentage point of cache hit rate is money. Going from 90% to 95% halves origin egress. This is why placement, segment sizing, and cache lifetimes get real engineering attention.
  2. The encoding ladder is a cost decision, not a quality decision. Adding a 4K rendition to every video adds storage and transcoding cost for every video, and is watched by a small fraction of viewers. Real platforms encode expensive renditions on demand for videos that prove popular, rather than for everything on upload. The optimisations below develop this.

Speed and cost optimisations

Four optimisations, each with the arithmetic that justifies it. This is the part that separates a design that works from one that is affordable.

Four optimisations that pay for themselvesMaking itaffordableParallel transcodePre-signed upload URLsTiered cold storagePopularity placementOn-demand renditions
Most videos are watched almost never, so transcoding every rendition up front is the largest wasted cost.

Parallel transcoding

Encoding a 10-minute video to six renditions serially, at roughly real-time speed per rendition, takes about an hour of wall-clock time. Nobody waits an hour.

The directed acyclic graph from the upload path makes it parallel. Split the video into 100 six-second segments; each segment × rendition pair is independent.

600 tasks, each taking ~2 seconds of CPU for a 6-second segment

On 100 workers in parallel: 600 ÷ 100 × 2 s = 12 seconds of encoding, plus split, merge, and queueing overhead — call it 1 to 2 minutes end to end.

From an hour to two minutes, using the same total compute. The only thing that changed is that the work was expressed as a graph instead of a script.

Pre-signed URLs on both paths

Covered for upload in the upload path, and it applies symmetrically to download. Playback URLs are signed and time-limited, so access control is enforced at URL-issuing time rather than on every byte. The application tier's involvement in a 600 MB transfer is one signature.

Tiered storage

Storage grows by 1 PB per day and nothing is ever deleted. But access is extremely skewed: most videos are watched heavily in their first days and rarely afterwards.

Model it: suppose after 90 days a video receives under 1% of its lifetime views. Then videos older than 90 days — which after a few years are the overwhelming majority of the catalogue — can move to a storage class costing perhaps a quarter as much, at the price of higher retrieval latency and a per-retrieval fee.

If 80% of a 365 PB/year archive moves to storage costing 25% as much, the storage bill for that portion drops by 60%.

The trade-off is honest: a cold video's first play after archival is slower. Mitigate by keeping one low-bitrate rendition on fast storage for every video, so playback can start immediately at lower quality while higher renditions are retrieved. Exact class pricing varies by provider and changes, so present the structure rather than specific figures.

Popularity-based placement, and on-demand renditions

Two related ideas that both follow from skewed viewership.

Do not push everything to the edge. Pre-populating every edge cache with every video is impossible — the catalogue is exabytes and an edge holds terabytes. Instead, push the predicted top slice (new uploads from large channels, regionally trending content) and let everything else arrive by origin pull on first request. The first viewer in a region pays a slower start; the rest hit cache.

Do not encode everything at every rendition. Generate the cheap, widely used renditions on upload; generate expensive ones (4K, and computationally heavy modern codecs) only once a video crosses a view threshold. If 95% of videos never cross it, this removes most of the expensive half of the transcoding bill.

Failure handling and follow-ups

What breaks, and the five questions this problem always attracts.

Resuming an upload that droppedUpload in 5MB chunksServerrecordsthe offsetConnectiondies midwayClient asksthe offsetResume, donot restartOne failed transcode task retries alone; the finished chunks stay finished.
Chunking makes failure cost a chunk instead of a file, and the same idea makes transcoding independently retryable.

Resuming a failed upload

The client uploaded 40 of 60 chunks and the connection died. On reconnect it calls a status endpoint, receives the list of chunks the server has, and sends the missing 20. Two details make this reliable:

  • Chunks are content-addressed or explicitly indexed, so "chunk 41" is unambiguous and a duplicate upload of the same chunk is idempotent — writing it twice produces the same result.
  • Incomplete uploads expire. An abandoned upload holds storage indefinitely otherwise. A lifecycle rule deleting incomplete multipart uploads after, say, seven days is a standard and frequently forgotten piece of housekeeping.

Retrying one failed transcode task

This is the payoff for the graph structure. Task 417 of 600 fails because a worker was terminated mid-encode.

The orchestrator re-queues that one task. The other 599 outputs are untouched. For this to be safe, each task must be idempotent and keyed by its inputs: the output object is named deterministically from (video_id, segment_index, rendition), so a re-run overwrites its own output and nothing else. Compare with a monolithic transcode job, where a failure 55 minutes in means starting over.

Persistent failures — a corrupt source segment, an unsupported codec path — go to a dead-letter queue after a bounded number of attempts, and the video is marked partially processed with the renditions that did succeed still playable.

Live streaming, and why it is a different system

Live inverts the latency budget. Batch processing tolerates minutes; live tolerates seconds end to end.

What changes:

  • No split-and-merge graph. Segments are encoded as they arrive, in order, continuously. There is no complete file to split.
  • A shorter ladder. Fewer renditions, because encoding must keep pace with real time.
  • Segment duration becomes the latency floor. Six-second segments impose at least six seconds of delay; low-latency variants of the streaming formats use shorter segments or partial-segment delivery to cut it, at the cost of more requests and lower cache efficiency.
  • The origin is under sustained write load rather than serving an immutable archive.

Say plainly that live shares the delivery network and shares almost nothing else.

Deduplication of re-uploads

The same video uploaded a million times — a viral clip, or a re-upload of copyrighted content — should not be stored a million times, and in the copyright case should not be published at all.

Exact-duplicate detection by cryptographic hash of the source file is trivial and nearly useless, because a one-frame trim changes the hash. What works is perceptual fingerprinting: derive a compact signature from the visual and audio content that survives re-encoding, cropping, and minor edits, then match against an index of known fingerprints. Describe it at this level; the algorithms are a specialised field and the deployed systems are proprietary.

Content moderation

Automated classification on upload flags a fraction for human review; the rest publish immediately. Because scanning is a side branch of the graph, it runs in parallel with encoding rather than delaying it. The policy question — publish-then-review or review-then-publish — is a product decision with a direct architectural consequence, and it is worth naming as such rather than assuming.