YouTube / Netflix Streaming

Handle video upload, transcoding pipelines, CDN delivery, and adaptive bitrate streaming at scale.

Last generated

Lesson 15 of 18 available15 practice questions

SPACED REPETITION Β· 15 practice questions

Make this lesson stick.

Try 3 questions now. No account needed. Sample answers aren't saved.

"Store the MP4 and serve it" meets a billion hours a day

Here is the design most people sketch in the first minute of "design YouTube": the creator's app POSTs the video file to your API server, the server saves it to object storage, and viewers download that same file from your servers. It works for a demo. Now put real load on it.

This lesson's running assumptions are sized like YouTube. YouTube said in February 2017 that people watch a billion hours a day; pair that with a round 500 hours of uploads every minute. Treat both as assumptions for the exercise, not current facts.

Request Naive plan What breaks first
Upload a 4 GB phone video over a 20 Mbps uplink One POST through your API server 27 minutes on one connection; a drop at 90% starts over. Your API tier carries 300 Gbps of uploads on average
Play it on a phone that gets 3 Mbps Serve the 10 Mbps original The file arrives 3.3Γ— slower than it plays, so it stalls, if the phone can decode that codec at all
41.7 million streams at once, on average Stream from your servers 104 Tbps. A server with a 10 Gbps network card serves 4,000 streams at 2.5 Mbps
A creator with 50 million subscribers premieres a video Caches that forward every miss 25,000 requests for the first segment reach storage within a quarter of a second
Count the views UPDATE videos SET views = views + 1 per play A viral video sends thousands of increments per second to one row

None of these is fixed by a bigger server. Each one is a pressure in the requirements, and each pressure has a standard design move. This lesson teaches six of them:

# Pressure Move The question it answers
1 Multi-GB files over flaky networks Resumable parts, straight to object storage How does the video get in?
2 Uploads no device plays well, and encoding takes CPU-hours Transcode into a bitrate ladder, in parallel, off the request path What do we store?
3 A viewer's bandwidth changes every few seconds Segments, playlists and adaptive bitrate (ABR) in the player How does playback survive a tunnel?
4 Over a hundred terabits per second, viewers everywhere A tiered CDN with immutable segments Where do the bytes come from?
5 Every play still asks your servers one question A cached metadata service and a separate search index What hits the app tier?
6 Billions of view events, one hot video An event log and windowed aggregation How do we count?

Here is the whole design as one map. Single lines are control requests; double lines carry video bytes. Every move below zooms into one lane.

UPLOAD (write path, once per video)
  creator --(1) POST /uploads ---------> Upload API --> Metadata DB (UPLOADING)
  creator ==(2) PUT parts, presigned ==> Object store  raw/{video_id}
  creator --(3) POST /complete --------> Upload API --> "uploaded" event
                                                             |
PROCESS (asynchronous, minutes)                              v
  coordinator --> task queue --> encoder workers ==> Object store  v/{video_id}/v1/...
  coordinator --> writes master playlist last (object store), then READY --> Metadata DB

PLAY (read path, billions of times a day)
  viewer --(4) GET /videos/{id}/playback --> Playback API --> cache --> Metadata DB
  viewer ==(5) GET playlists + segments ===> edge --> regional --> shield --> Object store
  viewer --(6) view + heartbeat events ----> ingest --> log --> aggregator --> counters
                                             ingest --> position store (user_id, video_id)

This lesson assumes the building blocks from Core Building Blocks, Databases & Storage and Messaging & Queues. The shared toolkit for social and media systems (fan-out, hot keys, counters) lives in Design Social & Streaming Systems. Here you apply it to video.

Before any move: which "YouTube" are you designing?

"Design YouTube" and "design Netflix" sound like one question. They are two different systems, and the questions you ask in Scope, the first phase of the interview framework, decide which one you build.

Ask Why the answer changes the design
Do users upload, or do we play a licensed catalog? A catalog has no public upload path. Its encoder can spend hours per title because millions will watch each encode
On-demand only, or live too? Live forbids encoding the future in parallel and changes the cache rules. Interviewers usually save it for a follow-up
How big and how long can uploads be? Sets part size, upload duration and encode time per video
Which devices and regions? Sets codecs (H.264 plays nearly everywhere; HEVC, VP9 and AV1 only where decoders exist), DRM and CDN footprint
How soon must a new upload be watchable? "Within a minute" forces publishing a small, fast subset of rungs first
Does money depend on the view count? A display count can be approximate. Payouts need an exact, auditable recount

The biggest fork is the first row:

YouTube-style (anyone uploads) Netflix-style (licensed catalog)
New content Millions of videos a day Titles arrive from studios on a schedule
Popularity A long tail: most videos are rarely watched A far smaller catalog; demand can be forecast
Encoding budget Cheap and fast first; spend more only on videos that earn views Spend heavily per title; the encode is reused by millions
Getting bytes near viewers Pull-through caching in tiers Push popular titles close to viewers before they ask

Netflix's Open Connect program is the push model taken to the limit: Netflix appliances sit inside ISP networks and receive nightly content fills. That only works because Netflix knows its catalog and can predict what people will watch tomorrow. Nobody can predict which of today's millions of uploads will go viral.

The numbers that drive the design

Say your assumptions out loud, then compute. Every figure below comes from these inputs: 500 hours uploaded per minute; 10 Mbps average source bitrate; 10-minute average video; the five-rung ladder from Move 2 (9.2 Mbps in total); a billion hours watched per day at an average delivered 2.5 Mbps; 10 minutes per play; a 2Γ— daily peak.

Quantity Arithmetic Result
Video uploaded 500 h/min Γ— 1,440 min 720,000 hours per day
Raw bytes 720,000 h Γ— 3,600 s Γ— 10 Mbps Γ· 8 3.24 PB per day
Upload bandwidth 3.24 PB Γ— 8 Γ· 86,400 s 300 Gbps average
Encoded ladder 9.2 Mbps Γ· 10 Mbps = 0.92 Γ— the source 2.98 PB per day
Stored per year (3.24 + 2.98) PB Γ— 365 about 2.3 EB, before redundancy
New videos 720,000 h Γ· 10 min 4.32 million per day, 50 per second
Concurrent streams 1 billion h Γ· 24 h 41.7 million average, 83 million at peak
Egress 41.7 million Γ— 2.5 Mbps 104 Tbps average, 208 Tbps at peak
Plays 1 billion h Γ· 10 min 6 billion per day: 69,000/s average, 139,000/s at peak

Predict first: which single number forces the biggest architectural decision?

Check your answer

Egress. You ingest 3.24 PB a day but send out 1 billion h Γ— 3,600 s Γ— 2.5 Mbps Γ· 8 = 1,125 PB a day, about 350 times as many bytes. At peak that is 208 Tbps. Even servers that each push 100 Gbps flat out would number about 2,100, and they would need to sit near viewers on every continent. That one number justifies a CDN before you draw a single box.

Storage is the runner-up: about 2.3 EB a year means object storage with cheaper tiers for the long tail, not disks you manage by hand.

πŸ’‘ Say it in the interview: "Reads dominate by bytes: about 350 bytes out for every byte uploaded. So I'll design two paths, an upload path that runs once per video and a playback path that runs billions of times, and make sure my application servers never touch video bytes on the playback path."

Move 1: Upload straight to storage, in resumable parts

Pressure: multi-gigabyte files arrive over phone networks that drop, and you don't want those bytes in your app tier.

The naive upload POSTs the whole file to your API server. A 4 GB file at 20 Mbps is 4 Γ— 8,000 Mb Γ· 20 Mbps = 1,600 s, about 27 minutes on one connection. A drop at minute 25 throws all of it away. Meanwhile your API servers spend their lives shovelling 300 Gbps of bytes they never look at. There is also a hard ceiling: a single S3 PUT accepts at most 5 GB.

Predict first: your API tier autoscales behind a load balancer. Why not keep streaming uploads through it and just add servers?

Check your answer

Adding servers fixes capacity, not the failure mode. Every upload is still one long all-or-nothing request, so a dropped connection still restarts from zero. You would also pay for a fleet whose only job is to copy bytes from one socket to another, and each server holds a connection open for half an hour per upload, which makes deploys and scale-in painful.

Split the control plane from the data plane

Your API decides what is allowed: who may upload, how large a file, and which object key it lands on. Object storage moves the bytes. The bridge between them is the presigned URL, a URL signed with your credentials that lets whoever holds it perform one operation on one object until it expires.

Object stores provide the resumable part too. With S3 multipart upload you start an upload, send numbered parts in any order and in parallel, then ask S3 to assemble them. A part is 5 MiB to 5 GiB (the last part may be smaller), an upload has at most 10,000 parts, and the assembled object may be up to about 50 TB. The tus protocol and Google Cloud Storage's resumable uploads solve the same problem with byte offsets instead of numbered parts. Any of them is a fine answer.

Client                      Upload API                  Object store (S3)
  | POST /uploads {size}          |                              |
  |------------------------------>| CreateMultipartUpload        |
  |                               |----------------------------->|
  |<-- upload_id, 16 MiB x 120 ---|                              |
  | POST /uploads/{id}/urls 1-20  |                              |
  |------------------------------>| signs 20 part URLs (1 hour)  |
  |<-- 20 URLs -------------------|                              |
  | PUT parts 1..20, four at a time, each part retried alone     |
  |------------------------------------------------------------->|
  |        ... app killed during part 17, reopened later ...     |
  | GET /uploads/{id}             | ListParts                    |
  |------------------------------>|----------------------------->|
  |<-- "parts 1-16 are stored" ---|                              |
  | fresh URLs, then PUT parts 17..120 ------------------------->|
  | POST /uploads/{id}/complete   | CompleteMultipartUpload      |
  |------------------------------>|----------------------------->|
  |<-- 202, status PROCESSING ----| publishes "uploaded" event   |

Four details in that diagram carry the design:

  • The server picks the object key (raw/{video_id}), never the client's filename. Each URL signs exactly one part of one key.
  • Record parts as they complete; ask storage when you lose track. The upload session (your ID, S3's upload ID, key, owner, expiry) lives in your database, and each part's number and ETag are recorded as the part completes. A client that lost its local progress asks your API, which calls ListParts (1,000 parts per page, so paginate) to find what is missing. AWS advises using that listing only for verification and building the CompleteMultipartUpload request from your own list of part numbers and ETags.
  • URLs are handed out in batches and can be re-signed. Presigned URLs are bearer tokens, so keep them short-lived. A URL signed with temporary credentials also dies when those credentials do (details).
  • Completion is the trigger. Nothing downstream starts until CompleteMultipartUpload succeeds, then one "uploaded" event starts the pipeline. Make /complete idempotent: a retried call finds the session already completed and returns the same answer.

The part size is arithmetic, not taste. Too small and big files exceed 10,000 parts; too large and a retry resends a lot:

import math

MiB, GiB = 2**20, 2**30
MAX_PARTS = 10_000                    # S3 multipart limit
MIN_PART, MAX_PART = 5 * MiB, 5 * GiB

def plan_parts(size_bytes, preferred=16 * MiB):
    """Pick a part size that respects S3's limits, then count the parts."""
    part = max(preferred, math.ceil(size_bytes / MAX_PARTS))
    part = math.ceil(part / MiB) * MiB    # round up to whole MiB
    if part > MAX_PART:
        raise ValueError("file too large for one multipart upload")
    return part // MiB, math.ceil(size_bytes / part)

for gb in (2, 30, 200):
    part_mib, parts = plan_parts(gb * 10**9)
    print(f"{gb:>4} GB -> {part_mib} MiB parts x {parts}")

It prints 16 MiB parts x 120, 16 MiB parts x 1789 and 20 MiB parts x 9537. A fixed 8 MiB part would cap an upload at 8 MiB Γ— 10,000 β‰ˆ 83.9 GB.

Failure modes

  • A part fails. Retry that part only. The other parts are already stored.
  • The user gives up. An unfinished multipart upload keeps its parts, and you keep paying for them, until someone aborts it. Add a lifecycle rule with AbortIncompleteMultipartUpload so storage cleans up by itself.
  • The file is not what it claims. The Content-Type header is a hint. The pipeline's first stage probes the file (container, codecs, duration) and runs policy, malware and copyright checks before anything is published.

Your turn: a drone-survey app uploads 30 GB files over a 25 Mbps uplink, in 16 MiB parts. The first version of the app asks for every part URL at the start, each valid for one hour. What goes wrong, and what do you change?

Check your answer

The upload takes 30 Γ— 8,000 Mb Γ· 25 Mbps = 9,600 s, about 2 hours 40 minutes, in 1,789 parts. Every URL for a part not sent in the first hour has expired by the time the client reaches it, so the upload fails at about 37% even if the connection never drops.

Hand out URLs in batches as the client needs them (or on demand per part), and let a resumed client call GET /uploads/{id} to learn which parts S3 already has and get fresh URLs for the rest. Don't fix it by signing URLs for a week: that widens the window in which a leaked URL can write to your bucket.

πŸ’‘ Say it: "Uploads go straight to object storage as presigned, resumable parts. My API only authorizes and tracks the session, which keeps 300 Gbps of upload traffic out of the app tier." Likely follow-up: "The user closes the app halfway. What happens?" Answer with ListParts, re-signed URLs and the abort lifecycle rule.

Move 2: Transcode once into a ladder, off the request path

Pressure: the uploaded file is the wrong thing to play. It might be ProRes from a camera, HEVC from a phone, a variable frame rate, 60 Mbps, or rotated. Most viewers can't decode it, can't download it fast enough, or both. Converting it takes CPU-minutes to CPU-hours.

The pipeline turns one source into a bitrate ladder: the same video at several resolutions and bitrates (each one a rendition, or rung), cut into short segments, plus playlists that tell players what exists. It also makes thumbnails and audio tracks. This lesson's ladder is H.264 with AAC audio; the rates are target averages including audio:

Rung Resolution Average bitrate
240p 426Γ—240 0.3 Mbps
360p 640Γ—360 0.7 Mbps
480p 854Γ—480 1.2 Mbps
720p 1280Γ—720 2.5 Mbps
1080p 1920Γ—1080 4.5 Mbps

Predict first: five renditions of a 10 Mbps source. Does storing them cost about five times the source?

Check your answer

No. Add the rungs: 0.3 + 0.7 + 1.2 + 2.5 + 4.5 = 9.2 Mbps, which is 0.92 Γ— the source. Storage scales with bitrate Γ— duration, and the low rungs are tiny. Multiplying the source size by the number of rungs gives 16.2 PB a day instead of 2.98 PB, more than five times too much. That is the classic slip in this estimate.

The rule that makes switching possible

A player switches rungs only at segment boundaries, and it can only start decoding a segment at a keyframe (an IDR frame, which needs no earlier frame to decode). So every rung must place its keyframes, and therefore its segment boundaries, at the same timestamps. Apple's HLS authoring spec recommends a keyframe every 2 seconds, 6-second target durations, and segment boundaries aligned across variants.

Encoders don't do this by default. They place keyframes at scene cuts and at a long maximum interval. With FFmpeg's default keyframe settings, on a 30-second test clip, the "6-second" HLS segments came out 8.33 s long, because x264's default maximum keyframe interval is 250 frames (8.33 s at 30 fps). Forcing keyframes fixes it:

ffmpeg -i source.mp4 \
  -vf scale=-2:720 -c:v libx264 -b:v 2400k -maxrate 3600k -bufsize 4800k \
  -force_key_frames "expr:gte(t,n_forced*2)" -sc_threshold 0 \
  -c:a aac -b:a 96k \
  -f hls -hls_time 6 -hls_playlist_type vod \
  -hls_segment_filename 720p/seg_%05d.ts 720p/index.m3u8

-force_key_frames puts a keyframe every 2 seconds regardless of frame rate, -sc_threshold 0 stops extra keyframes at scene cuts, and scale=-2:720 keeps the aspect ratio instead of stretching. Run once per rung with that rung's rates and height, every segment of the 30 fps test clip is 6.000 s (at 29.97 fps it would be about 6.006 s, and a video's last segment is usually shorter), and the keyframes of all five renditions land on the same timestamps. The 720p segments measured about 2.58 Mbps including audio and container overhead, close to the ladder's 2.5. libx264 encodes High profile unless told otherwise, which is what Apple's spec prefers for H.264.

Parallel, asynchronous, and idempotent

Encoding is slow, so it never happens inside the upload request. Assume one machine encodes the full ladder at 2Γ— real time. Then:

  • Fleet size: 30,000 hours of video arrive every hour, so about 15,000 machines are busy on average. That is an assumption to replace with a measurement; the shape of the answer (a large, elastic fleet fed by a queue) is what matters.
  • Latency for one long video: a 2-hour film takes 1 hour on one machine. Split it at keyframes into 60 two-minute chunks, and 60 machines finish in about a minute each, plus time to stitch and package.

Chunking is how you trade machines for latency. The cost: each chunk's encoder can't see the rest of the film, so bit allocation is a little less efficient, and you need a stitching step.

"uploaded" --> coordinator: probe the source, plan the work
                   |  10 chunks x 5 rungs = 50 tasks
                   v
             task queue (at-least-once delivery)
                   |
       +-----------+-----------+
       v           v           v
    worker      worker  ...  worker    each encodes one chunk at one rung and writes
       |           |           |       work/{video_id}/v1/{rung}/chunk_07.ts
       +--- "done" events -----+
                   v
             coordinator: all chunks of a rung done   -> package that rung
                          all rungs packaged          -> write master.m3u8 LAST
                          then status = READY, only if it is still PROCESSING

Queues deliver at least once (the details are in Messaging & Queues), so every step must survive running twice:

  • Deterministic output keys. A re-run writes the same key with equivalent bytes, so a retry overwrites instead of duplicating.
  • Completion is a set, not a counter. The coordinator records "chunk 3 of 720p is done". Recording it twice changes nothing.
  • Progress survives the coordinator. Keep that set in a database or a workflow engine, not in process memory. A coordinator that restarts must resume where it was; one that forgets waits forever for events it already consumed.
  • One writer flips READY, once. Only the coordinator runs UPDATE videos SET status = 'READY' WHERE video_id = ? AND status = 'PROCESSING', and only after the master playlist exists. If each worker ran that statement when it finished, the fastest rung would publish a video whose other rungs are missing.
  • The lease must outlast the task. In SQS, a received message stays hidden for the visibility timeout, 30 seconds by default and at most 12 hours. An encode that runs longer than the timeout reappears and a second worker starts the same task. Set a short timeout (a minute or a few) and have the worker extend it with ChangeMessageVisibility while it works (a heartbeat), so a dead worker's task comes back quickly.
  • Poison files go to a dead-letter queue. A corrupt upload fails on every attempt. After a few receives (maxReceiveCount), the message moves to a DLQ, the coordinator marks the video FAILED, and the uploader hears why.

Why not Kafka as the task queue? In a classic Kafka consumer group, each partition goes to one consumer and progress is one offset per partition. A 20-minute encode at the head of a partition holds up every task behind it, and parallelism can't exceed the partition count. Kafka is a fine home for the "uploaded" and "done" events; a work queue with per-message leases fits long-running tasks better. Newer Kafka releases add share groups (KIP-932, "Queues for Kafka"), which lease individual records with a time-limited acquisition lock; check your version's support before relying on them.

A task's life, with a crash in the middle

A 20-minute video, 10 chunks Γ— 5 rungs = 50 tasks. At the 2Γ— real-time assumption, the whole ladder for one 2-minute chunk takes about 60 s; say 6 s at 240p, 8 s at 360p, 10 s at 480p, 14 s at 720p and 22 s at 1080p. Leases are 60 s, and workers extend them every 20 s while they encode. Times are minutes:seconds.

Time Event Coordinator state
0:00 Upload completes; the coordinator enqueues 50 tasks PROCESSING, 0 of 50 done
0:02 50 workers each receive one task 0 of 50
0:08–0:12 The 240p, 360p and 480p tasks finish 30 of 50
0:15 The worker on chunk 7 / 1080p crashes before its first lease extension 30 of 50
0:16 The 720p tasks finish; the "done" event for chunk 3 / 720p is delivered twice 40 of 50 (the duplicate is already in the set)
0:24 The other nine 1080p tasks finish 49 of 50
1:02 Chunk 7 / 1080p's lease runs out (60 s after 0:02); another worker receives it 49 of 50
1:24 The re-run writes the same output key and reports "done" 50 of 50
1:30 Coordinator packages 1080p, writes master.m3u8, flips READY (1 row updated) READY

Without the crash the video would have been READY about 30 seconds after the upload, against 10 minutes on one machine. The crash cost one lease plus one re-encode of one chunk, about a minute, not the whole video; the lease length is the price of a crash, which is why leases are short and extended by heartbeats rather than set to hours. The duplicate event cost nothing. And no player could see the video before every rung it lists existed. The creator's app learns about READY by polling the video's status every few seconds or from a push notification; a job measured in minutes doesn't need a persistent connection.

Keep the source?

The source is 3.24 of the 6.22 PB you store each day, 52% of the bytes. Deleting it after encoding looks tempting. But the source is the only way to re-encode for a new codec, a better ladder or a fixed bug; without it you can never do better than your top rung. The usual answer: keep the source (or a high-quality mezzanine) in a cheap archive storage class, and delete only what is truly disposable.

Spend encoding effort where the views are

Better codecs cost more to encode. YouTube has said that VP9 takes about five times the compute of H.264 to encode, which is one reason it built a custom transcoding chip. With a long tail, most uploads are watched a handful of times, so a common strategy is: H.264 for everything, quickly; the expensive codecs (VP9, AV1) only for videos that start getting views, where the bandwidth savings repay the compute. A Netflix-style catalog turns that around and spends heavily on every title, because every title is watched by millions.

Your turn: requirement change. Creators must be able to watch their 20-minute upload within 2 minutes of it finishing, but full quality can follow later. What changes in the pipeline, and what do you give up?

Check your answer

Split the work into a fast path and a backfill: a small, fast subset of rungs first. Encode 360p and 720p first on a high-priority queue, in short chunks (say 1 minute, so 20 chunks Γ— 2 rungs = 40 small tasks run at once), and publish as soon as those two rungs are packaged. The master playlist lists only the rungs that exist. Then 240p, 480p and 1080p go through a lower-priority queue and the coordinator adds them to the playlist when they're done.

What you give up: early viewers get at most 720p. Smaller chunks add per-task overhead and lower compression efficiency. And the master playlist is no longer written once. It changes when rungs arrive, so it can no longer be cached forever. Serve it under a new versioned URL or with a short TTL (Move 4).

πŸ’‘ Say it: "Transcoding is an asynchronous, idempotent workflow: chunk Γ— rung tasks on an at-least-once queue, deterministic output keys, and a coordinator that writes the master playlist last and flips READY once." Likely follow-up: "A worker dies halfway through a task. Walk me through it."

Move 3: Let the player adapt: segments, playlists and ABR

Pressure: a viewer's throughput swings from 6 Mbps to 1.5 Mbps when the train enters a tunnel. One file at one bitrate either stalls there or wastes quality everywhere else.

Two protocols, one idea

Both HLS (Apple, .m3u8 playlists) and MPEG-DASH (an ISO standard, .mpd manifests) do the same thing: the video is a list of short segments per rung, fetched with plain HTTP GETs. Plain GETs of immutable files are exactly what CDNs cache best, and that is the whole trick. Packaged as CMAF (fragmented MP4 segments, instead of the MPEG-TS .ts files this lesson's commands write), one set of segment files can serve both an HLS playlist and a DASH manifest, so you don't store the video twice.

The player first fetches the master playlist (Apple now calls it the multivariant playlist), which lists the rungs. These values were measured from the five rungs encoded with the FFmpeg command above:

#EXTM3U
#EXT-X-VERSION:3
#EXT-X-STREAM-INF:BANDWIDTH=772000,AVERAGE-BANDWIDTH=750000,RESOLUTION=640x360,CODECS="avc1.64001e,mp4a.40.2"
360p/index.m3u8
#EXT-X-STREAM-INF:BANDWIDTH=355000,AVERAGE-BANDWIDTH=343000,RESOLUTION=426x240,CODECS="avc1.640015,mp4a.40.2"
240p/index.m3u8
#EXT-X-STREAM-INF:BANDWIDTH=1286000,AVERAGE-BANDWIDTH=1262000,RESOLUTION=854x480,CODECS="avc1.64001f,mp4a.40.2"
480p/index.m3u8
#EXT-X-STREAM-INF:BANDWIDTH=2650000,AVERAGE-BANDWIDTH=2584000,RESOLUTION=1280x720,CODECS="avc1.64001f,mp4a.40.2"
720p/index.m3u8
#EXT-X-STREAM-INF:BANDWIDTH=4690000,AVERAGE-BANDWIDTH=4645000,RESOLUTION=1920x1080,CODECS="avc1.640028,mp4a.40.2"
1080p/index.m3u8

None of these numbers is typed from the ladder table:

  • CODECS comes from each rung's encoded stream. libx264 chose High profile at a level that depends on resolution, so 720p is avc1.64001f (High, level 3.1) and 1080p is avc1.640028 (High, level 4.0). A string that understates the profile can make a device pick a stream it can't decode.
  • BANDWIDTH is the measured peak segment bit rate and AVERAGE-BANDWIDTH the measured average, audio and container overhead included. For VOD, Apple requires each to be within 10% of what the segments actually measure. They sit a little above the ladder's targets, most at the low rungs, where audio and MPEG-TS overhead are a bigger share. On this short, uniform test clip the peak is only a few percent above the average; real content with fast motion peaks much higher.
  • Order matters. Apple's players start on the first variant listed, so 360p goes first to match the start rule below. (Apple's own suggestion is a default near 2 Mbps; this design trades that for a faster first frame on slow links.)

Each rung's media playlist lists its segments in order, each preceded by its duration (#EXTINF:6.000000, then seg_00000.ts), and a VOD playlist ends with #EXT-X-ENDLIST. That tag tells players the list is final, so they load it once and never reload it.

For a 10-minute video that is 100 segments per rung and 500 segment files per video; at 4.32 million uploads a day, over 2 billion new objects a day. If object count hurts (per-request pricing, listing, metadata), store one file per rung and address segments as byte ranges (EXT-X-BYTERANGE in HLS, which needs playlist version 4 or later; an index box in DASH). That makes 5 objects per video instead of 500, at the price of range requests that your CDN must cache well.

Why the decision lives in the player

Only the player knows its own throughput and how many seconds of video it has buffered. Keeping the decision there leaves the server side as a dumb, cacheable file server, which is exactly what lets the CDN do the work.

Predict first: the viewer's link does 5 Mbps. Why not start at 1080p?

Check your answer

The first 6-second segment at 1080p is 4.5 Mbps Γ— 6 s = 27 Mb, which takes 5.4 s to download on a 5 Mbps link, all before the first frame. At 360p it is 4.2 Mb, 0.84 s. An Akamai-data study (Krishnan and Sitaraman, IMC 2012) found viewers start abandoning when startup takes more than 2 seconds, with each extra second adding about 5.8% to the abandonment rate. Start low, or make the first segments shorter, and climb once you've measured the link.

Here is a small player that follows four rules. Real players are more elaborate, but these four are what you should be able to state in an interview:

  1. No measurement yet: start at 360p.
  2. Budget = 80% of the harmonic mean of the last three throughput samples. The harmonic mean is dominated by the slowest sample, so it reacts quickly to a bad one.
  3. Step up at most one rung per segment, and only with at least 12 s of video buffered.
  4. Step down immediately to the highest rung that fits the budget.
LADDER = [("240p", 0.3), ("360p", 0.7), ("480p", 1.2), ("720p", 2.5), ("1080p", 4.5)]  # Mbps
SEG = 6.0          # seconds of video per segment
SAFETY = 0.8       # spend at most 80% of the estimated throughput
LOW_BUFFER = 12.0  # below this many seconds of buffer, never step up
MAX_BUFFER = 30.0  # stop downloading ahead beyond this

def estimate(samples):
    """Harmonic mean of the last 3 throughput samples: one slow sample pulls it down hard."""
    recent = samples[-3:]
    return len(recent) / sum(1 / s for s in recent)

def choose(samples, buffer_s, current):
    if not samples:
        return 1                                   # no data yet: start low (360p)
    budget = SAFETY * estimate(samples)
    fits = max((i for i, (_, mbps) in enumerate(LADDER) if mbps <= budget), default=0)
    if fits > current and buffer_s >= LOW_BUFFER:
        return current + 1                         # step up one rung at a time
    return min(fits, current)                      # step down at once if it doesn't fit

def simulate(network_mbps):
    samples, buffer_s, rung, playing = [], 0.0, 1, False
    for n, bw in enumerate(network_mbps):
        rung = choose(samples, buffer_s, rung)
        name, mbps = LADDER[rung]
        download_s = SEG * mbps / bw
        stall = max(0.0, download_s - buffer_s) if playing else 0.0
        buffer_s = max(0.0, buffer_s - download_s) + SEG if playing else buffer_s + SEG
        playing = True                             # playback starts after the first segment
        buffer_s = min(buffer_s, MAX_BUFFER)       # (idle time when full is not modelled further)
        samples.append(bw)
        print(f"{n:>3} {bw:>5.1f} {name:>6} {download_s:>6.2f} {stall:>5.2f} {buffer_s:>6.1f} {SAFETY*estimate(samples):>6.2f}")

print("seg  net   rung  dl(s) stall buffer budget")
simulate([6, 6, 6, 6, 6, 6, 1.5, 1.5, 6, 6, 6, 6, 6])

The run: a steady 6 Mbps link, a two-segment tunnel at 1.5 Mbps, then 6 Mbps again. "Budget" is the estimate after that segment's download, which feeds the next choice.

Segment Link (Mbps) Rung Download (s) Stall (s) Buffer after (s) Budget (Mbps)
0 6.0 360p 0.70 0 6.0 4.80
1 6.0 360p 0.70 0 11.3 4.80
2 6.0 360p 0.70 0 16.6 4.80
3 6.0 480p 1.20 0 21.4 4.80
4 6.0 720p 2.50 0 24.9 4.80
5 6.0 1080p 4.50 0 26.4 4.80
6 1.5 1080p 18.00 0 14.4 2.40
7 1.5 480p 4.80 0 15.6 1.60
8 6.0 480p 1.20 0 20.4 1.60
9 6.0 480p 1.20 0 25.2 2.40
10 6.0 480p 1.20 0 30.0 4.80
11 6.0 720p 2.50 0 30.0 4.80
12 6.0 1080p 4.50 0 30.0 4.80

Read it like an interviewer would:

  • Startup took 0.7 s because the first segment was 360p. The player waited for 12 s of buffer before its first step up (segment 3), then climbed one rung per segment.
  • The tunnel hit mid-download. Segment 6 was already requested at 1080p when the link fell, so it took 18 s. The 26.4 s buffer absorbed it without a stall. That buffer is why this player fetches up to 30 s ahead instead of one segment at a time.
  • The step down was immediate (segment 7 went straight to 480p), but the climb back was slow: the harmonic mean remembered the 1.5 Mbps samples for three segments, so full quality returned at segment 12. That asymmetry is deliberate. A dropped rung is mildly annoying; a stall is the thing viewers abandon over.

Production players weigh the same three goals (high quality, no stalls, few switches) with smarter rules. A Netflix-coauthored SIGCOMM 2014 paper by Huang et al. found that choosing the rung mainly from buffer level works in steady state, with throughput estimates needed mostly during startup. BOLA, a buffer-based algorithm by Spiteri, Urgaonkar and Sitaraman, ships in the dash.js reference player.

Your turn: run the same player with a harsher tunnel, [6, 6, 6, 6, 6, 6, 0.8, 0.8, 6, 6]. Predict whether it stalls, and where. Then name two changes that would help.

Check your answer

Segment 6 is already requested at 1080p when the link drops to 0.8 Mbps. It takes 27 Mb Γ· 0.8 Mbps = 33.75 s against 26.4 s of buffer: a 7.35 s stall, and the buffer ends at 6.0 s. Segment 7 is worse than it needs to be. The budget still includes two 6 Mbps samples (0.8 Γ— harmonic mean of 6, 6 and 0.8 β‰ˆ 1.52 Mbps), so the player picks 480p, which takes 9 s with only 6 s buffered: another 3 s stall. Only 360p (5.25 s) or 240p would have fit.

Two fixes that help. First, when the buffer is low, trust the newest sample rather than the average: with the budget taken from the newest sample alone whenever less than 12 s is buffered, segment 7's budget is 0.8 Γ— 0.8 = 0.64 Mbps, the player picks 240p (2.25 s) and doesn't stall, and the smooth-tunnel trace above comes out exactly the same. Second, abandon a download that is clearly too slow, then fetch that segment again at a lower rung. Shorter segments also help, because each wrong bet costs less.

Segment length is a trade-off

Short segments (2 s) Long segments (6 s)
Startup and switching Faster: less to download before each decision Slower
Live latency Lower (see the live section) Higher
Requests and objects 3Γ— as many Fewer
Compression No worse here: keyframes come every 2 s either way. Worse only if keyframes must get denser to match Leaves room for longer keyframe intervals

Apple recommends 6 seconds as a general default. Short-video apps and low-latency live streams choose shorter segments on purpose.

Failure modes on the player side

  • The edge serving you dies mid-session. The player retries the segment elsewhere. Apple's authoring spec recommends listing duplicate streams on different hosts in the master playlist for failover, and HLS content steering lets a server move players between delivery pathways; DASH manifests can list several BaseURLs. That is also how multi-CDN setups fail over.
  • The playlist lies. If BANDWIDTH understates real peaks, players pick rungs they can't sustain. Measure the packaged segments; don't copy the encoder's target.
  • Oscillation. Without the "one rung up, with buffer" rule, a player bounces between rungs every segment, which viewers notice more than a steady lower quality.

How you know playback is healthy. The player is the only component that sees what the viewer sees, so it reports it, through the same event pipeline as views (Move 6). Track time to first frame (as a percentile, not an average), rebuffering ratio (stall time Γ· watch time), playback failures, average delivered bitrate and CDN hit ratio, sliced by CDN, region, device and app version. "p95 time to first frame above 2 s in one region" points at a CDN or origin problem long before users complain.

πŸ’‘ Say it: "The server side is a dumb file server of immutable segments; the player picks a rung per segment from its throughput and buffer. It starts low for a fast first frame, climbs one rung at a time, and drops immediately when the buffer is at risk." Likely follow-up: "Why not decide the quality on the server?"

Move 4: Push bytes to the edge

Pressure: 104 Tbps on average, 208 Tbps at peak, viewers on every continent, and a player that makes a new HTTP request every few seconds. Every request that crosses an ocean pays a long round trip, and every request that reaches your storage costs you.

A CDN runs caching servers in hundreds of locations (points of presence, PoPs) close to viewers. It terminates the viewer's TCP and TLS connections nearby, answers from its cache when it can, and forwards only misses. How a viewer reaches a nearby PoP is the CDN's business: DNS-based steering answers with the address of a nearby PoP; anycast announces one address from many PoPs and lets internet routing pick the topologically closest; a steering service can hand the player specific server URLs in the playback response. DNS and CDN basics are in Networking Basics.

Tiers: small and near, big and far

viewers --> 500 edge PoPs --> 20 regional caches --> 1 origin shield --> origin
            small, near        bigger, fewer          one designated     object store,
            viewers                                   cache near origin  every segment
  A hit at any tier answers the request. Only a miss moves one step to the right.

Popular videos stay hot in the small edge caches. The long tail misses at the edge and hits in the bigger regional tier. The origin shield is one designated caching layer that every regional miss goes through, so the origin typically sees about one request per object while the shield holds it. CloudFront, for example, describes Origin Shield as an extra caching layer that consolidates requests "resulting in as few as one request going to your origin". "As few as" is the best case: a shield still evicts cold objects, and it merges requests only while a fetch is in flight.

Cache rules follow mutability

Most video objects never change after they are written, so name them so they never have to: put the encoding version in the path (v/{video_id}/v1/720p/seg_00042.ts). A re-encode writes v2/ paths and the playback API starts handing out the v2 master playlist. Nothing needs purging.

Object Changes after it is written? Cache-Control, for example Why
VOD segment Never; a re-encode writes a new path public, max-age=31536000, immutable Refetching identical bytes only adds origin load
VOD media playlist Never for that version Same Same argument
Master playlist Yes, when rungs are added or the version changes Versioned URL, or max-age=60 The one mutable VOD object
Live media playlist Every segment max-age=1 for 2 s segments Every second of staleness is a second of extra delay
Live segment Never Long: at least the replay window Written once, then read by every viewer

⚠️ TTL is not eviction. A one-year TTL says the copy is valid for a year. A small edge still evicts cold segments within hours to make room for popular ones. TTL answers "may I serve this copy?"; eviction answers "is there room to keep it?".

The premiere stampede

A creator with 50 million subscribers premieres a video. One million viewers press play within 10 seconds, spread over 500 edges, 20 regional caches and one shield. Assume a miss takes 250 ms to fill from the next tier, and every player starts at 360p, so they all want the same first segment.

Each edge sees 1,000,000 Γ· 10 s Γ· 500 = 200 requests per second, so 200 Γ— 0.25 s = 50 requests arrive at each edge while its first fetch is still in flight. What reaches the origin depends on whether each tier collapses concurrent misses for the same object into one upstream request:

Collapsing Requests for segment 0 reaching the origin
Nowhere 500 edges Γ— 50 = 25,000, within a quarter second
At the edges only 500: one per edge, each forwarded by its regional cache
At edges and regional caches 20: one per regional cache
At every tier, with a shield 1

25,000 requests in 0.25 s is 100,000 requests per second, all on one video's key prefix. S3 supports at least 5,500 GETs per second per prefix and scales up gradually, returning 503 Slow Down errors in the meantime. So the first viewers of the biggest premiere of the day get errors or long waits. And it repeats for every new segment and every rung until the caches are warm.

The fixes stack. Collapse requests at every tier. Route all misses through a shield. And for a scheduled premiere, pre-warm: request the first minutes of every rung through the regional tiers before the start time, so the stampede lands on warm caches.

Who may watch, and how to take a video down

  • Access control. Hand out segment and playlist URLs signed with a short-lived token (signed URLs or cookies) that the CDN edge verifies. Configure the CDN to check the token but leave it out of the cache key. Otherwise every viewer's URL is a different cache entry and your hit ratio drops to nearly zero.
  • DRM for premium content: segments are encrypted, and the player gets keys from a license server (Widevine, FairPlay, PlayReady) after the playback API authorizes it.
  • Takedowns. Invalidation is one call, not one per file per edge: CloudFront, for example, accepts a wildcard path such as /v/{video_id}/* as a single invalidation and forwards it to all edge locations within a few seconds. So: flip the video's status (the playback API stops issuing playlists and tokens), invalidate the video's whole path, and delete or block its origin objects so no miss can refill a cache. Players already streaming stop when their buffer runs out, within 30 s for the player in Move 3. Short-lived signed URLs add a second lock: a copied URL stops working when its token expires.

Failure modes

  • An edge PoP dies. Its viewers are steered to neighbours, whose caches are cold for that PoP's long tail. Misses surge into the regional tier and the shield, which is exactly the load those tiers exist to absorb.
  • The shield location fails. CloudFront, for example, retries at a secondary Origin Shield location. In your own design, name a fallback and rate-limit origin traffic during the switch.
  • The origin region fails. Replicate segments to a second region and configure origin failover at the CDN. Popular content keeps playing from caches meanwhile; the long tail is what suffers.
  • Short segment TTLs "to be safe". A 60-second TTL on immutable segments makes every edge refetch every hot segment every minute, for bytes that never change.

Your turn: legal sends a court order: remove a video within 5 minutes. It currently has 40,000 viewers. Its segments use public, unsigned URLs with one-year TTLs, and the playback API caches its responses for 10 minutes. What happens if you only set the status to REMOVED, and what do you change for next time?

Check your answer

Setting the status alone fails twice. The playback API keeps serving its cached response for up to 10 minutes, so new viewers can still start the video. And anyone holding the playlist URL can keep fetching segments from edge caches for up to a year, since public URLs need no permission.

Purging only the playlists wouldn't help either: a VOD player loads its playlist once and never reloads it, so the 40,000 current viewers keep fetching segments. Now: delete the playback cache entry, issue one wildcard invalidation for the video's whole path (/v/{video_id}/*, playlists and segments on every edge), and delete the origin objects so no miss can refill a cache. Current viewers stop within their buffer, about 30 s, well inside the 5 minutes. Next time: sign segment and playlist URLs with short-lived tokens so access depends on the playback API, and make the takedown path run all of this the moment the status changes.

πŸ’‘ Say it: "Segments are immutable and versioned, so they get year-long TTLs and a re-encode never needs an invalidation. The only mutable objects are master and live playlists, and I keep those small and short-lived. Misses go through a shield with request collapsing, so a premiere costs the origin about one request per segment per rung while the shield holds them." Likely follow-up: "A video goes viral in one minute. What does the origin see?"

Move 5: Metadata, search, and the request that still hits your servers

Pressure: the bytes are on the CDN, but every play still starts with one request to you: GET /videos/{id}/playback, which returns the title, the channel and a signed master-playlist URL. At 139,000 plays per second at peak, heavily skewed toward a few hot videos, this is the busiest request that reaches your metadata database. (Position heartbeats, below, are even more numerous, but they go to a separate, write-optimised store.)

The data itself is small. At about 3 KB per video (the row plus five rendition rows), 4.32 million new videos a day is 13 GB a day and under 5 TB a year. The challenge is read rate and skew, not size.

CREATE TABLE videos (
  video_id       TEXT PRIMARY KEY,    -- random ID; also the shard key
  owner_id       TEXT NOT NULL,
  title          TEXT NOT NULL,
  description    TEXT,
  status         TEXT NOT NULL,       -- UPLOADING, PROCESSING, READY, FAILED, REMOVED
  visibility     TEXT NOT NULL,       -- PUBLIC, UNLISTED, PRIVATE
  active_version INTEGER,             -- which encoding the playback API hands out
  duration_s     INTEGER,
  created_at     TIMESTAMP NOT NULL
);

CREATE TABLE renditions (
  video_id  TEXT    NOT NULL REFERENCES videos (video_id),
  version   INTEGER NOT NULL,         -- bumps on every re-encode
  codec     TEXT    NOT NULL,         -- h264, vp9, av1
  height    INTEGER NOT NULL,         -- 240 ... 1080
  avg_kbps  INTEGER NOT NULL,
  PRIMARY KEY (video_id, version, codec, height)
);

The rendition key includes version and codec, so an AV1 ladder and a v2 re-encode can live beside the original.

Decisions, each with its reason

  • A relational store, sharded by video_id. Lookups are by ID, the status transition wants a transaction, and hashing random IDs spreads load. Sharding by owner_id would put a mega-creator's whole catalog on one shard.
  • A cache in front, cache-aside. At a 99% hit rate, the database sees 1% of 139,000 = about 1,400 reads per second at peak. Cache-aside, invalidation and stampede protection are covered in Databases & Storage.
  • Two kinds of staleness. A title can be a minute stale. "May this video play?" cannot: on a takedown or a visibility change, delete the cache entry in the same code path that changes the status, and keep a short TTL as a backstop. The creator's own studio page reads from the primary, so they see their edits immediately.
  • Region failure. The primary database lives in one region. Keep an asynchronous replica in a second region and promote it if the first region fails: minutes of failover and possibly a few seconds of lost writes, which a title edit can tolerate. A standby in another availability zone of the same region does not survive a regional outage. Meanwhile, popular videos keep playing from the cache and CDN. One catch: status and visibility live in the same table, so a takedown committed seconds before the failure can be lost and the video becomes playable again. Re-apply recent takedowns after failover from the takedown service's own durable log, or write status changes synchronously to both regions.
  • The hot key. One viral video's metadata is one cache key on one cache node, which may receive tens of thousands of reads per second. Keep a copy in each API server's memory for a few seconds, or replicate hot keys across nodes. See the hot-key toolkit in Design Social & Streaming Systems.

Resume position and watch history

"Continue watching" needs the viewer's position. The heartbeat rate is bigger than it looks: 41.7 million concurrent streams each saving a position every 10 seconds is 4.2 million writes per second on average. That is more write traffic than uploads and views combined. Save every 30 seconds plus on pause and exit, and it drops to 1.4 million per second.

These writes don't belong in the relational metadata store. Use a horizontally scaled key-value or wide-column store keyed by (user_id, video_id), with last write wins: losing the last 30 seconds of a position is acceptable. Watch history is a different access pattern ("my recent videos, newest first"), so it gets its own table, partitioned by user and ordered by time. Two access patterns, two tables.

Search runs on a separate index (Elasticsearch or OpenSearch, for example) built from the metadata database, never the other way round. Changes flow from the database to the indexer through change events, such as change data capture or an outbox table (the outbox pattern is covered in Advanced Topics & Final Prep), so the index lags by seconds. Two consequences to state out loud:

  • The index can briefly show a video that was just removed. Filter on status and visibility at query time, and let the playback API, which reads the source of truth, have the final say when someone clicks.
  • Ranking signals such as view counts and freshness are refreshed in batches, not on every view. A search index that is reindexed on every play melts.

Your turn: the product team wants a video's view count shown in search results "in real time", updated on every view. What do you say?

Check your answer

Push back on "every view". At 139,000 plays per second at peak, a reindex per view is 139,000 index writes per second concentrated on the hottest documents. Offer a bound instead: counts in search refresh every few minutes from the aggregated counters (Move 6), while the watch page shows a fresher count from the counter cache. Name the cost: the number in search results can be minutes behind the number on the watch page.

Move 6: Count views without melting a row

Pressure: 69,000 plays per second on average, 139,000 at peak. Worse, a viral video gets 50 million views in its first day, and if 30% of them land in the first hour that is 15 million Γ· 3,600 s β‰ˆ 4,200 views per second on one row. If each locked increment holds the row for about a millisecond (commit round trip included), one row tops out near 1,000 increments per second. UPDATE videos SET views = views + 1 per play queues up behind that lock.

First, decide what a "view" is. That is a product rule, not a technical fact: for example, "the player reported at least 30 seconds of playback, or the whole video if it is shorter". The player sends one event per play with a unique play_id.

player -- "view" + play_id --> ingest API --> log, partitioned by video_id
                                                 |
                    +----------------------------+-----------------+
                    v                                              v
     stream aggregator, 10 s windows              raw events archived to object store
     one upsert per video per window                               | nightly batch
                    v                                              v
     counter store --> cache --> watch page       exact recount: dedupe, drop bots
                                                  -> payouts, corrections to counters

Aggregation turns a hot row into a trickle. The aggregator keeps an in-memory count per video and writes one upsert per video per 10-second window. The viral video goes from about 4,200 writes per second to one write every 10 seconds. The long tail barely aggregates (one view per window is one write), but long-tail videos aren't hot rows, and a sharded counter store absorbs their total.

The hot partition. Partitioning the log by video_id sends all of the viral video's events to one partition and one consumer. At 4,200 events per second that is easy for an in-memory counter. If one video ever outgrew a partition, split its key (video_id#0 to video_id#7) and sum the parts on read. More counter techniques are in Design Social & Streaming Systems.

Retries and crashes: where double counts come from

Delivery is at least once at two points. The client retries an event whose acknowledgement it never saw, so the same play_id arrives twice; the aggregator drops a play_id it has already counted recently. And the aggregator can crash after writing counts but before committing its position in the log. On restart it rereads those events and counts them again.

The fix for the second is to store the log offset in the same transaction as the counts. This runs the same eight events (one of them a client retry) through a naive consumer and a transactional one, each crashing mid-flush on its first batch:

import sqlite3
from collections import Counter

# View events as they sit in one log partition: (offset, play_id, video_id).
# The client retried play p3, so it appears twice.
LOG = [(0, "p1", "cat"), (1, "p2", "cat"), (2, "p3", "dog"), (3, "p3", "dog"),
       (4, "p4", "cat"), (5, "p5", "cat"), (6, "p6", "dog"), (7, "p7", "cat")]

class Crash(Exception):
    pass

def fresh_db():
    db = sqlite3.connect(":memory:")
    db.executescript("""
        CREATE TABLE views   (video_id TEXT PRIMARY KEY, n INTEGER NOT NULL);
        CREATE TABLE offsets (part INTEGER PRIMARY KEY, next INTEGER NOT NULL);
        INSERT INTO offsets VALUES (0, 0);""")
    return db

def consume_batch(db, size, atomic, crash=False):
    start = db.execute("SELECT next FROM offsets WHERE part = 0").fetchone()[0]
    batch = LOG[start:start + size]
    seen, counts = set(), Counter()
    for _, play_id, video in batch:
        if play_id not in seen:                  # drop retries within the batch
            seen.add(play_id)
            counts[video] += 1

    def add_counts():
        for video, n in counts.items():
            db.execute("INSERT INTO views VALUES (?, ?) ON CONFLICT(video_id) "
                       "DO UPDATE SET n = n + excluded.n", (video, n))

    def save_offset():
        db.execute("UPDATE offsets SET next = ? WHERE part = 0", (start + len(batch),))

    if atomic:
        with db:                                 # one transaction: both or neither
            add_counts()
            if crash:
                raise Crash()
            save_offset()
    else:
        with db:                                 # counts commit on their own...
            add_counts()
        if crash:
            raise Crash()                        # ...and the process dies here
        with db:
            save_offset()

for atomic in (False, True):
    db = fresh_db()
    try:
        consume_batch(db, 4, atomic, crash=True)  # first attempt dies mid-flush
    except Crash:
        pass                                     # restart: resume from stored offset
    consume_batch(db, 4, atomic)
    consume_batch(db, 4, atomic)
    print("atomic" if atomic else "naive ", dict(db.execute("SELECT * FROM views ORDER BY 1")))

It prints naive {'cat': 7, 'dog': 3} and atomic {'cat': 5, 'dog': 2}. There were seven plays: five of cat, two of dog. The naive consumer committed the first batch's counts, died before saving its offset, and counted that batch again after the restart. The transactional one rolled back the half-finished flush, so the replay counted each event once. Stream processors such as Flink get the same effect by checkpointing offsets together with their state and writing to the database transactionally or idempotently. Either way, the aggregator's set of recently seen play_ids is part of that state; if it is lost in a crash, a few retries slip through, which is acceptable for the display count and caught by the nightly recount.

Two paths: fast and approximate, slow and exact

The streaming count is what the watch page shows, seconds behind reality. It is not what you pay creators on. The raw events also land in object storage, and a nightly batch job recounts them exactly: it deduplicates by play_id across the whole day, removes bot traffic, and corrects the fast counters. Money comes from the slow, auditable path.

Two different questions: "how many people?" and "how many right now?"

  • Unique viewers of a video: a HyperLogLog estimates the number of distinct IDs added to it. Redis's version uses at most 12 KB per key with a 0.81% standard error, however many viewers there are.
  • Concurrent viewers of a live stream: a HyperLogLog can't do this alone, because nothing can be removed from it when a viewer leaves. Count heartbeats instead: players ping every 10 seconds, and "watching now" means "pinged in the last 30 seconds". One HyperLogLog per 10-second bucket, merging the last three buckets, gives that answer approximately.

Your turn: creators will now be paid per view each month, and finance wants every payout auditable. Which parts of the design change, and which stay?

Check your answer

The display path stays as it is: fast, approximate, seconds behind. What changes is the batch path, which becomes the system of record. Keep every raw event immutable in the archive, with the view rule applied at recount time. Deduplicate by play_id across the whole period, filter invalid traffic, and store each monthly result with the job version and input files that produced it, so finance can reproduce any number. The cost is latency: payouts are based on counts that are a day or more behind, and the public number can differ slightly from the paid number.

πŸ’‘ Say it: "Views are events, not row updates. The player emits one event per play with a unique ID; a stream aggregator turns thousands of events per second on a hot video into one write per window; a nightly batch recount over the raw log is the exact number that money uses." Likely follow-up: "The aggregator crashes after writing. Do you double count?"

Boss level: "Now make it live"

The interviewer has seen your on-demand design and asks the classic follow-up. Keep what still holds and name what changes:

On-demand Live
Ingest Resumable upload of a finished file A continuous RTMP or SRT stream from the broadcaster's encoder to an ingest server
Transcoding Chunk-parallel, can take minutes Real time: every rung must encode at least as fast as the video arrives. You can't split the future into chunks, so parallelism is per rung
Segments Written once, all before the video is published Written once, one every few seconds, with sequence numbers
Playlist Written once, ends with #EXT-X-ENDLIST A sliding window rewritten every segment
Cache rules Segments and media playlists immutable; the master playlist versioned or short-TTL Segments immutable; the media playlist lives about 1 s
Viewer count Views Concurrent viewers from heartbeats

Segments are still immutable. A live stream does not mean mutable files: seg_1841.ts is written once and then read by every viewer. Only the playlist changes.

Where live latency comes from

HLS clients should not start with a segment less than three target durations from the end of the playlist. So the player sits at least three segments behind the newest one, before encoding, packaging and download time:

Segment length Buffer behind the newest segment Plus encode, package, delivery
6 s at least 18 s about 20 s or more behind real time
2 s at least 6 s high single digits to about 10 s

To go lower, Low-Latency HLS publishes partial segments so players can fetch video before a whole segment is finished. Interactive, sub-second latency (a video call, an auction with live bidding) usually leaves segment-based HTTP delivery for WebRTC, which trades CDN caching for a different, more expensive fan-out architecture.

Origin load, computed properly

One million viewers, 2-second segments, five rungs, 500 edges, a playlist TTL of 1 second. Each player reloads its playlist about every target duration and fetches one segment per 2 seconds:

  • At the edges: 1,000,000 Γ· 2 = 500,000 playlist requests per second, plus 500,000 segment requests per second. Edges exist to take exactly this.
  • At the origin, if the 500 edges cached and collapsed but fetched straight from the origin (no regional tier, no shield): each edge fetches each rung's playlist once per second and each new segment once, so 500 Γ— (5 + 5 Γ· 2) = 3,750 requests per second.
  • With a shield: 5 + 2.5 = 7.5 requests per second.

Without a shield, the origin's load depends on the number of edge caches times the rungs; with a shield, on the rungs alone. Neither depends on the number of viewers. That is why a live event with a million viewers is a CDN problem, not an origin problem, as long as nothing in the path turns caching off.

Live failure modes

  • The ingest server dies. Broadcasters can send to a primary and a backup ingest endpoint, and the packager switches streams at a segment boundary.
  • The transcoder can't keep up. Falling behind real time is permanent lag, so shed load: drop the top rung rather than delay every rung.
  • The playlist TTL is set to zero "to be safe". Caches can no longer reuse a playlist between reloads, so origin load grows with the number of viewers again, limited only by whatever collapsing happens while a fetch is in flight. A 1-second TTL costs at most one second of staleness.

Final round: no label on the problem

Real prompts don't say which move they want. Read the requirements, find the pressure, then pick the move. Decide before you open each answer.

Challenge 1: lecture videos for a university

30,000 students. 400 lecturers each upload about 2 hours of recordings a week, and a lecture must be watchable by 8 am the next morning. Only enrolled students may watch. In exam week, assume a fifth of the students watch at once. Which moves do you keep, which do you drop, and what does playback need?

Check your answer

Keep the resumable direct-to-storage upload (2-hour recordings are large) and the ladder with aligned keyframes (students watch on phones and laptops). Drop the chunk-parallel machinery: a 2-hour lecture at 2Γ— real time takes 1 hour on one machine, against a deadline many hours away, so a managed transcoding service or a handful of workers is plenty.

Peak playback is 6,000 streams Γ— 2.5 Mbps = 15 Gbps. That is small for a commercial CDN and still awkward for your own servers, so use a CDN with signed URLs or cookies issued only to enrolled students. View counting can be a plain database increment; there is no hot row at this scale. The interview point is saying what you don't need and why.

Challenge 2: a swipe-based short-video feed

Clips are 15–60 seconds. Users swipe to the next clip after a few seconds, and the product lives or dies on the next clip starting instantly. Which part of the design changes most?

Check your answer

Startup dominates, so the player changes most. While one clip plays, prefetch the first segment of the next two or three clips at a low rung (360p, say), and use short segments (1–2 s) so that first fetch is small. At 360p, one 2-second segment is 0.7 Mbps Γ— 2 s = 1.4 Mb, about 0.18 MB. A three-clip window holds about 0.5 MB, and each swipe adds one new first segment, much of it wasted when users skip. That is the trade: bandwidth for instant start.

Short clips also mean many tiny objects (a 30-second clip at 2-second segments and 4 rungs is 60 files), so storing one file per rung with byte-range segments is attractive. Encode fast-first: a clip that doesn't start getting views within minutes may never get them. And redefine the view: "30 seconds of playback" makes no sense for a 15-second loop.

Challenge 3: live sports with in-play betting

5 million concurrent viewers for a match, and the betting feature needs the video within about 5 seconds of real time. What breaks in the standard live design, and what do you change?

Check your answer

Latency breaks first. Even 2-second segments put the player at least 6 seconds behind the newest segment before encoding and delivery, so standard HLS or DASH can't meet 5 seconds. Use low-latency HLS or low-latency DASH with partial segments, or WebRTC-based delivery if the budget is well under a second, and accept more requests, tighter tuning and less caching efficiency.

Capacity breaks second: assuming about 5 Mbps for HD sports, 5 million Γ— 5 Mbps = 25 Tbps for one event. Book it ahead with more than one CDN (a kick-off is not the moment to autoscale), steer players between CDNs, and put a shield in front of the packager so origin load depends only on the rungs, not on the number of edges or viewers. Count concurrent viewers from heartbeats.

Cheat sheet: pressure β†’ move β†’ cost

When the requirements say… Reach for What you give up
Multi-GB uploads over mobile networks Presigned multipart or tus/resumable upload straight to object storage; ListParts to resume A session to track; abandoned parts to clean up with a lifecycle rule
Many devices and network speeds A bitrate ladder with keyframes aligned across rungs CPU per upload; storage = sum of rung bitrates Γ— duration (0.92Γ— the source here)
"Watchable within minutes" Chunk Γ— rung tasks in parallel; publish a small, fast subset of rungs first Per-task overhead, a mutable master playlist, lower first quality
At-least-once task queue Deterministic output keys, completion as a durable set, heartbeated leases, a DLQ Duplicate work after crashes, bounded by the lease
Bandwidth that changes mid-video Short segments and player-side ABR (start low, up one rung, down at once) Quality ramps up over the first seconds; more requests with shorter segments
Over 100 Tbps, global viewers A tiered CDN; versioned, immutable segments with year-long TTLs Cold tiers for the long tail; paying for CDN egress
A premiere or viral spike Request collapsing at every tier, an origin shield, pre-warming An extra cache layer to pay for and operate
Hot metadata reads Cache-aside, sharding by video_id, per-server caching of hot keys Titles up to a minute stale; explicit invalidation for takedowns
Hot counters Event log plus windowed aggregation; offsets committed with counts Counts seconds behind; an exact batch recount for money
Live Real-time transcoding, sliding playlists with 1 s TTLs, heartbeat viewer counts Latency of about three segments unless you pay for low-latency delivery

Before moving on, run this design aloud from a blank page, in the phases of the interview framework. Scope: the questions that decide YouTube versus Netflix, and the two or three numbers that decide something (egress, plays per second, a hot video's views per second). Sketch: the three lanes of the map, tracing one upload and one play end to end. Deep dive: one move, say the pipeline or the CDN, with its failure modes. Wrap-up: the biggest trade-off and what changes at 10Γ—, or "now make it live". If you can do that without redrawing everything, you're ready.

Next: Design Social & Streaming Systems for the shared toolkit, Twitter/Instagram News Feed and Chat Application (WhatsApp) for the sibling designs, and Mock Interviews & Communication to practise this one under a timer.