YouTube / Video Streaming Platform
Sharpened prompt. Design a video platform ingesting 500 hours of video per minute, serving 1B hours of watch time per day, where a 4K source becomes playable in every resolution and codec within minutes, playback starts in under 2 seconds anywhere in the world, view counts are accurate enough to pay creators, and the 95% of videos nobody watches do not cost the same as the 5% that everyone does.
Two numbers govern everything: transcoding cost (compute) and egress cost (bandwidth). Every architectural decision here is really a decision about one of those two bills.
1. Problem framing
Functional requirements
- Upload video of arbitrary size and format; transcode into an adaptive ladder.
- Stream with adaptive bitrate (HLS/DASH) and instant seeking.
- Count views for creator analytics and monetisation.
- Serve thumbnails, captions, and multiple audio tracks.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Time to playable | < 5 min for a 10-min video | Parallel chunked transcoding, prioritised ladder |
| Startup latency | P95 < 2s | Edge caching plus a low-bitrate first segment |
| Rebuffer ratio | < 0.5% of watch time | ABR plus a deep CDN hierarchy |
| Durability | Originals never lost | Object storage with erasure coding, cross-region for masters |
| View accuracy | Auditable for monetisation | Deduplicated, fraud-filtered, reconcilable |
Back-of-the-envelope
Ingest: 500 h/min = 30,000 h/day. At ~5 GB/h of source = 150 TB/day of masters
Transcode: 30,000 h/day × ~8 renditions × ~0.5× realtime on GPU
= 120,000 GPU-hours/day -> ~5,000 GPUs running continuously
Storage: masters 150 TB/day + renditions ~2× = ~450 TB/day = 164 PB/year
Watch: 1B hours/day. At an average 2 Mbps = 1B × 3600 × 0.25 MB/s
≈ 900 PB/day of egress. At even $0.005/GB that is ~$4.5B/year.
-> Owning the CDN is not an optimisation, it is existential.
Views: 1B hours / ~10 min avg = ~6B view events/day = 70k/sec
The egress figure is the one to write on the board. It explains Google Global Cache, Netflix Open Connect, and every codec decision in the design.
2. High-level architecture
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Master storage | Object storage with erasure coding, cross-region | Masters are irreplaceable; renditions are regenerable |
| Splitter | GOP/keyframe-aligned segmentation | Chunks must be independently decodable or you cannot parallelise |
| Transcode | FFmpeg with NVENC on GPUs; AV1 on dedicated hardware or CPU | 10–20× faster than CPU for H.264/HEVC; AV1 is still expensive |
| Orchestration | Temporal or Celery with a priority queue | Long, multi-step, resumable workflows with per-chunk retries |
| Packaging | HLS (fMP4) + DASH from one CMAF source | One set of segments, two manifests — halves storage vs separate packaging |
| Quality gate | VMAF sampled per rendition | Catches transcode regressions that PSNR misses |
| CDN | Own PoPs + ISP-embedded caches | At 900 PB/day, third-party CDN pricing is not viable |
| View pipeline | Kafka → Flink → ClickHouse | Dedup and fraud filtering must be stateful and replayable |
4. The toughest parts
4.1 Transcoding: turning a 3-hour serial job into a 5-minute parallel one
Why it's hard. Transcoding is inherently sequential — inter-frame codecs reference previous frames — so a naive pipeline transcodes a 3-hour 4K video in hours per rendition, and you need eight renditions. The creator waits. Multiply by 500 hours uploaded per minute and no fleet size is sufficient.
Solution — split on keyframes, transcode chunks independently, stitch the results.
A GOP (Group of Pictures) starts with an I-frame that references nothing, so a GOP boundary is a safe cut point. Chunks are then embarrassingly parallel.
# 1. Split on existing keyframes without re-encoding. Fast, lossless, I/O bound.
ffmpeg -i master.mp4 -c copy -f segment -segment_time 6 \
-reset_timestamps 1 -segment_format mp4 chunk_%05d.mp4
# 2. Each chunk is an independent job. 180 min / 6 s = 1,800 parallel jobs.
# Force closed GOPs and a fixed keyframe interval so output segments align
# across renditions — ABR switching depends on that alignment.
ffmpeg -i chunk_00042.mp4 \
-c:v h264_nvenc -preset p4 -rc vbr -cq 23 -b:v 4500k -maxrate 6750k -bufsize 9000k \
-g 48 -keyint_min 48 -sc_threshold 0 -force_key_frames "expr:gte(t,n_forced*2)" \
-vf scale_cuda=1920:1080 -c:a aac -b:a 128k \
out_1080p_00042.mp4
# 3. Stitch and package once per rendition, then generate CMAF segments + manifests.
ffmpeg -f concat -safe 0 -i list.txt -c copy 1080p.mp4
Three details that matter. -sc_threshold 0 plus a fixed -g forces deterministic keyframe placement; without it, different renditions get keyframes in different places and the player cannot switch bitrates cleanly mid-stream. Closed GOPs ensure no chunk references frames outside itself. And chunk boundaries must be identical across renditions, which is why you force keyframes at fixed times rather than letting the encoder decide.
Prioritise the ladder. Produce 360p and 720p first and mark the video watchable; generate 1080p, 4K, and AV1 afterwards. Time-to-first-playable drops from "all renditions" to "one rendition," which is what the creator actually experiences.
Run the fleet on spot/preemptible instances — every chunk job is idempotent and cheap to retry, so a 60–90% discount costs only occasional re-runs. State the reasoning: "idempotent, short, retryable units of work are exactly the workload spot instances are for, and transcoding is the canonical example."
4.2 Per-title encoding: not every video deserves the same bitrate
Why it's hard. A fixed bitrate ladder (240p@400k, 480p@1M, 720p@2.5M, 1080p@4.5M…) is wrong in both directions. A screencast of static slides looks perfect at a fraction of that and you are wasting egress on every view. A high-motion sports clip looks blocky at 4.5 Mbps and you are shipping a bad experience. Multiply either error by 900 PB/day.
Solution — choose the ladder per title (ideally per shot) using a perceptual quality target.
def build_ladder(master, target_vmaf=93):
ladder = []
for height in (240, 360, 480, 720, 1080, 1440, 2160):
if height > master.height: break
# Binary search the lowest bitrate that reaches the quality target,
# measured on sampled representative segments rather than the whole file.
lo, hi = 100, 20_000
while hi - lo > 100:
mid = (lo + hi) // 2
score = vmaf(encode_sample(master, height, mid), master)
if score >= target_vmaf: hi = mid
else: lo = mid
ladder.append(Rung(height=height, bitrate=hi))
# Prune rungs that are too close together to be useful for ABR switching.
return [r for i, r in enumerate(ladder)
if i == 0 or r.bitrate > ladder[i-1].bitrate * 1.4]
Netflix reported roughly 20% average bitrate savings from per-title encoding at equal quality. At 900 PB/day, 20% is ~180 PB/day of egress — the largest single cost lever in the entire system. Do the arithmetic out loud; it converts a codec detail into a business argument.
The trade-off to acknowledge: the search costs extra encodes up front. Amortise it by sampling only a few representative segments rather than encoding the whole title repeatedly, and by skipping the analysis for videos predicted to get few views (see 4.6).
Codec selection follows the same logic: AV1 is ~30% more efficient than H.264 but far more expensive to encode, so it is worth it only for videos that will accumulate many views. Encode AV1 lazily, triggered by view velocity, not eagerly for everything.
4.3 Delivering 900 petabytes a day
Why it's hard. No commercial CDN can be bought at this scale at a price that works. Even if it could, the last mile is the problem: bytes must cross ISP peering links, which are congested at peak and expensive.
Solution — a hierarchy that terminates inside the viewer's ISP.
Origin (a few global sites)
└─► Regional PoP (dozens; large SSD/HDD caches, high hit rate)
└─► ISP-embedded cache (thousands; a box inside the ISP's network)
└─► Viewer
Cache hit rates compound: 60% at the ISP box, 30% of the remainder at the PoP,
so origin serves ~10% of what it otherwise would — and the ISP box's traffic
never crosses a paid transit link at all.
Three mechanisms make it work:
- Predictive pre-positioning. Push a new video from a large channel to ISP caches in its expected regions before the audience arrives. Popularity is predictable from the channel's history and the first minutes of view velocity, so pre-warming converts a cache-miss storm into a cache hit.
- Byte-range requests over segments. Players fetch 2–6 second segments by range, so seeking does not require the whole file and partial caching is natural. A segment is an independently cacheable object with a stable URL, which is exactly what a CDN wants.
- Cache the right renditions. 360p and 720p account for the bulk of views on mobile; 4K is a small fraction of views but a large fraction of bytes. Give the low rungs long TTLs at the edge and let 4K fall back to the PoP more often.
The other half of delivery is the player, and it is genuinely part of the system design:
// ABR: choose the next segment's rendition from throughput AND buffer occupancy.
function chooseRendition(ladder, throughputEstimate, bufferSeconds) {
const safe = throughputEstimate * 0.8; // margin for variance
let pick = ladder[0];
for (const r of ladder) if (r.bitrate <= safe) pick = r;
// Buffer-based override: a nearly empty buffer means drop immediately,
// regardless of what the throughput estimate claims.
if (bufferSeconds < 5) pick = ladder[0];
if (bufferSeconds > 30) pick = stepUp(ladder, pick); // plenty of runway, try higher
return pick;
}
Start playback at a low rendition and step up: a fast start beats a high-quality start, because abandonment is driven almost entirely by startup delay. That single heuristic is worth mentioning because it shows you understand the metric that matters (startup time and rebuffer ratio, not peak quality).
4.4 Uploading a 100 GB file from a laptop
Why it's hard. Large uploads fail. A single POST of 100 GB over a home connection will be interrupted, and restarting is unacceptable. The upload also must not traverse your API tier — 150 TB/day through application servers is pure waste.
Solution — resumable, chunked, direct-to-storage uploads with server-side assembly.
POST /upload/session
{"filename":"master.mov","size":107374182400,"sha256":"..."}
→ {"session_id":"u-8f14e","chunk_size":33554432,"upload_urls":"..."}
PUT <presigned-url>/part/47 ; direct to object storage, 32 MB
Content-Range: bytes 1543503872-1577058303/107374182400
→ 200 {"etag":"..."}
GET /upload/session/u-8f14e ; on resume: which parts landed?
→ {"received":[0,1,2,...,46,48], "missing":[47]}
POST /upload/session/u-8f14e/complete
{"parts":[{"n":0,"etag":"..."}, ...]}
→ 201 {"video_id":"v_9a2f"} ; assembly happens server-side
S3 multipart upload provides this natively. The properties that matter: parts upload in parallel (8–16 concurrent connections saturate a home link), each part is independently retryable, the client can resume after days, and the API tier only ever sees small JSON control messages.
Add a client-computed hash per part so corruption is detected at upload rather than discovered during transcode, and set a session TTL with cleanup so abandoned multipart uploads do not silently accumulate storage cost — a genuinely common production bug worth naming.
4.5 Counting views when money depends on it
Why it's hard. View counts drive creator payouts and advertiser billing, so they must be defensible. But "a view" is a product definition (30 seconds? 50% watched?), the same person can reload repeatedly, bots inflate counts deliberately, and 70k events/sec cannot each be a database transaction.
Solution — a stateful stream pipeline with explicit dedup, fraud filtering, and reconciliation.
views
.assignTimestampsAndWatermarks(
WatermarkStrategy.<View>forBoundedOutOfOrderness(Duration.ofMinutes(5)))
.keyBy(v -> v.videoId)
// 1. DEDUP: one countable view per viewer per video per 30-minute window.
// A Bloom filter or HLL per key bounds memory; exactness is not required
// for the public counter, only for the monetised one.
.process(new DedupWithin(Duration.ofMinutes(30)))
// 2. FRAUD: a rules + model pass over behavioural signals.
.filter(v -> !fraud.isSuspect(v))
// 3. AGGREGATE: 1-minute tumbling counts, exactly-once into the sink.
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new CountAndWatchTime())
.sinkTo(clickhouseExactlyOnceSink);
Run two counters with different guarantees, and say why:
| Counter | Guarantee | Latency | Use |
|---|---|---|---|
| Display count | Approximate, HLL-based | Seconds | The number on the page. Nobody can tell 1,203,441 from 1,203,502 |
| Monetised count | Exactly-once, auditable, reconciled | Hours | Creator payouts and advertiser billing |
Fraud signals worth naming: view duration distributions that are unnaturally uniform, IP and ASN concentration, absent or synthetic user-agent and TLS fingerprints, no correlated engagement (no scrubbing, no volume changes, no pauses), and impossible geographic velocity for one account. Flag rather than delete: keep suspect views in a separate ledger so a decision can be revisited during a dispute.
The reconciliation step is what makes it auditable: a nightly batch job recomputes counts from the raw event log in cold storage and compares against the streaming result. Discrepancies beyond a threshold page someone. Streaming gives you speed; batch gives you truth; comparing them gives you trust.
4.6 The long tail: 95% of videos, 5% of views
Why it's hard. View distribution is extreme — a small fraction of videos gets nearly all the traffic, while the vast majority are watched a handful of times ever. Storing every rendition of every video on fast storage forever means paying premium prices for petabytes nobody reads.
Solution — tier by observed popularity, and regenerate rather than store.
def storage_policy(video) -> Policy:
v30 = video.views_last_30d
if v30 > 100_000: # hot: keep everything close
return Policy(renditions="all", tier="ssd", cdn_ttl_days=30, precompute_av1=True)
if v30 > 1_000: # warm
return Policy(renditions="common", tier="hdd", cdn_ttl_days=7, precompute_av1=False)
if v30 > 10: # cool: keep only what is actually requested
return Policy(renditions=["360p", "720p"], tier="hdd", cdn_ttl_days=1,
on_demand=["1080p", "4k"]) # transcode lazily on first request
# cold: master in archive, renditions deleted, regenerated on demand.
return Policy(renditions=["360p"], tier="archive", cdn_ttl_days=0,
on_demand="all", first_play_latency_s=30)
The key insight: renditions are derived data and can be deleted. Only the master is irreplaceable. For a cold video, deleting the 4K rendition and regenerating it on the rare request trades 30 seconds of first-play latency — for a video watched twice a year — against a permanent storage saving across hundreds of millions of videos.
The same logic drives lazy AV1: encode the expensive codec only once view velocity proves it will pay for itself in egress savings.
4.7 Live streaming shares the pipeline but not the constraints
Why it's hard. Live has an end-to-end latency budget of a few seconds, so you cannot split-parallelise (there is nothing to split yet), cannot per-title encode (no complete title), and cannot pre-position (the content does not exist). Yet it must reuse the same delivery infrastructure.
Solution — a parallel ingest and packaging path that converges on the same CDN.
RTMP / SRT ingest → transcode in realtime (one pass, fixed ladder) → LL-HLS / LL-DASH
→ CMAF chunked transfer (partial segments published as they encode)
→ same PoP and ISP cache hierarchy
→ DVR window written to object storage, becomes VOD after the stream ends
The specific mechanisms worth naming: LL-HLS partial segments (publish 200 ms parts rather than waiting for a full 6 s segment), chunked transfer encoding so the CDN forwards bytes as they arrive rather than buffering a whole object, and a fixed ladder because there is no time to analyse content. Latency drops from ~30 s (classic HLS) to 2–5 s.
The trade to state: low-latency live gives up the encoding efficiency and quality optimisation that VOD enjoys, and it costs more per delivered byte. That is why platforms keep the two paths separate and convert live to VOD (with proper per-title encoding) after the fact.
5. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
| Viral video | Origin thundering herd on cold segments | Shield tier + request coalescing; predictive pre-positioning |
| Upload spike (a global event) | Transcode queue depth explodes | Priority by rendition; degrade to 360p+720p only; scale spot fleet |
| Spot capacity reclaimed en masse | Chunk jobs fail in bulk | Idempotent retries; a small on-demand baseline for the priority ladder |
| Bad encoder deploy | Silent quality regression | VMAF sampling gate before publish; canary a small fraction of jobs |
| Fraud campaign | Inflated counts, wrong payouts | Two-counter model; monetised count is reconciled nightly before payout |
| Regional peering congestion | Rebuffer ratio spikes in one ISP | ISP-embedded caches; ABR degrades gracefully; alert on per-ASN rebuffer rate |
6. Cheat sheet
- The numbers: ~900 PB/day egress and ~5,000 GPUs of continuous transcode. Every decision serves one of those bills.
- Transcode: GOP-aligned chunking → parallel jobs on spot → stitch. Fixed
-g,-sc_threshold 0, closed GOPs so renditions align. - Ladder: per-title (VMAF-targeted) encoding, ~20% bitrate saving; lazy AV1 gated on view velocity.
- Delivery: origin → PoP → ISP-embedded cache; predictive pre-positioning; CMAF segments serve both HLS and DASH.
- Player: start low and step up; buffer occupancy overrides throughput estimates.
- Upload: resumable multipart direct to object storage; API tier never touches the bytes.
- Views: two counters — approximate for display, exactly-once and reconciled for money.
- Storage: renditions are derived data; tier by popularity and regenerate the cold tail on demand.
- The one-liner: "Compute is chunked and parallel, delivery is hierarchical and terminates inside the ISP, and everything else is deciding which of the two bills to pay — because at 900 PB/day, a 20% encoding win is worth more than any latency optimisation."