Problem and scope
A video streaming platform lets people upload videos and lets other people watch them on phones, laptops and TVs. YouTube is the classic example of a platform where anyone can upload. Netflix is the classic example where a company uploads a curated catalogue and millions watch. The internal machinery is very similar: a video goes in once, gets converted into many versions, and is then served many millions of times from servers close to the viewers.
This problem is popular in interviews because it touches almost every building block at once. It has large binary files (object storage), slow CPU-heavy work (a processing pipeline with queues), extremely high read bandwidth (CDNs and caching), a metadata database, counters that see heavy write traffic, and a cost story that dominates the business. Interviewers typically probe four things:
- How do you upload a 4 GB file reliably over a flaky mobile connection?
- What happens between "upload finished" and "video is playable"?
- How does the player switch quality when the network gets slower?
- How do you serve tens of terabits per second without going bankrupt?
We will design the system end to end, with numbers, and then dig into those four questions.
Clarifying questions
Spend the first few minutes agreeing on scope. Good questions, with the answers we will assume:
| Question | Assumed answer |
|---|---|
| User-generated (YouTube) or curated catalogue (Netflix)? | User-generated, anyone can upload |
| Live streaming too, or only on-demand? | On-demand only; live is a follow-up |
| Which devices? | Web, Android, iOS, smart TVs |
| Maximum upload size and length? | Up to 10 GB, up to 4 hours |
| Do we need search, comments, likes? | Mention them, but focus on upload and playback |
| Recommendations? | High level only |
| Global or one country? | Global, with a large share of viewers in India |
| DRM (copy protection)? | Mention it; not a deep dive |
| Is a short delay before a new video is playable acceptable? | Yes, minutes are fine |
That last answer matters a lot. Because a delay of a few minutes is acceptable, processing can be asynchronous: the upload request returns quickly and a background pipeline does the heavy work.
Functional and non-functional requirements
Functional requirements (what the system does):
- A creator can upload a video, with title, description and thumbnail.
- An upload can be paused and resumed, and survives network drops.
- The system converts each upload into several resolutions and formats.
- A viewer can watch a video on any supported device, starting quickly.
- Playback adapts to the viewer's network speed.
- The system shows view counts, and a home feed with recommendations.
Non-functional requirements (how well it does it):
- Availability: playback should be highly available. A failed upload is annoying; a failed playback for millions of viewers is a crisis.
- Low start-up latency: the first frame should appear within about 2 seconds for most viewers.
- Smooth playback: few rebuffering events. Rebuffering means the player has run out of downloaded video and freezes while it waits.
- Durability: an uploaded video must never be lost. Losing a creator's only copy is unacceptable.
- Scalability: read bandwidth is enormous and spiky (a viral video, a cricket highlight).
- Cost efficiency: storage and bandwidth are the largest bills, so the design must keep them in check.
- Eventual consistency is fine for view counts, likes and recommendations.
Back-of-the-envelope estimates
Estimates tell you which part of the system is hard. State your assumptions clearly; the interviewer cares more about the reasoning than the exact inputs.
Assumptions:
- 1 million new videos uploaded per day.
- Average video length: 5 minutes (300 seconds).
- Original uploads average about 10 Mbps (megabits per second), a typical 1080p phone recording.
- 200 million daily active viewers, each watching 30 minutes per day.
- Average delivered bitrate: 2.5 Mbps (most people watch on phones at 480p or 720p).
- Each viewer starts about 10 videos per day, so 2 billion views per day.
Upload storage. One original video is 300 s × 10 Mbps = 3,000 megabits. Divide by 8 to get bytes: 375 MB. For 1 million uploads that is 375 TB of originals per day.
Transcoded storage. Suppose we produce a "bitrate ladder" of five renditions. A rendition is one encoded version of the video at a specific resolution and bitrate. Our ladder:
| Rendition | Bitrate |
|---|---|
| 240p | 0.4 Mbps |
| 360p | 0.8 Mbps |
| 480p | 1.4 Mbps |
| 720p | 2.8 Mbps |
| 1080p | 5.0 Mbps |
The ladder adds up to 10.4 Mbps. For one 300-second video that is 300 × 10.4 = 3,120 megabits, or 390 MB. So transcoded copies add another 390 TB per day, which is about 142 PB per year. Storage grows forever because videos are rarely deleted. This is why storage tiering (discussed later) matters.
Upload bandwidth. 375 TB per day into our system is 375 × 10¹² × 8 bits ÷ 86,400 seconds ≈ 35 Gbps on average. That is significant but manageable across many regions.
Playback bandwidth (the big one). 200 million viewers × 1,800 seconds × 2.5 Mbps = 9 × 10¹⁷ bits per day. Divided by 86,400 seconds, that is about 10.4 Tbps on average. Evenings are busier, so peak could be twice that, about 21 Tbps.
No single data centre serves that. This number alone tells you that a CDN is not optional. A CDN (content delivery network) is a large set of cache servers placed near viewers, often inside internet service providers' networks.
Origin load. If the CDN serves 95% of bytes from its caches (a 95% hit ratio), the origin still sends 5% of 10.4 Tbps ≈ 520 Gbps. At a 99% hit ratio, it drops to about 104 Gbps. Every percentage point of cache hit ratio is worth around 100 Gbps of origin traffic. That is why CDN design gets a deep dive.
View count writes. 2 billion views per day ÷ 86,400 ≈ 23,000 view events per second on average, with peaks perhaps three times higher. That is too many to apply as individual UPDATE videos SET views = views + 1 statements on a single hot row.
Transcoding compute. 1 million videos × 5 minutes = 5 million minutes of video per day. Suppose producing the whole ladder costs about 10 CPU-core-minutes per minute of video (a rough assumption; real costs vary hugely with codec and preset). That is 50 million core-minutes per day, or 50,000,000 ÷ 1,440 ≈ 35,000 CPU cores busy around the clock. Transcoding is a large, elastic batch workload.
Segments. Streaming formats cut each rendition into short segments. With 4-second segments, a 300-second video has 75 segments per rendition, so 375 segment files for five renditions. One 1080p segment is 4 s × 5 Mbps = 20 megabits = 2.5 MB; one 240p segment is about 200 KB.
Interview tip
Say out loud what the numbers imply: "Playback bandwidth is around 10 Tbps, three hundred times the upload bandwidth, so the design is dominated by delivery. I will put most of my effort into the CDN and the streaming format, and treat upload as a reliable but asynchronous pipeline."
API design
We need a small number of endpoints. Uploads use a separate, resumable protocol because a single huge HTTP request is fragile.
POST /v1/uploads
body: { "fileName": "trip.mp4", "sizeBytes": 734003200,
"contentType": "video/mp4" }
201 -> { "uploadId": "upl_7f3a", "chunkSizeBytes": 8388608,
"uploadUrls": "...signed URLs per part..." }
PUT <signed object-storage URL for part N> (raw bytes)
200 -> { "etag": "..." }
GET /v1/uploads/upl_7f3a
200 -> { "receivedParts": [1,2,3,5], "status": "IN_PROGRESS" }
POST /v1/uploads/upl_7f3a/complete
body: { "parts": [ { "n": 1, "etag": "..." }, ... ],
"title": "Goa trip", "description": "...",
"visibility": "PUBLIC" }
202 -> { "videoId": "vid_91c2", "status": "PROCESSING" }
GET /v1/videos/vid_91c2
200 -> { "title": "...", "status": "READY", "durationSec": 300,
"manifestUrl": "https://cdn.example.com/v/vid_91c2/master.m3u8",
"thumbnailUrl": "...", "views": 10482 }
POST /v1/videos/vid_91c2/views
body: { "sessionId": "...", "watchedSec": 42 }
202 -> {}
GET /v1/feed?cursor=...
200 -> { "items": [ ... ], "nextCursor": "..." }
A few design choices are worth explaining:
- Signed URLs. A signed URL is a temporary link to object storage that includes a cryptographic signature and an expiry time. The client uploads bytes directly to storage, so our API servers never carry the video data. They only hand out permissions.
202 Acceptedon complete means "I accepted the work, but it is not done yet". The client pollsGET /v1/videos/{id}or receives a notification when the status becomesREADY.- The manifest URL points at the CDN, not at our API. Playback traffic never touches application servers.
- View events are fire-and-forget (
202). The client does not wait for the count to update.
Data model and storage choice
There are three very different kinds of data, and each gets its own store.
1. Video bytes (originals, renditions, segments, thumbnails) go in object storage. Object storage (Amazon S3, Google Cloud Storage, Azure Blob, or an in-house equivalent) stores immutable blobs addressed by a key, with very high durability through replication or erasure coding. Erasure coding splits data into pieces plus extra parity pieces so that the object can be rebuilt even if several pieces are lost, using less space than full copies. Object storage is cheap per GB, scales without limit for practical purposes, and pairs well with CDNs.
Example key layout:
originals/vid_91c2/source.mp4
renditions/vid_91c2/720p/seg_00001.m4s
renditions/vid_91c2/720p/seg_00002.m4s
renditions/vid_91c2/720p/init.mp4
renditions/vid_91c2/720p/playlist.m3u8
renditions/vid_91c2/master.m3u8
thumbs/vid_91c2/default.jpg
2. Video metadata goes in a relational or wide-column database. Metadata is small, structured, and queried by id and by channel.
CREATE TABLE videos (
video_id CHAR(12) PRIMARY KEY,
channel_id CHAR(12) NOT NULL,
title VARCHAR(200) NOT NULL,
description TEXT,
status VARCHAR(12) NOT NULL, -- UPLOADING, PROCESSING, READY, FAILED, BLOCKED
visibility VARCHAR(10) NOT NULL, -- PUBLIC, UNLISTED, PRIVATE
duration_sec INT,
created_at TIMESTAMP NOT NULL
);
CREATE INDEX idx_videos_channel ON videos (channel_id, created_at DESC);
CREATE TABLE renditions (
video_id CHAR(12) NOT NULL,
rendition VARCHAR(8) NOT NULL, -- 240p ... 1080p
codec VARCHAR(8) NOT NULL, -- h264, vp9, av1
bitrate_kbps INT NOT NULL,
storage_key VARCHAR(255) NOT NULL,
PRIMARY KEY (video_id, rendition, codec)
);
At our scale (hundreds of millions of videos per year), a single relational database will not hold all rows. You would shard by video_id, or use a distributed SQL store or a wide-column store such as Cassandra or Bigtable. The access pattern is almost entirely "fetch by key", which shards cleanly. The channel_id index serves "list this channel's videos, newest first".
3. Counters and analytics go in a separate path. View counts, watch time and likes are written as events into a log (Kafka or similar), aggregated, and stored in a counter store. We cover this in a deep dive.
Caching. Metadata for popular videos is read far more often than it changes, so a cache such as Redis sits in front of the metadata database. The video page reads title, counts and manifest URL from the cache.
High-level design
UPLOAD PATH
+---------+ 1 create upload +-------------+
| Creator |------------------->| Upload API |
| app |<-- signed URLs ----| service |
+---------+ +------+------+
| 2 PUT parts | 4 status
v v
+-----------------+ 3 done +-------------+
| Object storage |----------->| Metadata DB |
| (originals) | event +-------------+
+-----------------+ |
v
+-------------------+
| Transcode |
| orchestrator |
+---------+---------+
| tasks
v
+-------------------+ +------------+
| Task queue |---->| Worker |
+-------------------+ | fleet |
+-----+------+
| segments,
| manifests
v
+-------------------+
| Object storage |
| (renditions) |
+---------+---------+
PLAYBACK PATH | origin
+---------+ page +-----------+ v
| Viewer |-------->| Video API | +-----------------+
| player | | + cache | | Origin shield |
+----+----+ +-----------+ +--------+--------+
| manifest + segments |
v v
+-----------------+ miss +-----------------------+
| CDN edge (near |--------->| CDN regional tier |
| the viewer) | +-----------------------+
+-----------------+
| view events
v
+-----------------+ +-------------+ +---------------+
| Event ingestion |--->| Kafka log |--->| Aggregators, |
+-----------------+ +-------------+ | counters, ML |
+---------------+
The design has two almost independent halves. The upload path is write-heavy, slow, and asynchronous. The playback path is read-heavy, fast, and served almost entirely by the CDN. Keeping them separate means a backlog of transcoding jobs never slows playback.
Request flows step by step
Upload flow
- The creator's app calls
POST /v1/uploadswith the file size. The upload service creates an upload record with statusUPLOADING, starts a multipart upload in object storage, and returns anuploadId, a chunk size (say 8 MB) and signed URLs for the parts. - The app splits the file into 8 MB parts and uploads them directly to object storage, several in parallel. Each successful part returns an ETag, a checksum-like identifier of that part.
- If the network drops, the app later calls
GET /v1/uploads/{id}to learn which parts arrived, and uploads only the missing ones. - When all parts are in, the app calls
completewith the list of parts and the video's title. Object storage stitches the parts into one object. The upload service creates thevideosrow with statusPROCESSINGand publishes aVideoUploadedevent. - The transcode orchestrator receives the event and builds a job graph (see deep dive). Workers process the tasks and write segments and manifests to object storage.
- When every required task has finished, the orchestrator marks the video
READY, writes therenditionsrows, invalidates the metadata cache, and notifies the creator.
Playback flow
- The viewer opens the video page. The client calls
GET /v1/videos/{id}. The video API answers from cache: title, counts, thumbnail and manifest URL. - The player downloads the master manifest from the CDN. This small text file lists the available renditions and their bitrates.
- The player picks a starting rendition, usually a modest one so the first frame appears quickly, and downloads that rendition's media playlist, which lists the segment URLs.
- The player downloads the first few segments from the CDN edge. If the edge has them cached (a hit), they come back in milliseconds. If not (a miss), the edge asks a regional tier, which may ask the origin shield, which finally reads from object storage.
- As segments arrive, the player measures download speed and buffer level and switches renditions up or down for the next segment.
- Every so often (for example at start, at 30 seconds, and at the end), the player sends a view or heartbeat event. These flow into the event pipeline, not into the metadata database.
Deep dive 1: resumable uploads
A 4 GB upload over a phone connection will almost certainly be interrupted. If the protocol is "one big HTTP POST", every interruption restarts from zero. Resumable uploads fix this with three ideas:
- Split into parts. The client cuts the file into fixed-size parts (5 to 64 MB is typical). Each part is uploaded independently and can be retried on its own. Object storage multipart upload APIs (for example S3 multipart upload) are designed for this.
- Server-side progress. The server, not the client, is the source of truth for which parts arrived. After a crash or app restart, the client asks and continues.
- Integrity checks. Each part carries a checksum, such as MD5 or SHA-256. Storage rejects a part whose bytes do not match, so corruption is caught early.
A worked example: a 734,003,200-byte file (700 MiB) with 8 MiB parts (8,388,608 bytes) has 734,003,200 ÷ 8,388,608 = 87.5, so 88 parts: 87 full parts and a final half part. If the connection dies after part 60, only 28 parts remain to send.
Two more details interviewers like:
- Abandoned uploads. Some users never finish. A lifecycle rule deletes incomplete multipart uploads after, say, 7 days, so orphaned parts do not pile up as storage cost.
- Upload directly to storage, near the user. Signed URLs can point to a storage endpoint in the nearest region, or to an upload acceleration service, so long-distance TCP round trips do not slow the transfer.
Common mistake
Routing upload bytes through your own API servers. They become a bandwidth bottleneck and need huge memory buffers. Let the API do authentication and bookkeeping, and let clients send bytes straight to object storage with short-lived signed URLs.
Deep dive 2: the transcoding pipeline
Transcoding means decoding the uploaded video and re-encoding it into different resolutions, bitrates and codecs. A codec is the compression algorithm, such as H.264 (plays almost everywhere), VP9, or AV1 (smaller files, but more expensive to encode and not supported by every old device). Creators upload in dozens of formats; viewers need a predictable set.
The pipeline as a DAG
The work is naturally a DAG (directed acyclic graph): a set of tasks where some tasks depend on others, with no cycles.
+-------------+
| validate & |
| probe |
+------+------+
|
+-----------+-----------+
v v
+-------------+ +---------------+
| split into | | extract audio |
| chunks | +-------+-------+
+------+------+ |
| |
+-------+--------+-----+ |
v v v v |
+-----+ +-----+ +-----+ +-----+ |
|240p | |480p | |720p | |1080p| |
|enc | |enc | |enc | |enc | |
+--+--+ +--+--+ +--+--+ +--+--+ |
+-------+-------+-------+------+
v
+------------------+
| package segments |
| + write manifests|
+--------+---------+
v
+------------------+ +-------------+
| thumbnails, | | content |
| captions (opt.) | | moderation |
+--------+---------+ +------+------+
+--------+-----------+
v
mark READY
Steps:
- Validate and probe. Check the file really is a video, read its duration, resolution and frame rate, and reject corrupt files early.
- Split into chunks. Cut the video into pieces of, say, 10 to 60 seconds at keyframe boundaries. A keyframe is a frame that can be decoded on its own without earlier frames, so cutting there keeps each chunk independently decodable.
- Encode in parallel. Each (chunk, rendition) pair is an independent task. A 300-second video cut into 10 chunks with 5 renditions gives 50 tasks that can run on 50 machines at once. That turns a long serial job into a short parallel one.
- Package. Join encoded chunks and cut them into streaming segments, then write the manifests.
- Side tasks. Thumbnails, automatic captions and content moderation (copyright matching, policy checks) run alongside.
Orchestrating the DAG
An orchestrator stores the DAG for each video in a database, and pushes ready tasks (tasks whose dependencies are all done) onto a queue. Workers pull tasks, run them, write output to object storage, and report completion. The orchestrator then enqueues newly unblocked tasks.
Important properties:
- Idempotent tasks. A task may run twice (a worker died after finishing but before reporting). Write outputs to deterministic keys such as
renditions/vid_91c2/720p/chunk_03.mp4, so a second run simply overwrites the same file with identical content. - Leases and retries. A worker takes a lease on a task (for example 10 minutes). If it does not report back before the lease expires, the task goes back on the queue. After a few failures, the task goes to a dead-letter queue for a human or a different worker type to inspect.
- Priorities. A popular creator's upload, or the first low-resolution rendition, can go on a higher-priority queue. A common trick: produce 360p first and mark the video playable at low quality, then add higher renditions as they finish.
- Elastic workers. Transcoding is CPU or GPU heavy and bursty. Workers can run on cheap interruptible machines (spot or preemptible instances) because the lease-and-retry model tolerates a machine vanishing.
Interview tip
Mention "per-title encoding": instead of one fixed ladder for every video, analyse the content first. A simple animation compresses well, so it gets lower bitrates for the same visual quality; a fast sports clip needs more. This saves storage and bandwidth at the cost of extra analysis.
Deep dive 3: adaptive bitrate streaming
Adaptive bitrate streaming (ABR) means the player can switch between renditions every few seconds, picking the best quality the current network can sustain. Two formats dominate:
- HLS (HTTP Live Streaming), defined by Apple. Manifests are
.m3u8text files. - MPEG-DASH (Dynamic Adaptive Streaming over HTTP), an international standard. Manifests are
.mpdXML files.
Both work the same way: each rendition is a series of short segments, and a manifest tells the player where they are. Modern deployments often use CMAF (Common Media Application Format), a fragmented MP4 segment format both HLS and DASH can reference, so you store one set of segments instead of two.
A simplified HLS master manifest:
#EXTM3U
#EXT-X-STREAM-INF:BANDWIDTH=400000,RESOLUTION=426x240
240p/playlist.m3u8
#EXT-X-STREAM-INF:BANDWIDTH=1400000,RESOLUTION=854x480
480p/playlist.m3u8
#EXT-X-STREAM-INF:BANDWIDTH=2800000,RESOLUTION=1280x720
720p/playlist.m3u8
#EXT-X-STREAM-INF:BANDWIDTH=5000000,RESOLUTION=1920x1080
1080p/playlist.m3u8
And a media playlist for one rendition:
#EXTM3U
#EXT-X-VERSION:7
#EXT-X-TARGETDURATION:4
#EXT-X-MAP:URI="init.mp4"
#EXTINF:4.0,
seg_00001.m4s
#EXTINF:4.0,
seg_00002.m4s
#EXT-X-ENDLIST
How the player decides
The player keeps a buffer: seconds of video downloaded but not yet shown. The decision loop runs once per segment:
buffer (seconds of video ready to play)
0s 10s 20s 30s
|---------|------------|------------|
danger steady comfortable
-> step -> keep -> try a
down quality higher one
Two families of algorithms exist, and real players combine them:
- Throughput-based: measure how fast the last few segments downloaded, and pick the highest rendition whose bitrate is safely below that (for example below 80% of measured throughput).
- Buffer-based: pick quality from the buffer level. A low buffer means play it safe; a full buffer means you can afford to try higher quality.
Worked example. The player is on 720p (2.8 Mbps). The last 4-second 720p segment is 1.4 MB, which is 11.2 megabits. It took 2 seconds to download, so measured throughput is 11.2 ÷ 2 = 5.6 Mbps. With an 80% safety margin, the usable budget is 4.48 Mbps. The ladder's highest rendition below 4.48 Mbps is 720p (1080p needs 5.0), so the player stays at 720p. If the next segment downloads in 1.5 seconds, throughput is 11.2 ÷ 1.5 ≈ 7.47 Mbps, and 80% of that is about 5.97 Mbps, which allows 1080p. A good player also checks the buffer before stepping up, so one fast segment does not cause flip-flopping.
Why segments, and how long?
Segments let the player switch quality at segment boundaries and let CDNs cache small, independent files with plain HTTP. Segment length is a trade-off:
| Shorter segments (2 s) | Longer segments (6 to 10 s) |
|---|---|
| Faster quality switching | Fewer requests and manifest entries |
| Faster start-up | Better compression efficiency |
| More HTTP requests, more objects to cache | Slower reaction to network changes |
Two to six seconds is a common range for on-demand video.
Common misconception
"Adaptive streaming is done by the server." In HLS and DASH, the server is a plain file server; the client decides which rendition to fetch next. This is exactly why the content can be cached by any standard HTTP CDN.
Deep dive 4: CDN distribution and cache hierarchy
At around 10 Tbps, delivery is the system. A CDN brings bytes close to viewers so they travel fewer network hops, start faster, and so the origin is protected.
The hierarchy
viewers -> edge PoPs (many, small, close to users,
| often inside ISP networks)
| miss
v
regional / mid-tier caches (fewer, bigger)
| miss
v
origin shield (one or a few per region)
| miss
v
object storage (the origin)
A PoP (point of presence) is a cluster of CDN servers in one location. The origin shield is a designated cache layer in front of the origin; all misses funnel through it, so if 500 edges miss the same segment at once, the origin sees roughly one request rather than 500.
Why the hit ratio is high
Video popularity follows a long-tail pattern: a small fraction of videos gets most views. Those popular segments stay hot in edge caches. Long-tail videos (watched a few times a month) mostly miss the edge and are served from regional tiers or the origin.
Useful techniques:
- Long cache lifetimes. Segments never change after publishing, so they can carry
Cache-Control: public, max-age=31536000, immutable. If a video is re-encoded, write new keys rather than overwriting. - Request coalescing. When many viewers miss on the same segment at the same moment, the edge sends one request upstream and serves everyone from that one response.
- Pre-positioning. For a curated catalogue (the Netflix model), you know in advance what will be popular tonight, so you push files to edge servers during off-peak hours. Netflix's Open Connect appliances, placed inside ISP networks, are a well-known public example of this approach.
- Multi-CDN. Large platforms use several CDN providers and steer each viewer to the one performing best in their region, which also gives resilience if one CDN has an outage.
Protecting content
Signed URLs or signed cookies let the CDN check that a viewer is allowed to fetch a segment without calling your API. Paid content adds DRM (digital rights management): segments are encrypted, and the player gets decryption keys from a licence server (Widevine, FairPlay or PlayReady, depending on the device).
View counts
23,000 view events per second on average, far more for a viral video, cannot each update one database row: the row becomes a hot spot with heavy lock contention.
The scalable pattern:
- The player sends a view event to an ingestion service, which writes it to a log such as Kafka, partitioned by
video_id. - Stream aggregators read the log, filter fraud (bots, repeated refreshes from the same session), and sum counts per video in memory for a short window, for example 10 seconds.
- Every window, the aggregator writes one increment per video: "video X +4,213". That converts thousands of writes into one.
- The displayed count is read from cache and may lag by seconds or minutes. That is acceptable, and is why some platforms visibly freeze or lag view counts on new videos while validation runs.
For extremely hot videos, a single counter key can still be a hot spot. You can shard it into N sub-counters (views:vid_91c2:0 to views:vid_91c2:15), increment a random one, and sum them on read.
Exact versus approximate
Views that drive creator payments need exact, auditable counting from the durable event log, recomputed in batch. The number on the page only needs to be roughly right and fresh. Keeping the two paths separate is a classic pattern sometimes called a lambda-style design: a fast approximate path plus a slower accurate one.
Recommendations at a high level
Recommendations are their own large system; in this interview you only need the outline:
- Collect signals: watches, watch time, likes, skips, searches, subscriptions. These come from the same event log as view counts.
- Candidate generation: from millions of videos, cheaply pick a few hundred likely candidates. Sources include "videos similar to ones you finished", subscriptions, trending in your region and language, and embedding-based nearest neighbours. An embedding is a list of numbers representing a user or video, where similar items have nearby vectors.
- Ranking: a heavier machine-learning model scores each candidate for this user (for example predicted watch time) and sorts them.
- Re-ranking and filtering: remove watched or blocked videos, enforce diversity so the feed is not ten videos of the same creator, and apply policy rules.
- Serving: precompute feeds for active users in batch, and refresh the top of the feed in near real time.
Scaling and bottlenecks
| Component | Bottleneck | How to scale |
|---|---|---|
| Playback bandwidth | Tbps of egress | CDN tiers, pre-positioning, multi-CDN, efficient codecs |
| Origin | Cache misses on long-tail content | Origin shield, request coalescing, regional storage replicas |
| Transcoding | CPU/GPU hours, bursts | Chunk-level parallelism, autoscaled spot fleets, priority queues |
| Metadata DB | Reads for the video page | Cache in front, read replicas, shard by video_id |
| View counts | Hot rows | Event log, windowed aggregation, sharded counters |
| Storage cost | Grows forever | Tiering, delete unneeded renditions, better codecs |
Cost considerations
Cost is a first-class requirement here, and interviewers like candidates who raise it unprompted.
- Bandwidth is the largest cost. Better codecs cut it directly: if AV1 delivers the same quality at, say, 30% fewer bits than H.264 (the real saving depends on content and settings), a popular video's delivery bill drops by about the same fraction. That is why platforms spend extra encoding CPU on popular videos: encode once, save on every view.
- Encode lazily for the long tail. Many uploads get very few views. You can produce only a basic ladder in H.264 at upload time, and produce expensive AV1 renditions only after a video crosses a view threshold.
- Storage tiering. Move originals and rarely watched renditions to colder, cheaper storage classes after some time. Cold tiers are cheap to keep but slower or costlier to read, so keep popular renditions hot.
- Keep or drop the original? Keeping the original lets you re-encode later with better codecs. Many platforms keep it, in cold storage.
- CDN hit ratio. As computed above, each percentage point of hit ratio saves around 100 Gbps of origin egress at our scale.
Failure handling
- Upload interrupted: resume from the last confirmed part; abandoned uploads expire.
- Transcode worker crash: lease expires, task is retried; idempotent outputs make retries safe. Poison videos (corrupt files that crash the encoder every time) go to a dead-letter queue after N attempts, and the video is marked
FAILEDwith a helpful message to the creator. - Orchestrator crash: the DAG state lives in a database, so a new orchestrator instance resumes from it.
- CDN PoP outage: DNS or anycast routing sends viewers to the next nearest PoP; multi-CDN steering moves traffic to another provider.
- Origin region outage: replicate popular renditions to a second region; the origin shield fails over.
- Metadata DB outage: the video page can still serve from cache for popular videos; uploads queue up or return a retryable error.
- Viewer on a bad network: ABR drops to 240p; the player retries failed segment requests, possibly against a backup CDN host listed in the manifest.
- Event pipeline lag: view counts freeze briefly but playback is unaffected, because the two paths are decoupled.
Trade-offs and alternatives
| Decision | Option A | Option B | Our choice and why |
|---|---|---|---|
| Upload protocol | Single POST | Multipart resumable | Resumable, because large files and flaky networks |
| Transcode timing | Eager full ladder | Basic ladder now, extra codecs later | Hybrid, to save cost on the long tail |
| Streaming format | HLS only | DASH only | CMAF segments with both manifests, for device coverage |
| Segment length | 2 s | 6 to 10 s | About 4 s, balancing switching speed and request count |
| Delivery | Own CDN | Commercial CDN | Commercial first; own edge appliances only at very large scale |
| View counts | Synchronous DB increment | Event log plus aggregation | Event log, because of hot rows and write volume |
| Metadata store | Sharded SQL | Wide-column NoSQL | Either works; access is key-based. Pick what the team operates well |
| Original file | Delete after transcode | Keep in cold storage | Keep, so you can re-encode with future codecs |
What interviewers probe
"How would you add live streaming?" Live changes the latency budget. The broadcaster sends a stream (for example over RTMP or SRT) to an ingest server, which transcodes in real time and publishes segments as they are produced. The manifest becomes a sliding window that the player refreshes. Low-latency HLS and DASH use partial segments to get latency down to a few seconds. There is no time for chunk-parallel encoding, so you need dedicated real-time encoders per stream.
"A video goes viral in ten minutes. What breaks?" Edge caches warm up quickly because the same segments are requested repeatedly. The risks are the first burst of misses (request coalescing and origin shield absorb it), the metadata cache key for that video (replicate hot keys or use a local in-process cache), and the view counter (sharded counters).
"How do you make playback start in under 2 seconds?" Start at a low rendition, use short initial segments, keep manifests and first segments hot at the edge, preconnect to the CDN, and prefetch the first segment of the video a user is likely to click next (for example the next item in autoplay).
"How do you prevent someone from hotlinking your segments?" Signed URLs or cookies with short expiry, checked by the CDN, plus DRM for premium content.
"Where is the single point of failure?" The orchestrator and metadata DB are the obvious ones. Run them replicated with failover; the DAG state is persisted, so orchestrator instances are stateless.
"How would you handle copyright?" A fingerprinting step in the DAG computes audio and video fingerprints and matches them against a reference database before or shortly after publishing. Matches can block, monetise for the owner, or flag for review.
Interview questions
Q1. Why upload directly to object storage instead of through the API servers?
Video files are huge, and routing them through application servers wastes their bandwidth, memory and connections on byte shuffling. With signed URLs, the API only authenticates and records the upload, while object storage, which is built for large transfers, receives the bytes. It also simplifies scaling, since the API tier stays small and stateless.
Q2. What is a bitrate ladder?
It is the set of renditions you encode for each video, such as 240p at 0.4 Mbps up to 1080p at 5 Mbps. The player moves up and down this ladder depending on network conditions. The ladder can be fixed for all videos or tuned per title based on how compressible the content is.
Q3. How does adaptive bitrate streaming work?
Each rendition is cut into short segments listed in a manifest. The player downloads segments one at a time and, before each request, chooses which rendition to fetch based on measured throughput and how many seconds of video it has buffered. Because the decision is on the client and the files are plain HTTP objects, any CDN can cache them.
Q4. HLS versus DASH?
Both are HTTP-based adaptive streaming formats with segments and manifests. HLS comes from Apple and is required for native playback on Apple devices; DASH is an international standard widely used on Android, browsers and TVs. Using CMAF segments lets you store one copy of the media and publish both an HLS and a DASH manifest pointing to it.
Q5. Why split a video into chunks for transcoding?
Encoding a long video serially on one machine takes a long time. Cutting at keyframes produces independent chunks, so each (chunk, rendition) pair becomes a separate task across many machines, which cuts wall-clock time dramatically. The packaging step stitches the results back together.
Q6. How do you make transcoding tasks safe to retry?
Make each task idempotent: given the same input, it writes the same output to a deterministic storage key, so running it twice is harmless. Combine this with leases, so a task whose worker disappears is re-queued, and a retry limit with a dead-letter queue for inputs that always fail.
Q7. What is an origin shield and why use it?
It is a cache layer between the many CDN edges and your origin storage. All edge misses go through it, so concurrent misses for the same object collapse into one origin request. This protects the origin during spikes and raises the effective cache hit ratio.
Q8. Why not increment the view count in the database for every view?
At tens of thousands of views per second, and much more on one viral video, a single row becomes a write hot spot with lock contention. Writing events to a log and aggregating them in short windows turns thousands of increments into one write per video per window. The trade-off is that the displayed count lags by seconds or minutes.
Q9. How would you reduce bandwidth cost?
Use more efficient codecs (VP9, AV1) for popular videos, tune bitrates per title, keep the CDN hit ratio high with long cache lifetimes and coalescing, and place caches inside ISP networks at very large scale. Avoid spending expensive encodes on long-tail videos that few people watch.
Q10. What segment duration would you choose and why?
Around 2 to 6 seconds for on-demand video. Shorter segments react faster to network changes and start faster, but create more requests and objects. Longer segments compress slightly better and reduce request overhead, but make quality switching sluggish.
Q11. How do you store petabytes of video cheaply but durably?
Use object storage with erasure coding or multi-zone replication for durability. Apply lifecycle rules to move originals and cold renditions to cheaper storage classes, and delete abandoned multipart uploads. Keep only hot renditions in the standard tier.
Q12. What happens if a transcoding worker dies halfway through a task?
Its lease expires because it stops renewing it, and the orchestrator puts the task back on the queue for another worker. Partially written output is overwritten, because the task writes to deterministic keys and only reports success after the output is complete. The video simply becomes ready a little later.
Q13. How is the "video is ready" state decided?
The orchestrator tracks every task in the video's DAG. When all tasks required for playback are complete (for example at least one low rendition and the manifests), it updates the video status in the metadata database and invalidates the cache. Optional tasks like extra codecs can finish later and update the manifests.
Key takeaways
- The system splits into an asynchronous, write-heavy upload path and a read-heavy playback path served almost entirely by the CDN.
- Playback bandwidth (around 10 Tbps in our estimate) dominates the design and the cost.
- Use resumable multipart uploads straight to object storage via signed URLs.
- Model transcoding as a DAG of idempotent tasks, parallelised by chunk and rendition, with leases, retries and a dead-letter queue.
- HLS and DASH serve short segments listed in manifests; the client chooses quality using throughput and buffer level.
- A multi-tier CDN with an origin shield, immutable cache headers and request coalescing keeps the hit ratio high.
- Count views through an event log with windowed aggregation, not direct row updates.
- Cost levers: codec choice, per-title encoding, lazy encoding for the long tail, storage tiering and cache hit ratio.
Next lesson
Continue with Design a file sync and storage service.

