YouTube is the canonical example of a user-generated video platform: anyone can upload, uploads arrive constantly, most videos are watched rarely, and a few are watched by hundreds of millions of people. Google has said that more than 500 hours of video are uploaded every minute. Designing for that is a different problem from designing a subscription streaming catalogue, where a studio delivers a few thousand titles that can be prepared at leisure.

This article designs such a platform from first principles and uses the parts of YouTube's real architecture that are public, such as Vitess for MySQL sharding, Google Global Cache in ISP networks and the Argos video coding unit, to check the reasoning. Adaptive bitrate playback and bitrate ladders are covered in the video streaming architecture article and are only summarised here; the focus is on what user-generated content changes: ingest, processing at write time, metadata at scale, counting and the long tail.

Advertisement

Requirements and what they imply

Functional requirements: creators upload videos of any length and format from any network; videos become watchable within minutes; viewers browse, search, watch, like and comment; the platform shows view counts, recommends videos and runs ads; rights holders and policy systems can block or claim content. Non-functional: playback starts in about a second and rarely stalls, uploads survive flaky mobile connections, no accepted upload is ever lost, and cost per hour watched keeps falling.

Three properties drive the design. Writes are heavy but rare per object: an upload is gigabytes and triggers minutes of compute, then is mostly read. Reads are extremely skewed: a small share of videos gets most views, while the long tail must still be stored and served. And processing cost is paid for every upload, including the majority that will never be popular, so the pipeline must decide how much work a video deserves.

Upload, process, store, serve: the four planes of a user-generated video platformCreator appresumable uploadUpload frontendssession + chunksOriginal storeblob, replicatedWorkflow engineper-video DAGupload doneTranscode: split, encode, joinThumbnails, captions, ASRRights match and policy checksMetadata DBsharded MySQLRendition storesegments per formatpublishstateViewerplayerWatch APImetadata + URLsEdge cachesin ISPs and PoPsCounters, logs, MLviews, recs, ads1 watch2 segmentsmisseventsWrites are rare and heavy; reads are constant and cheap only if the edge absorbs them
Four planes with different scaling rules: upload frontends scale with concurrent uploads, processing with upload hours, metadata with object count and request rate, delivery with hours watched.

Upload: resumable, chunked and verified

A 3 GB upload over a phone connection will be interrupted. The upload protocol must let the client resume from the last byte the server confirmed, rather than restart. Google APIs, including the YouTube Data API, use a resumable protocol: the client opens a session with the metadata, receives a session URI, sends the file in chunks, and after an interruption asks the server which range it holds.

# Resumable upload, the protocol shape used by Google APIs including the YouTube Data API.
# 1. Start a session: metadata only, server returns a session URI in the Location header.
POST /upload/youtube/v3/videos?uploadType=resumable&part=snippet,status
X-Upload-Content-Length: 3221225472
X-Upload-Content-Type: video/mp4
{ "snippet": {"title": "Launch talk"}, "status": {"privacyStatus": "private"} }
-> 200 OK, Location: <session URI>

# 2. Send bytes in chunks (multiples of 256 KiB except the last).
PUT <session URI>
Content-Range: bytes 0-8388607/3221225472
-> 308 Resume Incomplete, Range: bytes=0-8388607

# 3. After a network failure, ask the server what it has, then continue from there.
PUT <session URI>
Content-Range: bytes */3221225472
-> 308 Resume Incomplete, Range: bytes=0-1157627903

Server side, upload frontends are stateless; the session URI encodes or points to the session state, so any frontend can accept the next chunk. Chunks are written to a durable blob store and the session records the confirmed byte range only after the bytes are stored. When the last chunk arrives the server checks length and checksum, writes the original to replicated storage, creates the video's metadata row in state UPLOADED and enqueues processing. The original is kept: new codecs and better encoders mean the platform will re-encode from the original for years. The file upload service design covers chunked upload mechanics in more depth.

Advertisement

Processing: a per-video DAG

Processing is a workflow, not a single job. Each video gets a directed acyclic graph of tasks: probe the input, split it into chunks at keyframes, encode each chunk into each rendition in parallel, join and package the results into segments and manifests, then generate thumbnails, captions from speech recognition, rights fingerprints and policy classifier scores. A workflow engine tracks each task's state so a failed encode of chunk 37 at 720p is retried alone.

def process_video(video_id, original):
    meta = probe(original)                       # duration, resolution, frame rate, audio
    chunks = split_at_keyframes(original, target_seconds=10)

    # Fan out: one task per (chunk, rendition). Cheap formats first, so the video
    # becomes watchable quickly; expensive codecs are added later if views justify it.
    first_ladder = ladder_for(meta, codec="h264")
    run_parallel(encode(c, r) for c in chunks for r in first_ladder)
    for r in first_ladder:
        join_and_package(video_id, r)            # concatenate chunks, write segments + manifest

    run_parallel([thumbnails(video_id, original),
                  speech_to_captions(video_id, original),
                  rights_match(video_id, original),       # fingerprint vs reference library
                  policy_classifiers(video_id, original)])

    mark_state(video_id, "PROCESSED")            # one row update, guarded by version
    schedule_if_popular(video_id, codecs=["vp9", "av1"])  # re-encode when demand justifies it

Chunking is what makes latency independent of length. A two-hour video split into 10-second chunks becomes 720 independent encodes per rendition, so with enough workers it finishes in roughly the time of one chunk plus overhead. The price is quality at chunk boundaries, since each encoder starts without knowledge of the previous chunk, and a join step that must produce clean segment boundaries.

The second idea is cost-aware encoding. Newer codecs such as VP9 and AV1 save bandwidth for the same quality but cost much more compute to encode. Spending that compute on a video watched ten times is waste; spending it on a video watched ten million times pays back quickly in egress. So the pipeline produces a widely supported baseline first and adds expensive formats when a video's demand crosses a threshold. At YouTube's scale even the baseline is expensive enough that Google built a custom ASIC, the Argos video coding unit, which it reported as 20 to 33 times more compute-efficient than its previous server-based encoding system, with ten encoder cores each able to encode 2160p at 60 frames per second in real time.

Rights and policy run in the same pipeline

A user-generated platform must decide what it is allowed to show. YouTube's Content ID system matches uploads against reference files that rights holders provide, and policy classifiers flag content for review. Architecturally these are just more tasks in the DAG, with one important property: their result can change a video's state after it was published. A rights holder can add a reference file next month that matches a video uploaded today.

So rights and policy decisions are stored as separate, versioned records attached to the video, and the watch path evaluates them at request time, for example per country, rather than baking them into the rendition files. That keeps a block, a claim or a monetisation change to a metadata write instead of a reprocessing job.

Metadata: sharded relational storage

Every video has a row with title, owner, state, privacy, durations, rendition list and counters, plus channels, playlists, comments and subscriptions around it. Relational queries and transactions are valuable here, so YouTube kept MySQL and, starting in 2010, built Vitess to shard it horizontally: a routing layer that sends each query to the shard holding the key, pools connections and allows resharding without application changes. Vitess is now an open-source CNCF project used well beyond YouTube.

The sharding key decides which queries are cheap. Sharding videos by video id spreads write load evenly and makes the watch page a single-shard lookup, but listing a channel's uploads becomes a scatter across shards unless you maintain a secondary index or a separate table sharded by channel id. The usual answer is to keep both: the canonical video row by video id and a channel-to-videos index by channel id, kept consistent asynchronously. Hot reads such as a viral video's metadata are absorbed by caches in front of the database, with short time-to-live values or explicit invalidation when the owner edits the title. The sharding architecture article explains resharding and cross-shard query options.

Counting views at scale

A view count looks like a single integer and behaves like a distributed systems problem. A viral video can receive tens of thousands of view events per second; one database row cannot take that write rate, and counts must exclude bots and repeated refreshes, which you cannot judge from one event in isolation.

import random

SHARDS = 64

def record_view(video_id, viewer, event_log):
    # 1. Durable, append-only event first: the source of truth for audits and recounts.
    event_log.append({"video": video_id, "viewer": viewer, "ts": now()})
    # 2. Fast approximate counter: spread hot keys across shards to avoid one hot row.
    shard = random.randrange(SHARDS)
    kv.incr(f"views:{video_id}:{shard}")

def display_count(video_id):
    cached = cache.get(f"views:{video_id}")
    if cached is not None:
        return cached
    total = sum(kv.get(f"views:{video_id}:{s}") or 0 for s in range(SHARDS))
    cache.set(f"views:{video_id}", total, ttl_seconds=30)
    return total

# A batch job later replays event_log, drops invalid views (bots, repeats, spam)
# and writes the audited count back; the public number converges to it.

The design separates three things. A durable append-only event log is the source of truth. A sharded approximate counter, spread across many keys and read through a cache, gives a fast and good enough number on the watch page. And a batch validation job replays the log, removes invalid views and writes back an audited count, which is what creators are paid on. The public number can lag or be adjusted downward, which is an explicit product decision: correctness for payouts and abuse resistance beat a perfectly live counter. Likes and subscriber counts follow the same pattern.

Delivery: getting the bytes close to viewers

Watch traffic dwarfs everything else, so delivery is the main cost centre. The player asks the watch API for metadata and a manifest, then fetches segments from edge caches over HTTP. Google runs caches inside its own points of presence and also places Google Global Cache servers inside ISP networks, so popular segments are served from the viewer's own ISP without crossing the internet.

The long tail shapes cache policy. Popular videos stay in the edge; a rarely watched video misses at the edge, misses at a regional tier and is fetched from the rendition store, so its first second is slower. Cache admission matters as much as eviction: letting every one-off request evict a popular segment lowers hit rate for everyone. Segment-level caching also helps, because most viewers watch only the opening minutes, so the first segments of a video are far hotter than the last. The CDN design article covers tiers and admission policies.

Worked example: storage and egress per day

Use the public figure of 500 hours uploaded per minute, which is 720,000 hours a day. Assume, for illustration, an original averaging 8 Mbit/s and a set of renditions across resolutions and codecs across about ten renditions that together average 12 Mbit/s.

QuantityCalculationResult
Original storage per hour8 Mbit/s x 3,600 s / 83.6 GB
Rendition storage per hour12 Mbit/s x 3,600 s / 85.4 GB
New data per day720,000 h x 9 GBabout 6.5 PB before replication
With 3 copies6.5 PB x 3about 19 PB per day
Encode work per day720,000 h x 10 renditions7.2 million rendition-hours

Even with rough assumptions the conclusions are firm. Storage grows by petabytes a day, so cold originals and unpopular renditions must move to cheaper storage classes and erasure coding rather than triple replication. Encoding is millions of rendition-hours a day, which is why encoder efficiency justified custom silicon. And because egress scales with hours watched, a codec that saves 30 percent of bits on the most-watched tenth of the catalogue saves more money than anything done to the other nine tenths.

Failure modes

FailureEffectDesign response
Upload interruptedCreator restarts a 3 GB uploadResumable sessions; confirm ranges only after durable write
Encoder task fails or hangsVideo stuck in processingPer-task retries with timeouts; DAG resumes from failed task
Bad encoder releaseCorrupt renditions at scaleCanary new encoders; keep originals to re-encode
Metadata shard overloadedWatch pages fail for one key rangeCaching, read replicas, resharding through the routing layer
Viral spike on one videoHot counter row, cache misses at edgesSharded counters; request coalescing at cache tiers
Late rights claimPublished video must be blocked in some countriesEvaluate rights at request time from versioned records
Regional cache outageTraffic shifts to other sites, origin load spikesSteering with capacity limits, origin shielding

Trade-offs

Process at upload time or on demand: preparing renditions at upload makes playback fast but spends compute on videos nobody watches; producing expensive formats only after demand appears saves compute at the cost of slightly worse quality for early viewers. Relational metadata with sharding versus a native distributed database: a routing layer over MySQL kept YouTube's existing schema and queries, while a distributed SQL store would trade migration effort for built-in resharding. Exact versus eventual counts: live exact counters are costly and gameable, while audited eventual counts are fair but surprise users when they lag. And building custom hardware only makes sense at a scale where a 20-fold efficiency gain outweighs years of chip design.

Related reading: recommendation system architecture for what decides which videos become hot.

What to do next

  1. Write the requirements for your platform with numbers: upload hours per minute, hours watched per day, and the share of views going to the top 1 percent of videos.
  2. Implement or adopt a resumable upload protocol and test it by killing the connection at random offsets.
  3. Model processing as a per-video DAG with chunk-level retries, and measure time from upload to first playable rendition.
  4. Define a demand threshold for expensive codecs and compute its payback from egress price and encode cost.
  5. Choose the metadata sharding key from your top three queries, and list which ones become scatter-gather.
  6. Split view counting into an event log, a sharded fast counter and an audited batch count.
  7. Run the capacity table above with your own bitrates and replication policy, and decide when renditions move to colder storage.
Key takeaway: A user-generated video platform has four planes that scale differently: resumable uploads, a per-video processing DAG, sharded metadata and edge delivery. Chunked parallel encoding makes processing time independent of video length, and cost-aware codecs spend compute only where views pay it back. Keep originals, evaluate rights and policy at request time, shard metadata by the key your hottest queries use, count views with an event log plus sharded counters and an audited batch job, and push popular segments as close to viewers as possible.