ASCELLA
Logo

Excellence In Execution
Precision In Performance

Ascella Group | Technology | 2026-06-20

The System Design Behind YouTube

YouTube handles 500 hours of video uploaded every minute, 1 billion hours watched daily, and 2.7 billion monthly active users through a combination of microservices architecture, adaptive video transcoding, Google Media CDN with over 3,000 edge locations, MySQL scaled horizontally via Vitess, Bigtable for petabyte-scale metadata, and a two-stage machine learning recommendation engine. Google acquired YouTube in 2006 for $1.65 billion. It is now worth over $150 billion and generates $38.7 billion in annual advertising revenue. The engineering decisions that make this scale possible are deliberate, well-documented, and instructive for any team building distributed systems.

Introduction: The Scale Problem Is Not Theoretical

Most distributed systems are designed to handle tens of thousands of requests per second. YouTube handles a different problem altogether.

Every minute, 500 hours of video are uploaded to YouTube. Every day, users generate one billion hours of watch time. YouTube Shorts alone accumulates 70 billion daily views. The platform serves 2.7 billion monthly active users across every device type, connection speed, and geography on the planet. Q3 2025 advertising revenue reached $10.26 billion for a single quarter.

At the beginning of 2025, YouTube hosted over 5.1 billion videos. That number has nearly doubled since 2021. The infrastructure required to store, process, index, serve, and monetise this volume of content is not an extension of ordinary web architecture. It is a distinct discipline of engineering.

When Google acquired YouTube for $1.65 billion in 2006, the platform was already struggling under its own growth. The engineering decisions made in the years following that acquisition, and the architectural principles that guide YouTube today, are among the most studied examples in distributed systems design. Understanding how YouTube works reveals patterns that apply to any system that needs to grow beyond what a single machine, database, or data centre can support.

This article covers the full system design behind YouTube: how video gets from a creator's device to a viewer's screen, what databases handle what types of data, how the recommendation engine works, and how the CDN achieves near-instant delivery to billions of people simultaneously.

Part 1: The Upload Pipeline

When a creator uploads a video, the process that follows is more complex than the progress bar suggests.

Chunked upload and ingestion

YouTube does not wait for a full video file to arrive before processing begins. The client splits the file into chunks, typically a few megabytes each, and uploads them in parallel. This chunked approach handles interruptions gracefully. If a connection drops mid-upload, only the incomplete chunks need to be re-sent rather than restarting the entire file transfer.

When chunks arrive at YouTube's ingestion layer, the upload service logs metadata to the MySQL database (managed through Vitess, covered in detail below). An event is emitted to a message queue, triggering the transcoding pipeline. The original file lands in Google Cloud Storage as the raw source.

Video transcoding

A single uploaded video file cannot serve the full range of YouTube's users. A viewer on a 5G connection in a city and a viewer on a 2G connection in a rural area cannot receive the same stream. A desktop browser and a mobile phone cannot play the same file format equally well.

YouTube addresses this through transcoding. Every uploaded video is converted into multiple formats and multiple quality levels. The standard output set includes resolutions from 144p up to 4K (2160p), and formats including H.264/MPEG-4 AVC and VP9. VP9 is YouTube's preferred codec for modern browsers because it achieves HD and 4K quality at roughly half the bandwidth required by older codecs, saving Google billions annually in infrastructure costs.

Transcoding is a compute-intensive process. YouTube runs it as a distributed batch job across its server infrastructure. Popular videos receive priority processing to minimise the time between upload and availability. For less popular content, transcoding may complete in lower-priority queues. Each video at upload time is assigned a unique identifier and processed by an automated pipeline that handles encoding, thumbnail generation, metadata extraction, transcript generation, and monetisation status assignment in parallel.

Without smart encoding and adaptive quality delivery, YouTube would need significantly more storage and bandwidth to serve the same content. The transcoding layer is one of the most consequential engineering decisions in the entire system.

Part 2: Storage Architecture

YouTube stores two fundamentally different types of data: the video files themselves (binary, large, rarely changed after upload) and everything else (metadata, user data, engagement signals, search indexes). These require different storage systems.

Video file storage

Raw video files and all transcoded renditions are stored in Google Cloud Storage. This is a distributed object store designed for durability and scale. Google's private fibre network, which connects its data centres globally, allows video files to move between storage locations without using the public internet. This network is one of the largest private fibre networks in the world.

Popular videos are copied to edge cache nodes (covered in the CDN section) so they can be served from a location close to the viewer. Less popular videos stay in centralised storage and are fetched on demand when requested.

Relational data: MySQL via Vitess

YouTube uses MySQL as its primary relational database for structured data: user accounts, channel information, subscriptions, video metadata, and relationships between entities. The challenge is that MySQL was not designed to handle the query volume YouTube generates at planetary scale.

YouTube solved this by building Vitess. Vitess is an open-source sharding middleware layer that sits between application code and MySQL. It handles query routing, connection pooling, and horizontal sharding transparently. Application code talks to Vitess as if it were a single MySQL instance. Vitess distributes queries across many MySQL shards behind the scenes, allowing the system to scale to millions of queries per second without changing application logic.

YouTube open-sourced Vitess in 2012. It is now a Cloud Native Computing Foundation (CNCF) project and is used by companies including Slack, Pinterest, HubSpot, and GitHub. The fact that YouTube built it and then released it to the open-source community reflects how central it is to YouTube's own operational stability.

Metadata and activity data: Bigtable

Google Bigtable is a wide-column NoSQL database designed for massive-scale time-series data. YouTube uses it for watch history, engagement events (views, likes, comments, shares), user activity logs, and analytics data. Bigtable handles high-volume write workloads efficiently, making it well-suited for recording the constant stream of interactions YouTube generates across billions of sessions simultaneously.

Caching: Memcache and Redis

Serving every request from the database would be prohibitively slow and expensive. YouTube maintains caching layers using Memcache for high-speed retrieval of frequently accessed data: view counts, session data, feed results, and autocomplete suggestions. The cache hit rate on hot data reaches approximately 90%, meaning nine out of ten requests for popular data return from memory rather than touching the database. Zookeeper handles node coordination across the distributed cache infrastructure.

Part 3: Video Delivery and the CDN

Getting a video from storage to a viewer's screen in milliseconds, at consistent quality, for billions of simultaneous viewers, is the core delivery challenge YouTube has spent years solving.

Google Media CDN

YouTube's content delivery network is not a third-party service. It is Google Media CDN, built on Google's own global network infrastructure. Google owns the fibre connecting its data centres and edge nodes, which means most of a video's journey to a viewer bypasses the public internet entirely. This gives YouTube far more control over latency and reliability than companies that depend on external CDN providers.

Google Media CDN spans over 3,000 edge locations globally. The system achieves cache hit ratios of 98.5 to 99% for video-on-demand content. In practice, this means that when a viewer presses play, the video data is almost certainly already sitting at a node near them rather than being fetched from a central server halfway around the world.

The CDN does not simply store copies of every video. It uses machine learning to predict which videos will be in demand based on trending signals, time of day, geography, and historical patterns. Content is pre-positioned at edge nodes before viewers request it, reducing origin fetch rates and improving first-play latency.

YouTube operates over 500,000 servers globally to support its infrastructure. This includes specialised hardware for video encoding, content serving, database operations, and machine learning workloads.

Adaptive Bitrate Streaming

Even with the CDN delivering content from nearby edge nodes, network conditions vary. A viewer's connection may be fast during the first ten seconds of a video and slow during the next thirty. Playing a fixed-quality stream under these conditions results in buffering.

YouTube uses two adaptive streaming protocols to handle this: MPEG-DASH (Dynamic Adaptive Streaming over HTTP) and HLS (HTTP Live Streaming). Both protocols work on the same principle. The video is divided into short segments, typically two to ten seconds each. The player continuously monitors available bandwidth and selects the appropriate quality level for each segment. If bandwidth drops, the next segment loads at lower quality. If bandwidth improves, quality rises again. The viewer sees this as automatic quality adjustment rather than a loading spinner.

MPEG-DASH is the primary protocol for modern browsers and devices. HLS is used for Apple devices and some legacy environments. Supporting both ensures YouTube works well across the full range of hardware and operating systems its users run.

Part 4: The Recommendation Engine

YouTube's recommendation system is one of the most studied and most consequential machine learning systems in production. It is also one of the primary reasons users spend an average of 49 minutes per day on the platform.

The recommendation system operates in two stages.

Stage 1: Candidate generation

YouTube has over 5.1 billion videos. Evaluating all of them for every user on every page load is not feasible. The candidate generation stage narrows this down to a manageable set of a few hundred videos that are broadly relevant to a specific user.

This stage uses the user's watch history, search history, demographic data, and engagement signals as inputs to a neural network trained on billions of historical interactions. The output is a ranked shortlist of candidate videos. The key goal here is recall: include everything that might be relevant, even if it means including some that ultimately will not be shown.

Stage 2: Ranking

The ranking model takes the candidate set from stage one and scores each video against a more detailed set of signals: predicted watch time, click-through likelihood, freshness of the content, diversity considerations, and user-specific engagement patterns. The output is the ordered list of videos shown in the feed, sidebar, and homepage.

YouTube uses TensorFlow as the machine learning framework for training and serving these models. The recommendation system processes billions of user interactions in real time to keep signals current. A video that starts trending in one geography will see its ranking signals update within minutes as views accumulate.

One documented finding from YouTube's own research is that the system optimises primarily for watch time rather than clicks. A video a user watches for ten minutes contributes more positive signal than one they click on and abandon after five seconds. This design decision shapes the entire content recommendation landscape on the platform.

Content moderation at scale

The recommendation system also has a safety dimension. In 2025, YouTube removed approximately 7 to 9 million videos per quarter that violated its community guidelines. The vast majority were removed before accumulating significant views. YouTube uses approximately 20,000 content moderators worldwide through a combination of direct employees and contractors, alongside machine learning classifiers trained on labelled content. The automated systems flag content at the point of upload, reducing the volume that reaches human review queues.

Part 5: Microservices and Event-Driven Architecture

YouTube does not run as a single application. It is composed of many independent services, each responsible for a specific function, communicating through well-defined interfaces.

The core services include:

  1. Upload service: handles chunked file ingestion and storage

  2. Transcoding service: manages encoding jobs across multiple workers

  3. Metadata service: reads and writes video and channel information via Vitess

  4. Search service: handles query processing and result ranking

  5. Recommendation service: serves personalised feed content

  6. Notification service: sends alerts for new uploads, comments, and live streams

  7. Analytics service: processes engagement events and surfaces creator metrics

  8. Streaming service: manages adaptive delivery to players

Each service can be scaled independently based on demand. During a major live event, the streaming and notification services handle significantly higher load than normal while the transcoding service remains at baseline. In a monolithic architecture, the entire system would need to scale together. With microservices, each layer scales proportionally to the actual demand it is experiencing.

Apache Kafka handles event streaming between services. When a video upload completes, an event is published to a Kafka topic. The transcoding service, metadata service, and notification service each consume that event and act on it independently and in parallel. This decouples the services from one another and ensures that the failure of one does not block the others.

YouTube's backend services run on Google's internal compute infrastructure, using Borg-style orchestration (the precursor to Kubernetes) for container scheduling and workload management. Distributed tracing follows the pattern established by Google's Dapper system, allowing engineers to trace individual requests across services and diagnose latency issues in production.

Part 6: Search and Indexing

YouTube is the second-largest search engine in the world. More searches are conducted on YouTube than on Bing, Yahoo, and every other search engine except Google itself.

The search system maintains an inverted index over video titles, descriptions, transcripts, tags, and engagement signals. When a user submits a query, the search service retrieves matching candidates, ranks them using a combination of relevance signals and personalisation, and returns results within milliseconds.

Autocomplete suggestions are served from cache rather than the search index directly. When a user begins typing, the autocomplete service queries a pre-built set of common search phrases stored in memory. This allows results to appear with sub-100 millisecond latency without hitting the search index for every keystroke.

Video transcripts improve search quality significantly. When a creator uploads a video, YouTube generates an automatic transcript using speech recognition. These transcripts are indexed alongside other metadata, meaning a search for a phrase spoken in a video can surface that video even if the creator never typed those words in the title or description

Part 7: Global Infrastructure and Reliability

YouTube's service level expectations require near-continuous availability for users in every country. An outage that takes YouTube offline for one hour affects more than 2.7 billion potential users and has significant advertising revenue consequences.

The architecture distributes this risk across multiple layers.

Regional data centres serve as the origin for content not yet at edge nodes. Multiple Availability Zones within each region ensure that the failure of a single data centre does not take the regional service offline. YouTube's Google-owned fibre network allows traffic rerouting at the infrastructure level rather than the application level, reducing the impact of physical network failures.

Read operations are served from replicated database instances. Write operations go to the primary MySQL shard and replicate to secondaries. There is an acceptable window of temporary inconsistency in some signals: a slight discrepancy in a video's view count between the primary and a read replica is acceptable. An incorrect search result or a missing video from a channel page is not. The system design reflects these different consistency requirements by choosing the right database and replication strategy for each type of data.

YouTube's cache infrastructure absorbs a very large portion of read traffic. With a 90% hit rate on hot data, the majority of requests never reach the database at all. This acts as a natural circuit breaker. During traffic spikes, the cache absorbs the surge before it reaches the storage layer.

Key Design Principles Worth Carrying Forward

The architecture behind YouTube contains lessons that apply broadly to any engineering team building systems intended to grow.

Design for the read path first. YouTube serves vastly more read requests than write requests. Its database strategy, caching layers, and CDN are all optimised for read performance. Teams building data-heavy systems should understand which operations dominate and design storage and caching accordingly.

Accept eventual consistency where it is appropriate. Not every number on YouTube needs to be precisely accurate at every moment. View counts can lag. Watch history can be slightly delayed. Knowing which data can tolerate eventual consistency and which cannot allows the architecture to scale specific components without enforcing strict consistency across everything.

Move data close to the user. Google Media CDN's 3,000-plus edge locations exist to reduce the distance data travels. Every millisecond of latency reduced at the CDN layer improves the experience for billions of users. The principle applies at smaller scales too: caching frequently accessed data closer to the application tier reduces load on the database and speeds response times.

Build purpose-built tools when generic ones cannot scale. Vitess exists because MySQL at YouTube's query volume needed custom sharding middleware. Building and open-sourcing Vitess was a significant investment. It also solved a real problem that no existing tool addressed adequately. When a team encounters a genuine scaling wall in a standard tool, the right answer is sometimes a new layer rather than a bigger server.

Decouple services through event-driven communication. Kafka's role in YouTube's architecture means that the upload service, transcoding service, and notification service operate independently. Failures in one do not cascade to others. New services can be added by subscribing to existing events without changing the services that produce them.

Frequently Asked Questions

How does YouTube handle 500 hours of video uploaded per minute? YouTube accepts uploads as chunked file transfers, stores raw files in Google Cloud Storage, and triggers a distributed transcoding pipeline via a message queue. Multiple transcoding workers process files in parallel, converting each upload into multiple quality levels and formats. Popular videos receive priority in the processing queue to minimise time-to-availability.

What database does YouTube use? YouTube primarily uses MySQL for relational data such as user accounts, subscriptions, and video metadata. MySQL is scaled horizontally through Vitess, an open-source sharding middleware developed by YouTube and released to the public in 2012. Bigtable handles large-scale NoSQL data including watch history and activity logs. Memcache provides in-memory caching with approximately a 90% hit rate for frequently accessed data.

What is Vitess and why did YouTube build it? Vitess is a database sharding middleware layer that sits between application code and MySQL. It allows YouTube to distribute queries across many MySQL instances transparently, enabling millions of queries per second without changing application logic. YouTube built it because standard MySQL could not handle its query volume. Vitess is now a CNCF open-source project used by Slack, Pinterest, GitHub, and many others.

How does YouTube deliver video so quickly to users worldwide? YouTube uses Google Media CDN, which spans over 3,000 edge locations globally on Google's own fibre network. Popular videos are cached at nodes close to viewers, achieving cache hit ratios of 98.5 to 99%. Machine learning algorithms pre-position trending content at edge nodes before users request it, reducing first-play latency.

How does YouTube's recommendation system work? The recommendation system works in two stages. Stage one uses neural networks trained on billions of interactions to narrow billions of videos down to a few hundred relevant candidates for each user. Stage two ranks those candidates using predicted watch time, click-through rate, content freshness, and personalisation signals. TensorFlow handles model training and serving. The system optimises for watch time rather than clicks, a design choice that shapes which videos get recommended.

Why does YouTube use both MPEG-DASH and HLS? MPEG-DASH is the primary adaptive streaming protocol used for most modern browsers and Android devices. HLS is used for Apple devices and some legacy environments. Both protocols divide video into short segments and allow the player to select quality levels dynamically based on current bandwidth. Supporting both protocols ensures consistent performance across the full range of devices YouTube users run.

Conclusion

YouTube's system design is not a single clever idea. It is the accumulated result of many careful decisions made under real pressure, at real scale, over two decades of operation.

The upload pipeline converts raw files into dozens of formats and quality levels within minutes. Google Cloud Storage and Bigtable handle petabytes of new data every day. Vitess makes MySQL work at a scale MySQL was never intended for. Google Media CDN puts video data within milliseconds of billions of viewers simultaneously. The two-stage recommendation engine processes billions of interactions to surface the most relevant content from a library of over five billion videos. Microservices and Kafka decouple the system so each component can fail, scale, or evolve independently.

The numbers behind YouTube in 2026 are extraordinary: 2.7 billion monthly users, 1 billion hours watched daily, $38.7 billion in advertising revenue in 2025. Every one of those numbers is supported by engineering decisions that chose the right tool for each layer of the system and designed each layer to be extended without being rebuilt.

That is what good system design does. It makes the next order of magnitude of growth a series of measured investments rather than an emergency rewrite.

At Ascella Group, we help engineering teams design systems that scale with confidence. Whether you are building an upload pipeline, choosing a database strategy, or designing a content delivery architecture, the principles behind YouTube are worth understanding before you start.