System design · Transcoding pipeline, object storage, CDN

How to design YouTube

YouTube is two systems that share a database. One takes in large video files, slowly and in the background, and turns each one into a dozen playable versions. The other serves those versions to a huge audience, fast, from servers close to every viewer. The interview is about keeping those two paths apart.

Updated · 7 min read

Requirements

Agree the scope first. YouTube has hundreds of features; pick the video path and say clearly what you're leaving out.

  • Functional: upload a video; watch a video on any device and network; see title, description and view count.
  • Functional, optional: search, comments, likes, subscriptions, recommendations.
  • Non-functional: playback should start within about two seconds and rarely stall, even on a slow mobile network.
  • Non-functional: an uploaded video must never be lost, but it can take minutes to become watchable. Availability of playback matters more than freshness of metadata.

Capacity estimates

These are assumptions, stated out loud. The point is the order of magnitude, which decides the architecture.

QuantityAssumptionResult
Upload volume500 hours of video uploaded per minute≈ 720,000 hours a day
Original storage≈ 1 GB per hour of uploaded source video≈ 720 TB a day of originals
Transcoded storagethe full rendition ladder adds about the same again≈ 1.4 PB a day ≈ 500 PB a year, before replicas
Transcoding compute≈ 5 CPU-hours to encode one hour of video into every rendition30,000 video hours an hour × 5 ≈ 150,000 cores busy around the clock
Views1B views a day≈ 11,600/s average, ≈ 35,000/s at a 3× peak
Egress5 minutes watched per view at ≈ 3 Mbps≈ 110 MB per view ≈ 110 PB a day ≈ 10 Tbps average

Three conclusions. Storage grows by petabytes a day, so video bytes belong in object storage, never in a database. Transcoding is a large, steady compute job that must run asynchronously on a worker fleet. And 10 Tbps of egress can't come from one region - a CDN isn't an optimisation here, it's the delivery system.

API design

Clients never stream large files through your application servers. The API hands out signed URLs and the bytes go straight to and from storage.

POST /api/videos
→ { "title": "My trip", "description": "...", "size_bytes": 2147483648 }
← 201 { "video_id": "v_8f3k", "upload_url": "https://uploads.example.com/...signed", "part_size": 16777216 }

PUT {upload_url}?part=1..n            multipart, resumable
POST /api/videos/{id}/complete        → 202 { "status": "processing" }

GET /api/videos/{id}
→ 200 { "title": "My trip", "status": "ready", "views": 10432,
        "manifest_url": "https://cdn.example.com/v_8f3k/master.m3u8" }

POST /api/videos/{id}/views           → 202

Multipart, resumable upload matters: a 2 GB file over a phone connection will be interrupted, and restarting from zero loses users. With S3 this is a presigned multipart upload, so retries are per part.

Data model

Metadata is small and relational-ish; video bytes are huge and immutable. Keep them in different stores. Metadata fits PostgreSQL (RDS) sharded by video ID, or a key-value store such as DynamoDB. Video files live in object storage such as S3, keyed by video and rendition.

videos
  video_id        string     primary key
  owner_id        bigint
  title           text
  description     text
  status          enum       uploading | processing | ready | failed
  duration_s      int
  created_at      timestamp

renditions
  video_id        string     partition key
  rendition       string     sort key, e.g. "1080p-h264"
  manifest_key    string     object storage key of the playlist
  segment_prefix  string     object storage prefix of the segments

view_counts
  video_id        string     primary key
  views           bigint     updated in batches, not per view

Keep views out of the videos row. A popular video gets thousands of views a second, and incrementing one row that often turns it into a hot spot that also blocks reads of the title.

High-level design

Two paths: a slow, durable upload path and a fast, cached watch path.

  • Upload: the API service creates the metadata row with status "uploading" and returns a presigned URL. The client uploads the original straight into an uploads bucket in S3.
  • Processing: when the upload completes, an event goes onto a queue (SQS or Kafka). Transcoding workers - EC2 instances running FFmpeg, or a managed service such as AWS Elemental MediaConvert - pull jobs, produce each rendition as short segments, and write them to a media bucket.
  • When every rendition is written, the worker sets status to "ready". Metadata is cached in Redis, because one video's title is read millions of times.
  • Watch: the player fetches metadata from the API, then fetches the manifest and segments from a CDN such as CloudFront, which pulls from the media bucket on a miss.

Where it breaks

Transcode synchronously inside the upload request, and the API servers spend minutes of CPU per video while holding a connection open. A burst of uploads exhausts them, uploads time out, and the watch path - on the same servers - goes down with them.

The fix is the queue. Accept the upload, enqueue a job, return 202. Workers drain the queue at their own pace and scale on queue depth. A burst becomes a longer wait before a video is watchable, not errors. The new risk is backlog: if the queue grows faster than workers drain it, "processing" turns into hours.

The second break is egress. Serve video from origin and a single viral video saturates the bucket's region. The CDN absorbs that, as long as the hit rate is high - so segments must be cacheable, with long cache lifetimes and stable URLs.

The transcoding pipeline

  • Split before encoding. Cut the original into chunks of a few seconds at keyframes, fan the chunks out as separate jobs, then stitch the results. A two-hour video finishes in minutes instead of hours, and a failed chunk retries alone.
  • Make jobs idempotent. A queue delivers at least once, so a worker that crashes mid-job will see the job again. Write outputs to deterministic keys so a retry overwrites instead of duplicating.
  • Use a dead-letter queue for files that fail repeatedly - corrupt uploads exist - so one bad video can't block the pipeline.
  • Prioritise. Encode a low rendition first and mark the video watchable, then fill in higher resolutions.

Adaptive bitrate streaming

Each video is encoded into a ladder of renditions - say 240p up to 1080p or 4K - each cut into segments of two to six seconds. A manifest (HLS .m3u8 or DASH .mpd) lists them. The player measures its own download speed and picks the next segment's quality accordingly, so a viewer on a train drops to 360p instead of stalling.

The trade-off to name: more renditions mean smoother playback and more storage and compute. A common answer: encode the full ladder only for videos that get traffic.

View counts and recommendations

Don't write one database update per view. Send view events to a stream, aggregate them in memory or in a stream processor, and flush a sum per video every few seconds. The count is slightly stale, which nobody notices; the database sees thousands of times fewer writes.

Recommendations are usually out of scope. If asked, say they're computed offline from watch history into a per-user list, stored in a key-value store, and served from there - and move on.

What interviewers look for

  • Separating video bytes (object storage plus CDN) from metadata (a database plus cache).
  • An asynchronous transcoding pipeline: queue, workers that scale on backlog, idempotent jobs, retries and a dead-letter queue.
  • Direct-to-storage, resumable uploads instead of streaming files through application servers.
  • Adaptive bitrate streaming explained in a sentence or two, and the storage-versus-quality trade-off of the rendition ladder.
  • Batched view counts, and a clear statement of what you scoped out.

Frequently asked questions

Why not store videos in the database?

+

Videos are large, immutable files read sequentially. Object storage such as S3 is built for that: cheap per gigabyte, durable, and easy to put a CDN in front of. A database is built for small rows and queries, and would be slow and expensive at petabyte scale.

Why does transcoding need a queue?

+

Transcoding takes minutes of CPU per video and uploads arrive in bursts. A queue lets the upload return immediately, lets workers process jobs at a steady rate, and retries jobs when a worker fails. Without it, a burst of uploads overloads the servers handling requests.

What is adaptive bitrate streaming?

+

The video is encoded at several qualities and cut into short segments. The player downloads a manifest listing them, measures its bandwidth, and picks a quality for each segment. When the network slows, it switches down instead of stalling.

How does a CDN help a video platform?

+

Video egress is far larger than any single region can serve. A CDN caches segments at edge locations close to viewers, so popular videos are served from the edge and origin only handles misses.

How do you count views at scale?

+

Send each view as an event to a stream, aggregate counts per video for a few seconds, and write the totals in batches. That avoids a hot row per popular video, and the small delay in the displayed count is acceptable.

Now break one yourself.

The first challenge takes about two minutes. No signup.