Design Netflix
Stream a fixed catalog to 40 M people at once, 200 Tbps at peak, starting in under two seconds, from cache servers placed inside the viewers’ own internet providers.
Last updated 2026-09-30. Difficulty: hard. Patterns: cdn, proactive-caching, adaptive-bitrate, steering, resilience. Reported at Netflix, Amazon, Google, Meta, Apple.
The interviewer asks, the candidate answers and draws, and you press Next. Pause to answer yourself at the key decisions, and ask the AI Mentor anything along the way.
Functional requirements
- Play a title. Press play and video starts in under two seconds, then adapts quality to the connection without stalling.
- Resume where you left off. Across devices: stop on the TV, continue on the phone at the same second.
- Personalised homepage. Rows such as "Continue watching" and "Because you watched…", different for every profile.
- Many devices. TVs, phones, browsers and consoles, each with different codecs, screen sizes and DRM systems.
- Content protection. Studios require DRM: the player gets a decryption license only if the account is entitled to the title in that country.
- Out of scope. Encoding and transcoding the catalog (see Design YouTube), billing, search, downloads for offline viewing, live events.
Non-functional requirements
- Scale (40 M concurrent streams at peak). Evening peaks in each region, roughly 100 M daily viewers watching two hours each.
- Throughput (~200 Tbps at peak). This is the requirement that drives the hardest trade-off: no transit link or commercial CDN deal handles this economically, so the bytes have to be served from inside the ISPs, close to viewers.
- Startup time (p90 < 2 s from press to first frame). Manifest, license and the first video segment all sit on this path.
- Rebuffering (< 0.1 % of viewing time). Stalls are the metric viewers feel most. Quality can drop; playback must not stop.
- Availability (99.99 % for playback). An AWS region can fail and viewers in it must keep watching. The control plane is multi-region active-active.
- Catalog churn. The catalog is known in advance and changes slowly (new releases are scheduled), which is what makes proactive caching possible.
Back-of-envelope estimates
- Peak egress: ~200 Tbps. 40 M concurrent streams × ~5 Mbps average (a mix of HD and 4K, with efficient codecs) = 200 Tbps. A single 100 Gbps link carries 0.05 % of that.
- Cache servers needed: ~2,000 minimum, several thousand in practice. A modern flash appliance serves ~100 Gbps. 200 Tbps ÷ 100 Gbps = 2,000 at full load; with headroom for failures and spread across thousands of ISP sites, several times that.
- Catalog size, all encodings: ~1 PB. About 50 k hours of content × ~40 Mbps summed across every bitrate, codec and audio track ≈ 50 k × 3 600 s × 5 MB/s ≈ 0.9 PB. Every copy of the full catalog is about a petabyte.
- Share of catalog one appliance holds: ~25 % of bytes, ~95 % of views. An appliance holds ~250 TB, about 25 % of a 1 PB catalog. Viewing is heavily skewed toward recent and popular titles, so that quarter serves the vast majority of plays from that site.
- Playback heartbeats per second: ~1.3 M. Each stream reports progress every 30 s: 40 M ÷ 30 ≈ 1.3 M events/s. Most of them only overwrite the resume position.
- Play starts per second at a big launch: ~100 k. If 6 M people press play in the first minute of a big release: 6 M ÷ 60 ≈ 100 k/s hitting playback, steering and licensing at once. The control plane is sized for this spike, not the average.
Components
- Player (TV · phone · browser): Asks the control plane where and how to play, then fetches video segments over HTTP straight from cache servers, choosing the bitrate segment by segment from its buffer level and measured throughput.
- ISP-embedded caches (Open Connect appliances): Cache servers Netflix gives to ISPs to install inside their networks, free. They hold the titles most likely to be watched in that area and serve most of the bytes without crossing the ISP’s border.
- Exchange-point caches (larger OCA clusters at IXPs): Bigger clusters at internet exchange points that hold a wider slice of the catalog. They serve ISPs without embedded appliances, act as the fallback tier, and fill the embedded ones.
- Origin storage (S3 · every encoding): The master copy of every encoded file in AWS. Viewers never stream from it; it only fills the exchange-point tier.
- Fill planner: Each day, predicts what will be watched where and computes, per appliance, the exact set of files it should hold. Appliances download changes during their off-peak window.
- Steering service: Chooses which cache servers a given player should use: the ones that hold the file, are healthy and lightly loaded, and are closest in network terms (learned from the ISP’s routing via BGP). Returns a ranked list, not one server.
- API gateway (Zuul-style edge): The front door to the control plane in AWS: authenticates the device, routes to services, sheds load and fails over between regions. Never carries video.
- Playback service: Checks the account is entitled to the title in this country, picks encodings the device supports, fetches the resume position, asks steering for servers and returns a manifest of segment URLs.
- DRM license service: Issues a short-lived decryption license bound to this device and session after checking entitlement again. On the startup path, so it must be fast and must survive launch spikes.
- Homepage service: Assembles the homepage for a profile from precomputed rows, adds live rows such as "Continue watching", and falls back to popular rows if personalisation is unavailable.
- Precomputed rows (per profile): Each profile’s ranked rows and titles, recomputed in batch every few hours and on significant events. Reading a homepage is a lookup, not a model run.
- Recommendation jobs: Batch and near-real-time pipelines that score titles for every profile from viewing history and write the ranked rows. Also feed the fill planner with expected popularity by region.
- Playback events (Kafka): Heartbeats, quality switches, rebuffers and errors from every player, about 1.3 M a second at peak. Feeds viewing history, recommendations, quality monitoring and the fill planner.
- Viewing history: Keeps each profile’s resume position per title and the list of what it has watched. Serves "Continue watching" and the resume point at play time.
- History store (Cassandra · key = profile): Wide rows per profile: the latest position per title (overwritten, so it stays small) and a compacted log of what was watched. Write-heavy and multi-region.
User flows
- Press play: first frame in two seconds. Two round trips to the control plane and then pure HTTP to a cache a few milliseconds away. Everything on this path is tuned for startup time.
- The player asks the playback service to play a title.
- Playback checks entitlement and fetches the resume position. Entitlement is per country, because licensing deals are. The resume point is one key read from the profile’s history row.
- Steering returns a ranked list of cache servers that hold the files, close to this viewer. Steering knows which appliance holds which files (from the fill plan), their health and load, and which are reachable inside the viewer’s ISP (from BGP routes the appliances report). It returns two or three choices so the player can fail over without asking again.
- The player requests a DRM license for this session. Done in parallel with fetching the first segments, so the license round trip overlaps the download. The license is bound to the device and expires with the session.
- The player fetches the first segments at a conservative bitrate from the top-ranked server, then ramps up. Starting low (one or two seconds of video at ~1 Mbps) gets the first frame up fast; the adaptive bitrate logic then climbs as the buffer fills. Segments are plain HTTP range requests on large pre-encoded files.
- Every 30 seconds the player sends a heartbeat with its position. The history service overwrites the resume point for (profile, title). Losing one heartbeat costs 30 seconds of resume accuracy, so these are fire and forget.
- Open the app and see your homepage. Personalisation is computed ahead of time, so the homepage is a read. If the recommendation system is down, the page still renders.
- The app requests the homepage for a profile.
- The homepage service reads the profile’s precomputed rows. Rows are written by batch jobs every few hours and refreshed after notable events (finishing a series). Reading them is a single key lookup.
- "Continue watching" is built live from viewing history. This row must reflect what you watched ten minutes ago on another device, so it is not precomputed: it reads the latest positions directly.
- If personalisation is unavailable, the service falls back to regional popular rows. Graceful degradation: a generic but working homepage beats an error page. The fallback rows are cached everywhere and cheap to serve.
- A hit season launches at midnight. The scale-breaking case. Because the release is scheduled, the bytes are already in every ISP before anyone presses play; the risk moves to the control plane.
- Days before release, the planner marks the new episodes as high demand everywhere they will launch. Expected demand comes from the size of the previous season’s audience and pre-release interest. The files are added to the file sets of thousands of appliances.
- Appliances download the episodes during their nightly off-peak windows. Filling uses the ISP’s spare capacity at 4 a.m., not its busiest hours. Appliances in the same site fill from each other where possible, so each file crosses the ISP’s border once.
- At midnight, millions press play within minutes. About 100 k play starts a second. The control plane was scaled up ahead of the scheduled release, and non-essential work (some personalisation, prefetching) is shed first if it struggles.
- Steering spreads viewers across every appliance that holds the episode. Load is part of the ranking, so a busy appliance drops down the list. Because every site already holds the file, no request has to travel far for it.
- A cache server fails mid-stream. The player handles it alone: a buffer to ride through, a list of alternative servers, and a lower bitrate if the network is the problem.
- Segment requests to the appliance start timing out. The player has 30 to 60 seconds of video buffered, so a short outage is invisible as long as it reacts within that window.
- The player switches to the next server in its list. No call to the control plane is needed, which matters when the failure is widespread and everyone would otherwise ask at once. Segment URLs differ only by host.
- If throughput drops instead, adaptive bitrate steps down rather than stalling. The bitrate choice for each segment is driven mainly by how full the buffer is: as it drains, the player picks smaller encodings. Dropping from 4K to 1080p is a quality blip; a stall is a failure.
- The appliance’s health reports stop; steering removes it from new lists. Steering also learns from player error reports flowing through the event stream, which catch failures that health checks miss, such as a congested link in the ISP.
- Quality monitoring flags the rise in switches and rebuffers for that ISP. Rebuffer rate and startup time per ISP and per appliance are the operational dashboards. A regional dip is often the first sign of an ISP problem the ISP has not noticed yet.
- The nightly fill. Proactive caching: decide in advance what every appliance should hold, and move the bytes when networks are idle.
- Recommendation and viewing data produce expected demand per title per region. Demand is predictable: what was popular yesterday, what is new, what is being recommended prominently. The planner works per file (title × encoding), because a 4K file is only worth caching where many 4K devices watch.
- The planner computes each appliance’s file set to maximise the share of predicted bytes it can serve. A knapsack-style problem per appliance: given 250 TB, choose files with the highest expected bytes served per byte stored. Appliances in one site are planned together so they hold complementary sets.
- During its off-peak window each appliance downloads what it is missing and deletes what is no longer planned. Fills come from peers in the same site first, then the exchange-point tier, then origin, so the expensive paths are used least. Files are immutable, so there is nothing to invalidate.
- The appliance reports its new contents and health to steering. Steering only sends a player to an appliance that has confirmed it holds the files, so a fill that failed halfway is harmless.
Deep dives
Why put servers inside ISPs
Why not just buy a commercial CDN?
200 Tbps at peak is a large fraction of all internet traffic in the evening. Through a commercial CDN, every byte is billed, and it still crosses the ISP’s border at a peering point, which is often where congestion happens. The ISP has to carry it across its own backbone too.
The catalog is fixed and known in advance, and demand is predictable. That makes it possible to put the right bytes physically inside each ISP before anyone asks for them, which helps both sides: the ISP’s backbone and transit carry far less, and viewers get bytes from a few milliseconds away.
- Own appliances embedded in ISPs, plus larger clusters at exchange points chosen
- Commercial CDN situational: early in a service’s life, or as overflow during unexpected spikes
- Serve from your own cloud regions rejected
The answer: A private CDN in two tiers. Appliances (dense flash or disk servers, ~100 Gbps each) are given to ISPs to install inside their networks; they serve the predicted popular slice for that region. Larger clusters at internet exchange points hold more of the catalog, serve ISPs without embedded appliances, and are the fallback and fill source. Origin in S3 holds everything but serves only the exchange tier. This is Netflix Open Connect, which serves essentially all of its video traffic; the control plane (accounts, playback decisions, recommendations) runs separately in AWS.
Why would an ISP agree to host your servers?
Because it saves them money and complaints. Without the appliance, the same video crosses their transit or peering links and their backbone at peak hours. With it, the traffic stays on the last mile. Netflix provides the hardware for free; the ISP provides space, power and a network port.
How would a small startup do this?
It would not. Use a commercial CDN and pay per byte until the bill clearly exceeds the cost of building your own. The design principle to keep is the same: immutable files, aggressive caching, and steering clients to good servers.
What happens in a region with no embedded appliances?
Viewers are steered to the nearest exchange-point cluster, reached over peering. Startup time and peak-hour quality are somewhat worse, which is exactly the data used to decide where to place the next appliances.
Proactive fill instead of caching on demand
A normal CDN fills on a miss. Why decide in advance what each server holds?
A pull-through cache fetches a file the first time someone asks for it. For video that has two costs: the first viewer in each place waits for a multi-gigabyte file to be pulled, and the fill traffic happens exactly when demand is highest, at peak hours, on the links you were trying to protect.
Netflix’s catalog does not change minute to minute, and what people will watch tomorrow evening is highly predictable from what they watched today and what is launching. That makes it possible to fill at 4 a.m. instead of 9 p.m.
- Planned daily fill during off-peak windows, from peers first chosen
- Pull-through caching with LRU eviction situational: for the long tail at exchange-point clusters, and for unexpected hits
- Replicate the whole catalog everywhere rejected
The answer: Every day a planner turns expected demand per title, per encoding, per region into a file set for each appliance, choosing files with the most predicted bytes served per byte stored, and planning servers in the same site to hold complementary sets. Each appliance has a configured off-peak window in which it fetches missing files (peer, then exchange tier, then origin) and deletes ones no longer planned. Files are immutable, so there is never an invalidation. Steering only uses an appliance for files it has confirmed it holds. Exchange-point clusters additionally keep pull-through caching for the long tail.
A small documentary suddenly goes viral. What happens?
Tonight it is served from the exchange tier, and from appliances that happened to hold it, with somewhat longer paths. Steering spreads the load and players adapt bitrate. Tomorrow’s plan sees the demand and places it widely. For extreme cases, the planner can trigger an urgent fill outside the window at a throttled rate.
How do you measure whether the plan is good?
The share of bytes served from the viewer’s own ISP (the "offload" rate) per region, and how much fill traffic it took to get there. If offload drops, the prediction or the placement is wrong, and the per-title misses show which.
Why plan per encoding, not per title?
Because a title exists in dozens of files: each resolution, codec and audio language. In a region where most devices are phones, the 4K HEVC file is not worth the space; the 720p AV1 file is. Planning per file uses each terabyte where it serves the most bytes.
Choosing a server for each viewer
With thousands of cache servers, how does a player know which one to stream from?
DNS-based CDNs pick a server from the resolver’s location, which is often wrong (a viewer using a public DNS resolver looks like they are somewhere else) and coarse. Here the choice matters a lot: the right appliance is inside the viewer’s own ISP, a few milliseconds away; the wrong one is across a congested peering link.
The control plane already sees every play request, knows exactly which files each appliance holds, and can learn the ISP’s routing from the appliances themselves. So it can make a precise decision per request instead of relying on DNS.
- Per-request steering in the control plane, using appliance-reported routes, contents, health and load; return a ranked list chosen
- DNS-based geolocation rejected
- Anycast to the nearest server situational: for small objects such as images and API traffic, not long video streams
The answer: Each appliance runs a BGP session with its ISP to learn which client prefixes it can reach directly, and reports those routes, its health, load and file set to steering. On a play request, steering takes the client IP and the files needed, filters appliances that hold them and can reach that prefix, ranks by network proximity, health and load, and returns the top few. The manifest contains those hosts; the player uses the first and falls back down the list itself. Player-side error and throughput reports feed back into ranking, catching problems health checks miss.
Why give the player several servers instead of one?
So failover is local and immediate. If the first server stalls, the player switches in hundreds of milliseconds using its buffer, without a control-plane call. During a large failure, that also prevents millions of players stampeding the steering service at once.
How do you avoid overloading the best appliance in a site?
Load is a ranking input, and steering spreads requests across appliances with the file in proportion to spare capacity. Planned complementary file sets also mean popular files are on several servers in each site.
Steering is down. Can anyone start watching?
Players that already have a manifest keep streaming. For new plays, the playback service falls back to a cached default mapping from client network to a few large exchange-point clusters, which are less optimal but always hold the popular files. Startup and quality dip slightly; playback keeps working.
How fresh does steering's view of each box need to be?
Health and load within seconds, because they change quickly; file sets within minutes, since they only change during fills; routes as BGP updates arrive. Stale load data is the most dangerous, so boxes push it continuously and steering treats a silent box as unhealthy.
Starting fast and never stalling
How does the player pick a bitrate, and how do you get the first frame up in under two seconds?
Each title is encoded at a ladder of bitrates, from a few hundred kilobits to over 15 Mbps for 4K. Picking too high causes stalls when the network dips; too low wastes a good connection. Networks change second to second, especially on phones.
Startup has its own budget: the manifest request, the license request and enough video to show the first frame all sit between the press and the picture. Starting at the top bitrate means downloading megabytes before anything shows.
- Buffer-based adaptive bitrate, starting low; license fetched in parallel; per-title encoding ladders chosen
- Throughput-only adaptation situational: for the startup phase before the buffer has data, combined with buffer-based logic afterwards
- Fixed bitrate chosen by the user rejected
The answer: The player starts at a conservative bitrate so the first segment is small, and requests the DRM license in parallel with the first download. Once playing, it chooses each segment’s bitrate mainly from buffer occupancy (a full buffer allows climbing, a draining one forces stepping down), using throughput only as a guide. Encodings use per-title, even per-shot, ladders: a cartoon looks perfect at a third of the bitrate a grainy action film needs, which Netflix has written about extensively. Device capability decides the codec (AV1 or HEVC where supported, H.264 otherwise).
Why is a buffer-based approach better than estimating bandwidth?
Bandwidth estimates on mobile networks are noisy and lag reality. The buffer level is a direct measure of whether you are keeping up: if it grows you can afford more, if it shrinks you cannot. Using it as the main signal avoids both stalls and pointless oscillation.
How would you cut startup time further?
Prefetch the manifest and license when the user hovers over a title or opens its details page, keep a warm connection to the likely appliance, and start from a cached low-bitrate opening segment. Measure startup at p90 per device type, since old TVs are usually the slowest.
Why are segments a few seconds long?
Shorter segments let the player react faster to network changes and start sooner, but each request has overhead and hurts compression at segment boundaries. Two to six seconds balances responsiveness against efficiency; the first segments can be shorter to speed up the start.
How do you handle a viewer on a train going through tunnels?
The buffer is the defence: build a large one when the network is good, step down aggressively when throughput collapses, and keep requesting low-bitrate segments rather than waiting for a high one. Mobile players also cap the top bitrate to save data unless the user opts in.
Recording 1.3 million heartbeats a second
Every stream reports its position every 30 seconds. Where does that go, and how does "Continue watching" stay correct?
At peak, 40 M streams × one heartbeat per 30 s is about 1.3 M writes a second. Nearly all of them say "this profile is now at second 1,864 of this title", which only matters until the next heartbeat. A relational database storing each as a row would be enormous and slow.
Two different questions are being answered: "where do I resume?" needs only the latest position per (profile, title), and "what has this profile watched?" needs a compact history. Treating them separately makes the write path cheap.
- Heartbeats through Kafka into a wide-column store: overwrite latest position, compact history chosen
- One row per heartbeat in a relational database rejected
- Position only on the client rejected
The answer: Players send heartbeats into Kafka keyed by profile. The history service consumes them and writes to a Cassandra table keyed by profile, with one column per title holding the latest position and timestamp: an overwrite, so the row does not grow with viewing time. A separate compacted record of titles watched (with completion) feeds recommendations and the "watched" markers; old detail is rolled up into summaries. "Continue watching" reads the profile row directly. Netflix has described running viewing history on Cassandra for exactly this write-heavy, per-member pattern.
You pause on the TV and immediately open the phone. Will it resume at the right place?
Usually within a few seconds, which is what users expect. The player also sends an extra heartbeat on pause and on exit, so the latest position is written immediately rather than up to 30 seconds late.
Two devices on the same profile play the same title at once. What position wins?
The most recent write by timestamp, which is the device that reported last. That is acceptable because it is rare; a stricter rule (only the device that is actually playing may advance the position) can be enforced with the session id if needed.
How do you delete a profile's viewing history for privacy?
Delete the profile's row in the history store (a single partition, so it is cheap) and write a tombstone event so downstream systems such as recommendations purge their copies. Raw events in Kafka expire by retention; long-term analytics should only ever hold pseudonymised, aggregated data.
What happens if Kafka is unavailable for a few minutes?
Players buffer recent heartbeats locally and resend; the history service applies only the newest position, so a burst of late events is harmless. Resume points may lag by minutes for that window, which users rarely notice.
Surviving a region failure
An AWS region goes down at peak. What happens to the people watching?
Streams already playing keep going: the video comes from cache servers in ISPs, not from AWS, and players have tens of seconds buffered. The risk is the control plane: new plays, licenses, homepages, and heartbeats all depend on services that run in a region.
If one region fails and its traffic moves elsewhere, the other regions must absorb a sudden extra load, with caches that are cold for those users.
- Active-active control plane in several regions, with evacuation and graceful degradation chosen
- Active-passive failover rejected
- Single region with many zones rejected
The answer: Run the control plane active-active in three or more regions, with member data replicated across them (Cassandra multi-region), and route each device to the nearest healthy region at the edge. When a region is unhealthy, shift its traffic to the others by DNS and edge routing; each region is provisioned with headroom for this, and evacuations are rehearsed regularly. Services declare fallbacks (popular rows instead of personalised ones, cached entitlement) so the essential path, press play and watch, keeps working when anything else is slow. Netflix popularised regularly breaking production on purpose (Chaos Monkey and region evacuation exercises) to prove this works.
Which features would you sacrifice first under overload?
Anything not needed to start or continue playback: personalised artwork, some recommendation rows, profile animations, non-essential telemetry. Licensing, entitlement and steering are protected with priority and load shedding at the gateway.
How do you know evacuation will work on the bad day?
By doing it on good days: shift a region’s traffic away on a schedule, watch error rates and latency in the receiving regions, and fix what breaks. A failover path that is never exercised is a failover path that does not work.
What data is hardest to keep consistent across regions?
Anything with a single correct answer at a moment: account and payment state, and entitlement. Those are written in one home region and replicated, with reads tolerating short staleness. Viewing history and preferences can be eventually consistent everywhere because a few seconds of lag is invisible.
How do you avoid a retry storm when a region comes back?
Clients use exponential backoff with jitter, the gateway limits concurrency per service, and traffic is shifted back gradually rather than all at once. Circuit breakers in callers stop hammering a dependency that is still warming up.
The theory behind it
- CDNs and edge computing: How a content delivery network works, push vs pull, cache keys and TTLs, what to put at the edge, edge functions and KV, and the pitfalls of caching dynamic content.
- Caching: Where to cache (browser, CDN, application, database), cache-aside vs write-through vs write-back, eviction policies, invalidation, hot keys, thundering herds and cache stampedes.
- Load balancing: Layer 4 vs layer 7 load balancers, routing algorithms, health checks, sticky sessions, global load balancing and how to talk about them in a system design interview.