Netflix: Global Streaming, the API Edge and Regional Evacuation
This page is one interview loop in three rounds, built on a real company's public engineering history. All three rounds design the same system. Each round is an era of Netflix's history that answered a crisis, and each opens with the interviewer raising the scope, so the design from the round before has to evolve.
| Round 1: Mid-level | Round 2: Senior | Round 3: Architect | |
|---|---|---|---|
| Era | About 2008–2012: from our own data center to AWS | About 2012–2018: hundreds of microservices behind one front door | About 2013–today: our own CDN, and surviving the loss of a region |
| Level (Amazon) | SDE II (L5) | Senior SDE (L6) | Principal (L7) |
| Members (published) | 9.4M subscribers at the end of 2008, DVD and streaming together | More than 83M (2016); 125M (2018) | More than 325M paid memberships (end of 2025) |
| Traffic we plan for | 1 billion API requests a day (Netflix's 2011 figure; it reported about 20,000/s at peak), planned at about 35,000/s at the evening peak (assumption) | 500,000 edge requests/s per region at peak (assumption); EVCache at 30M operations/s worldwide (published, 2016) | 3M edge requests/s worldwide (assumption); about 100 Tbps of video at the global peak (our estimate) |
| Footprint | One AWS region, 3 AZs, third-party CDNs for video | One region per geography, 3 AZs; Zuul, Eureka, EVCache | Three active-active AWS regions; Open Connect appliances; partnerships with more than a thousand ISPs |
| Target | No single point of failure; Netflix's internal goal is 99.99% | Contain failures; many deploys a day | Move all traffic out of a failed region in under 10 minutes (Netflix: 8 minutes, 2018) |
| Reading time | ~35 min | ~40 min | ~45 min |
You can start at any round. Rounds 2 and 3 open with a "Where we left off" summary that catches you up.
How to read a case study. Every claim about what Netflix actually did comes from Netflix's own engineering blog, talks, open-source documentation or financial filings, and each round ends with a Sources list. We mark those claims Netflix published, What Netflix built (published) or (published). Where Netflix hasn't published the details, we say so and show a design that fits, labeled as ours. Numbers marked assumption are ours, chosen to make the arithmetic concrete. They are not Netflix's internal figures.
Loop Opener: Why Netflix?
Two Systems: a Small Brain and Huge Pipes
When you press Play on Netflix, two very different systems go to work.
- The brain (the control plane). A handful of API calls: who are you, what can you watch, where did you stop, and which server should send you the video? These requests are small, a few kilobytes each, but every one of them must work.
- The pipes (the data plane). Then the video itself arrives: gigabytes an hour, from a server as close to your home as possible.
Netflix published this split plainly: its control plane runs on AWS, and its video is delivered by Open Connect, Netflix's own content delivery network, whose servers sit inside internet providers' networks and at internet exchange points.
Synthesizing vector architecture diagram...
The thin arrows carry kilobytes and must never fail; the thick arrow carries almost all the bytes and must be close to the viewer. The two planes fail, scale and cost money in completely different ways.
A few words we'll use all page:
| Word | What it means on this page |
|---|---|
| Control plane | The services that decide: sign-in, the catalog, recommendations, playback decisions. Small requests, high value. |
| Data plane | The system that moves the video bytes. Huge volume, simple requests. |
| CDN | Content delivery network: many cache servers near viewers that serve files so the origin doesn't have to. |
| Manifest | The small document a player gets before playing: the list of video and audio files, their bitrates, and the URLs to fetch them from. |
| Region / AZ | An AWS region is a geographic area (for example us-east-1 in Virginia); an Availability Zone is one or more separate data centers inside it. |
| Edge gateway | The front door: the first server of ours a device's API request reaches. Netflix's is called Zuul. |
| Evacuation | Moving all members' traffic out of one AWS region into the others. |
What Makes It Hard
- The brain must never be down. If the control plane is down, nobody can press Play, even though the video servers are fine.
- The pipes are enormous. Netflix gave as its reason for building its own CDN that it had grown to be "a significant portion of overall traffic" on consumer ISP networks. Evening is when everyone watches, and evening is when the internet is busiest.
- Everything changes all the time. Hundreds of services, hundreds of deploys a day and more than a thousand device types mean something is always broken somewhere.
The Question the Whole Loop Answers
How do we keep the control plane always available, and push the video bits as close to the viewer as possible?
The answer grows every round:
- Round 1: leave the single data center, make every service stateless and spread it across Availability Zones, keep member data in Cassandra, and let CDNs carry the video.
- Round 2: put one programmable front door (Zuul) in front of hundreds of services, find instances through a registry (Eureka), contain slow dependencies with circuit breakers and concurrency limits, cache hot data in every zone (EVCache), and break things on purpose (chaos engineering).
- Round 3: build our own CDN inside ISPs, filled overnight; run three AWS regions active-active; and practice moving all traffic out of a region in minutes.
Round 1 · Mid-level · "Era 1: From a Data Center to the Cloud"
~35 min · SDE II (L5) · 1 AWS region, 3 AZs · 1B API requests/day, ~35,000/s planned for the evening peak (assumption; Netflix reported ~20,000/s at its 2011 peak) · video through third-party CDNs · no single point of failure · 99.99% (Netflix's internal goal)
R1.1 Establish Design Scope
The interviewer sets the scene: "It's 2008. We rent DVDs by mail, and a year ago we started streaming. Everything runs as one big application on one relational database in our own data center. In August the database got corrupted. Design where we go from here."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| What does the streaming service do? | Members browse the catalog, press Play, and resume where they stopped, on the web and a growing list of devices. | Four operations: sign in, browse, get playback info, save position (R1.2). |
| Where does the video come from today? | Third-party CDNs. Our servers never send video. | Our API only has to tell players where to fetch; it never carries video (step 1.2). |
| What broke, exactly? | Netflix published it: in August 2008 "a major database corruption" meant that for three days Netflix could not ship DVDs to members. | The single database is a single point of failure for the whole business (step 1.1). |
| How big are we? | At the end of 2008 Netflix reported about 9.4M subscribers, DVD and streaming together. Streaming is growing fast and we can't predict it. | We need capacity we can add in minutes, not months (step 1.4). |
| How much API traffic? | Plan for 1 billion API requests a day. (That's the figure Netflix published for its API at the end of 2011.) | 11,574 requests/s on average. Netflix reported about 20,000/s at its 2011 peak; we plan for 3× the average, about 35,000/s, to leave room for growth (assumption) (R1.7). |
| What's the goal? | No single point of failure. Netflix's stated internal goal is 99.99% availability. | Every tier must survive losing a server, and then a whole AZ. |
Out of scope for this round: device-specific APIs, our own CDN, and running in more than one region.
What 99.99% allows. A 30-day month has 43,200 minutes; 0.01% of that is 4.32 minutes of downtime a month.
R1.2 Functional Requirements, Derived Step by Step
| Phrase from the problem | Operation |
|---|---|
| "Members sign in" | login(credentials) returns a session token |
| "Browse the catalog" | getHome(profile) returns rows of titles; getTitle(id) returns details |
| "Press Play" | getPlayback(profile, title, device) returns a manifest: files, bitrates, CDN URLs, DRM license info |
| "Resume where you stopped" | saveBookmark(profile, title, position) every so often while playing; getPlayback returns the last position |
Not yet: a different API shape for each device family, our own CDN, several regions.
R1.3 Non-Functional Requirements: the Questions
We name each quality first; the numbers come in R1.7.
- Availability first. Can members press Play if one server, one database node or one data center fails? Today, the answer for the database is no.
- Scale for evenings. Traffic follows the evening: many times more people watch at 9 p.m. than at 5 a.m. Can we add capacity for the peak and give it back afterward?
- Keep video off the API. Video bytes and API calls must not share servers, links or budgets.
- Write availability for bookmarks. Position updates arrive constantly while people watch. If we can't save one, the member resumes at the wrong place tomorrow.
R1.4 The API
This API is illustrative: it's a design that fits the problem, not Netflix's actual interface.
Get a playback manifest
httpPOST /v1/playback/manifest HTTP/1.1 Host: api.example-streaming.com Authorization: Bearer <session token> Content-Type: application/json { "profile_id": "p_3019", "title_id": "t_81049281", "device": { "type": "tv", "model": "brandX-2008", "drm": ["playready"] } }
httpHTTP/1.1 200 OK Content-Type: application/json { "playback_id": "pb_01HZX894", "resume_position_s": 1452, "streams": [ { "bitrate_kbps": 1000, "url": "https://cdn-a.example.net/t_81049281/1000.ismv?exp=1223476800&sig=..." }, { "bitrate_kbps": 2200, "url": "https://cdn-a.example.net/t_81049281/2200.ismv?exp=1223476800&sig=..." } ], "fallback_cdn_base": "https://cdn-b.example.net/t_81049281/", "license": { "url": "https://license.example-streaming.com/playready", "challenge_token": "..." } }
Save the position
httpPUT /v1/bookmarks/p_3019/t_81049281 HTTP/1.1 Authorization: Bearer <session token> Content-Type: application/json { "position_s": 1512, "playback_id": "pb_01HZX894", "client_time": "2008-10-08T20:41:07Z" }
httpHTTP/1.1 204 No Content
- The manifest points to CDN URLs, signed with an expiry, so only this member can fetch the files and only for a while. The API never sends a video byte.
- The bookmark is a
PUT: sending the same position twice leaves the same result, so the player can retry freely.
Status codes
| Code | Meaning |
|---|---|
200 OK / 204 No Content | Manifest returned / position saved |
401 Unauthorized | Session expired; sign in again |
403 Forbidden | This title isn't available to this member or country |
429 Too Many Requests | The device is calling far too often |
503 Service Unavailable | We're shedding load; retry after the Retry-After delay |
Recap
- Four operations; the API returns where video lives, never the video.
- Playback and bookmark calls are safe to retry.
R1.5 Design Evolution: From One Database to Many Machines
Each step is a problem, your turn to think, the answer, and what it costs us.
Step 1.0: The Baseline
One Java application behind a hardware load balancer, on a few large servers in our own data center, with one large relational database for everything: members, the catalog, DVD queues, viewing history. Video comes from third-party CDNs.
Synthesizing vector architecture diagram...
Every request, for every feature, ends at one database. When it's corrupted, the whole business stops.
It works until it doesn't, and in August 2008 it didn't.
Step 1.1: The Database Got Corrupted and the Business Stopped
The problem: the one database is corrupted. For three days we can't ship DVDs. Streaming is small today but growing fast, and it will depend on the same kind of database. What would you do?
Step 1.2: Should Video Go Through Our Servers?
The problem: a new engineer proposes serving video from our new AWS API servers too: "one system, one set of logs, and we control everything". What would you do?
Primitive: Distributed Cache Patterns and Eviction · Loop: Design a Video Streaming Platform (transcoding, segments and adaptive bitrate, which this page doesn't re-teach)
Step 1.3: Viewing History Grows Forever and Must Always Be Writable
The problem: every player saves its position every minute while playing, and every title a member watches becomes a history record. In the relational database this is the biggest, fastest-growing table, and a write that fails means a member resumes at the wrong place. What would you do?
Primitive: Database Sharding and Partition Keys · Loop: Design a Distributed Key-Value Store (quorums and last-writer-wins, from the inside)
Step 1.4: Evenings Bring Three Times the Traffic
The problem: requests at the evening peak are about three times the daily average, and about nine times the early-morning trough (our assumptions). What would you do?
Round 1 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.0 | (baseline) | One application, one relational database, one data center | A single point of failure for the business |
| 1.1 | Database corruption stopped the business | Move to AWS; stateless services in 3 AZs | A seven-year migration |
| 1.2 | Should video go through our servers? | Control plane on AWS; data plane on CDNs; signed manifest URLs | Dependence on third-party CDNs |
| 1.3 | History grows and must stay writable | Cassandra, RF 3, LOCAL_QUORUM, keyed by profile | Eventual consistency; wide rows |
| 1.4 | Evenings bring 3× traffic | Auto Scaling, scheduled pre-scaling, catalog caches | Warm-up time; staleness |
R1.6 Architecture v1
Synthesizing vector architecture diagram...
Read it as two planes. The control plane is stateless servers in three AZs with state in a replicated database; the data plane is somebody else's CDN. Nothing in the picture exists only once.
Sources for this round
- Izrailevsky, Vlaovic and Meshenberg, Completing the Netflix Cloud Migration, Netflix, February 2016: the August 2008 corruption and three days without DVD shipments; the move away from vertically scaled single points of failure; completion in January 2016.
- Ciancutti, Four Reasons We Choose Amazon's Cloud as Our Computing Platform, Netflix TechBlog, December 2010.
- Netflix, Q4 2008 results press release, January 2009: about 9,390,000 subscribers.
- Schmaus, Making the Netflix API More Resilient, Netflix TechBlog, December 2011: "a billion requests a day", and around 20,000 requests/s at peak.
- Duvedi, Li, Garg and Fisher-Ogden, Scaling Time Series Data Storage, Part I, Netflix TechBlog, January 2018: Cassandra for viewing history, the 9:1 write-to-read ratio, the EVCache layer, the redesign.
- Netflix, How Netflix Works With ISPs Around the Globe, March 2016: streaming launched in 2007; third-party CDNs before Open Connect.
- Meshenberg, Gopalani and Kosewski, Active-Active for Multi-Regional Resiliency, Netflix TechBlog, December 2013: the 99.99% internal availability goal.
- Gonigberg et al., Open Sourcing Zuul 2, Netflix TechBlog, May 2018: slow warm-up of new instances.
Everything else in this round (instance counts, bitrates, the table schemas, the API shapes) is our design or our assumption.
R1.7 Numbers
Every figure here is an assumption except the ones marked published. The point is the arithmetic, not the exact values.
API traffic
| Quantity | Arithmetic | Result |
|---|---|---|
| Average requests | 1,000,000,000 a day (published, 2011) ÷ 86,400 s | 11,574/s |
| Evening peak | 3 × average (assumption; Netflix reported about 20,000/s at its 2011 peak, and we leave room for growth) = 34,722 | ≈ 35,000/s |
| Early-morning trough | ⅓ × average (assumption) = 3,858 | ≈ 3,900/s |
API fleet. We assume one instance handles 500 requests/s at our target CPU.
| Arithmetic | Instances | |
|---|---|---|
| Needed at the peak | 35,000 ÷ 500 | 70 |
| Survive losing one of 3 AZs | The 2 remaining AZs must hold 70, so each AZ holds 70 ÷ 2 = 35 | 35 per AZ, 105 in all |
| Needed at the trough | 3,900 ÷ 500 = 7.8, round up | 8 |
| Trough, surviving an AZ | 8 ÷ 2 = 4 per AZ | 4 per AZ, 12 in all |
Autoscaling moves the fleet between 12 and 105, about 9 to 1.
Bookmarks. We assume 300,000 people watch at the evening peak, and each player saves its position every 60 seconds.
- Writes: 300,000 ÷ 60 s = 5,000 bookmark writes/s, about 14% of the peak API traffic (5,000 ÷ 35,000).
- With 3 copies each, Cassandra performs 5,000 × 3 = 15,000 replica writes/s. On 6 nodes (2 per AZ), that's 2,500 writes/s per node, well within what one Cassandra node handles (a comfortable budget of 5,000/s per node, assumption).
- Reads: at the published ratio of 9 writes per read, about 5,000 ÷ 9 ≈ 560 history reads/s, and most hit the cache.
History storage. Assume 10M members, 1,000 records each over the years, 200 bytes a record: 10,000,000 × 1,000 × 200 B = 2 TB. Three copies make 6 TB, 1 TB per node on 6 nodes.
Bytes: why the API must not carry video
| Arithmetic | Result | |
|---|---|---|
| Video at the peak | 300,000 streams × 2 Mbps (assumption) | 600 Gbps |
| API at the peak | 35,000 req/s × 10 KB (assumption) = 350 MB/s × 8 | 2.8 Gbps |
| Ratio | 600 ÷ 2.8 | ≈ 214× |
Video is more than 99% of the bytes (600 ÷ 602.8 = 99.5%). That's the whole argument for step 1.2 in one line.
R1.8 Trade-Offs
Our own data center vs AWS
| Build more data centers | Move to AWS (chosen) | |
|---|---|---|
| Capacity | Buy and install months ahead of a growth we can't predict | Add or remove in minutes |
| Failure model | Few, big, reliable machines | Many, small machines that can vanish; the software must expect it |
| Engineering focus | Power, cooling, hardware, networking | The product |
| Cost | Cheaper per server at steady load | Pay for flexibility; buy commitments later for the steady base |
Netflix's own summary of the outcome in 2016: "the cloud allows one to build highly reliable services out of fundamentally unreliable but redundant components".
Relational database vs Cassandra for viewing history
| Relational (one primary) | Cassandra (chosen) | |
|---|---|---|
| Writes | One primary takes every write; failover pauses writes | Any node can coordinate a write; losing an AZ keeps writes flowing |
| Scale | Bigger box, then manual sharding | Add nodes; data spreads by partition key |
| Queries | Any join or filter | Only what the keys allow; we design tables per screen |
| Consistency | Strong by default | Tunable; we choose quorum writes and accept eventual reads |
Third-party CDNs vs our own CDN (not yet)
| Third-party CDNs (now) | Our own CDN (Round 3) | |
|---|---|---|
| Time to start | Days | Years: hardware, software, ISP partnerships |
| Control | Their placement, their caching rules | We decide what is where, and when |
| Cost | Per GB, forever | Up-front hardware; our traffic becomes an efficiency project |
At 2008's scale, buying is right. Round 3 revisits it when Netflix is a significant share of ISP traffic.
R1.9 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| An API instance dies | A few failed requests | The load balancer's health check removes it; the Auto Scaling group replaces it; devices retry the safe GET and PUT calls. |
| An AZ is lost | A third of API instances and one Cassandra copy of every row disappear | The remaining two AZs were sized for the full load (35 each at the peak). Cassandra still has 2 of 3 copies, so LOCAL_QUORUM reads and writes succeed. |
| The memcached cluster fails | Every read goes to Cassandra | Cassandra sees about 560 history reads/s plus catalog misses, which it can absorb at this scale; we alarm on the miss rate. At Round 2's scale this becomes dangerous (R2.8). |
| A Cassandra node is slow | Some writes slow down | With LOCAL_QUORUM, a write needs 2 of 3 copies, so one slow copy doesn't slow the write; the missed write reaches it later by hinted handoff (the coordinator stores the write and replays it when the node returns, for a limited window) or repair. |
| One CDN degrades in one city | Buffering for some viewers | Players switch to the fallback CDN URL in their manifest; we lower that CDN's share in new manifests. |
R1.10 Pillar Check
| Pillar | What Round 1 covers |
|---|---|
| Reliability | No single point of failure: stateless services in 3 AZs, Cassandra with 3 copies and quorum writes; retries on idempotent calls REL 10 · REL 11 |
| Performance Efficiency | Video from CDNs near viewers; catalog cached in memory; each screen reads one partition PERF 3 · PERF 4 |
| Security | Session tokens on every API call; signed, expiring CDN URLs; TLS for the API SEC 9 |
| Cost Optimization | Autoscaling between 12 and 105 instances instead of 105 all day; no video bytes through our servers COST 9 · COST 8 |
| Operational Excellence | Skipped this round: health checks, alarms on cache miss rate and on error rates. |
| Sustainability | Skipped this round: capacity follows demand rather than sitting idle at night. |
R1.11 Round 1 Rubric and Follow-Ups
What a strong mid-level (L5) answer shows
- Names the single point of failure and removes it with replication and statelessness, not with better backups.
- Separates the control plane from the data plane, and backs it with a byte count.
- Picks a store for viewing history from the access pattern: write-heavy, keyed by member, availability over consistency.
- Designs Cassandra tables per screen, with the partition key the app already has.
- Sizes the fleet for losing an AZ, not just for the peak.
- Autoscales, and knows new instances need warm-up.
Follow-up questions
-
"Why not keep the relational database and add read replicas?" Answer: replicas scale reads, but viewing history is 9 writes for every read, and every write still goes to one primary. A primary failure still stops writes until failover finishes.
-
"Two devices on the same profile save different positions at the same moment. What happens?" Answer: Cassandra keeps the write with the later timestamp. The member resumes at one of the two positions, which is acceptable for a bookmark. It would not be acceptable for money, which is why we don't keep money this way.
-
"Why sign the CDN URLs if the video is encrypted anyway?" Answer: encryption (DRM) stops people from watching without a license; signing stops strangers from using our CDN bill to download files at all. Different jobs.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "Better backups fix it" | A restore takes days and the single database remains. |
| "Serve video from the API servers" | Video is about 214 times the API's bytes; it would starve the control plane. |
| "One relational table for history" | One primary takes every write and pauses them on failover. |
| "Size for the peak all day" | The fleet would sit about 90% idle every early morning. |
| "Survive the peak" | We must survive the peak with an AZ missing. |
Round 2 · Senior · "Era 2: Hundreds of Microservices Behind One Front Door"
~40 min · Senior SDE (L6) · one region per geography, 3 AZs · 500,000 edge requests/s per region at the peak (assumption) · 10M device connections per region (assumption) · hundreds of services and deploys a day · contain every failure
R2.0 Where We Left Off
This is what the candidate says aloud in the first 60 seconds of Round 2. If you're starting here, it's everything you need from Round 1.
Round 1 in 60 seconds. "In August 2008 Netflix's one relational database was corrupted and DVDs didn't ship for three days. Netflix decided to leave vertically scaled single points of failure and moved to AWS, a migration that finished in January 2016. We made every service stateless and spread it across three Availability Zones, sized so two AZs carry our planned evening peak of 35,000 API requests a second, 35 instances per AZ, autoscaling down to 12 instances at night. Member data like bookmarks and viewing history lives in Cassandra, keyed by profile, with three copies (one per AZ) and quorum writes, because history is written nine times for every read and availability matters more than instant consistency. Video never touches our servers: the API returns a manifest with signed URLs on third-party CDNs, because video is about 214 times the API's bytes. Open costs: every device calls one generic API, services find each other by configuration, a slow dependency can stall everything, one cache cluster, and we learn about failures from members."
Architecture v1, compact
Synthesizing vector architecture diagram...
Round 1 in one picture: nothing exists once, and video goes around our servers.
Round 1 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.1 | Database corruption stopped the business | AWS, stateless services, 3 AZs | Years of migration |
| 1.2 | Video through our servers? | Control plane vs data plane; signed CDN URLs | CDN dependence |
| 1.3 | History must stay writable | Cassandra, RF 3, quorum writes | Eventual consistency |
| 1.4 | 3× evening traffic | Autoscaling, caches | Warm-up time |
Open costs: one generic API for every device, hand-configured service addresses, no isolation between dependencies, one cache cluster, failures found by members.
R2.1 The Scope Raise
Interviewer: "It's a few years later. The monolith has become hundreds of services owned by different teams. We support well over a thousand device types, and each wants different data in a different shape. Last week a service nobody considers critical slowed down and took the home page with it. Teams deploy many times a day. And we want devices to keep a connection open, so we can push to them."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| How many services? | Hundreds of microservices, each owned by one team. (Netflix, 2016.) | Nobody can hand-maintain the addresses between them (step 2.2). |
| How many device types? | More than 800 in mid-2012; "over a thousand" by the end of that year. (Netflix.) | One generic API fits none of them well; each device family gets its own endpoints (step 2.1). |
| How much traffic reaches one region? | Plan for 500,000 requests a second at the edge at the peak. | About 150 gateway instances per region (R2.6). |
| What exactly happened last week? | The ratings service slowed from 30 ms to 5 seconds, and the home page, which calls it, stopped loading. | Every dependency is isolated and has a fallback (step 2.3). |
| How often do we deploy? | "Hundreds of deploys every day." (Netflix, 2015.) | Change is our main cause of failure: routing must support canaries, and failures must be contained. |
| Persistent connections? | Yes. We want millions of devices connected at once. | Idle connections must be cheap: thread-per-connection won't do (step 2.1). |
| Who finds failures? | We must find them before members do. | Break things on purpose, in production, safely (step 2.5). |
Scope change
| Round 1 | Round 2 | |
|---|---|---|
| Services | A few, one API | Hundreds, many teams |
| Devices | Browsers and a few devices | More than a thousand types |
| Edge traffic per region | ~35,000 requests/s | 500,000 requests/s |
| Connections | Short requests | Millions of long-lived connections |
| Deploys | Occasional | Hundreds a day |
| Availability target | 99.99% | 99.99%, with any single dependency failing |
The "Not yet" list from R1.2 comes back: device-specific APIs are in scope now.
R2.2 What Breaks in the Round 1 Design
| Round 1 piece | What breaks at the new scale |
|---|---|
| One generic API | Netflix described the old model: a device like the PS3 made many requests to a one-size-fits-all REST API to start one screen, each returning fields it didn't need. More device types make it worse. |
| Addresses in configuration | Autoscaling and deploys replace instances all day. A config file of IP addresses is wrong within minutes. |
| A load balancer in front of every service | Every call between services crosses one more box, which is one more thing to fail and one more per-GB bill (R2.6). |
| One shared thread pool per server | A slow dependency holds threads until none are left for fast ones (step 2.3). |
| One cache cluster | Readers in two of three AZs cross an AZ boundary on every read, and losing that AZ loses the cache. |
| Thread-per-connection servers | 10M connected devices across 150 instances is 66,667 connections per instance: 66,667 threads (R2.6). |
| Testing before release | Hundreds of deploys a day change the system faster than any test environment can copy it. |
R2.3 New Requirements and API Additions
A front door for all device traffic. Every device request enters through one gateway fleet, which authenticates it, applies routing rules, and sends it to the right service. Netflix published the gateway's structure, which the diagram shows: inbound filters run before the request is proxied (authentication, routing, decorating the request); an endpoint filter either answers directly or proxies to the backend ("origin" in Zuul's terms); outbound filters run on the response (compression, metrics, headers).
Synthesizing vector architecture diagram...
Netty handlers do the network work on both sides; filters hold all the logic. Change the filters and the same gateway solves a different problem, which is why Netflix could deploy the same core at the edge and, with fewer filters, for internal traffic.
Device-specific endpoints. In 2012 Netflix moved from a one-size-fits-all REST API to "a platform for API development" where each UI team writes endpoints for its own devices. An illustrative endpoint (our shapes):
httpGET /tv/v2/home?profile_id=p_3019&rows=12 HTTP/1.1 Host: api.example-streaming.com Authorization: Bearer <device token>
json{ "rows": [ { "id": "continue_watching", "items": [ { "title_id": "t_81049281", "box_art": "https://img.example.net/t_81049281/tv_1080.jpg", "resume_position_s": 1512 } ] }, { "id": "top10_country", "items": [ { "title_id": "t_80100172", "box_art": "https://img.example.net/t_80100172/tv_1080.jpg" } ] } ], "degraded": ["ratings"] }
One request builds the whole TV home screen with TV-sized artwork; a phone's endpoint returns different fields. degraded tells the UI that ratings came from a fallback (step 2.3).
Registration and discovery. Every instance registers itself in a registry and renews a lease; callers look services up by name. Eureka's REST interface looks like this (simplified):
httpPOST /eureka/v2/apps/RATINGS HTTP/1.1 Content-Type: application/json { "instance": { "instanceId": "i-0a1b2c3d", "app": "RATINGS", "ipAddr": "10.1.2.3", "vipAddress": "ratings", "port": { "$": 8080, "@enabled": "true" }, "status": "UP" } }
httpPUT /eureka/v2/apps/RATINGS/i-0a1b2c3d HTTP/1.1
The first registers; the second is the heartbeat, sent every 30 seconds.
Every dependency declares its fallback (our configuration shape):
yamldependencies: ratings: timeout_ms: 300 concurrency_limit: adaptive fallback: fail_silent # home page shows no stars subscriber: timeout_ms: 200 concurrency_limit: adaptive fallback: custom # use the member data in the signed request identity playback_license: timeout_ms: 800 concurrency_limit: adaptive fallback: fail_fast # no safe default: return 503 fast
R2.4 Design Evolution: A Front Door, a Registry and Containment
Step 2.1: Devices Talk to Dozens of Services
The problem: a TV starting up calls the member service, the catalog, recommendations, ratings, artwork and more, each through its own public load balancer. Each device type needs different fields. We want to add authentication once, route 1% of traffic to a canary, and keep millions of devices connected for push. What would you do?
Primitive: API Gateway and Reverse Proxy
Step 2.2: Service Addresses Change Constantly
The problem: the ratings service runs on 60 instances today, 110 tonight, and a completely new set after this afternoon's deploy. How does the API layer know where to send requests? What would you do?
Primitive: Gossip Protocol and Failure Detection (heartbeats, timeouts and why a registry must not panic)
Step 2.3: One Slow Service Took Down the Home Page
The problem: ratings is optional: the home page is fine without stars. But when ratings slowed from 30 ms to 5 seconds, the home page stopped loading for everyone. What would you do?
Primitive: Circuit Breaker, Bulkhead and Fault Tolerance · Drill: The slow recommendation service that took down checkout · Loop: Distributed rate limiter, step 3.4 (load shedding vs rate limiting)
Step 2.4: Hot Data Must Be Fast in Every AZ
The problem: services read member data (viewing history, recommendations, profile settings) millions of times a second. Cassandra can't serve that at sub-millisecond latency, and our one memcached cluster lives mostly in one AZ. What would you do?
Primitive: Distributed Cache Patterns and Eviction
Step 2.5: Failures Surprise Us in Production
The problem: we have fallbacks, breakers and three AZs, on paper. But last month an instance failure took down a service whose fallback had a bug nobody knew about, because the fallback had never run. What would you do?
Round 2 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | Devices talk to dozens of services | Zuul edge gateway with filters; device-specific API; Zuul 2 on Netty; edge authentication and Passport | A shared critical component; harder async debugging |
| 2.2 | Addresses change constantly | Eureka registry, 30 s heartbeats, a documented 90 s lease; client-side load balancing | ~3–5 minutes of stale entries after a crash |
| 2.3 | A slow optional service took the page down | Bulkheads, circuit breakers, fallbacks (Hystrix, now in maintenance); adaptive concurrency limits | Degraded pages; fallbacks to maintain |
| 2.4 | Hot data fast in every AZ | EVCache: a copy per AZ, local reads, cache-aside, TTLs | Staleness; 3× RAM; cross-AZ writes |
| 2.5 | Failures surprise us | Chaos Monkey, chaos experiments, ChAP | Bounded, chosen risk |
R2.5 Architecture v2
Synthesizing vector architecture diagram...
Follow a request left to right: one front door, a device-shaped API, then direct calls between services, each wrapped in a breaker and a limit, over addresses from Eureka. Dotted lines are lookups, misses and injected failures, not the request path. The load balancer at the front is our AWS mapping; Netflix fronted Zuul with AWS load balancers, and the exact type is not something we rely on.
Sources for this round
- Meshenberg, Gopalani and Kosewski, Active-Active for Multi-Regional Resiliency, Netflix TechBlog, December 2013: Zuul opened to the community in June 2013.
- Jacobson, Embracing the Differences: Inside the Netflix API Redesign, Netflix TechBlog, July 2012: 800+ device types; device-specific endpoints.
- Cockcroft, A Closer Look at the Christmas Eve Outage, Netflix TechBlog, December 2012: "over a thousand different streaming devices".
- The Cloud Gateway team, Zuul 2: The Netflix Journey to Asynchronous, Non-Blocking Systems, Netflix TechBlog, September 2016: Netty, connection scaling, the measured results, 83M members.
- Gonigberg et al., Open Sourcing Zuul 2, Netflix TechBlog, May 2018: 80+ clusters, ~100 origins, more than 1M requests/s, filter types, self-service routing, load balancing features.
- Netflix/zuul on GitHub (checked September 2026): active, Netty-based
zuul-core. - Netflix TechBlog, Edge Authentication and Token-Agnostic Identity Propagation, February 2021: authentication in Zuul filters, the Passport, the move to NodeQuark backends for frontends.
- Eureka wiki, Understanding Eureka client/server communication and Server self-preservation mode: 30 s heartbeats, 90 s eviction, 30 s delta fetches, the 15% threshold.
- Vroom, Mulcahy, Yuan and Gulewich, Zero Configuration Service Mesh with On-Demand Cluster Discovery, Netflix TechBlog, August 2023: why client-side load balancing, Eureka's degraded mode, the move to Envoy.
- Schmaus, Making the Netflix API More Resilient, December 2011: breakers, bounded thread pools, 10 s windows, three fallback kinds.
- Netflix/Hystrix README: maintenance mode; resilience4j and adaptive limits.
- Landau, Thurston and Bozarth, Performance Under Load, Netflix TechBlog, March 2018: adaptive concurrency limits.
- The EVCache team, Caching for a Global Netflix, Netflix TechBlog, March 2016: EVCache scale, per-region redundancy, "hundreds of microservices".
- Netflix/EVCache wiki: reads from the client's own zone, writes to all zones, zone fallback.
- The Netflix Simian Army, July 2011; Basiri, Hochstein, Thosar and Rosenthal, Chaos Engineering Upgraded, September 2015; ChAP: Chaos Automation Platform, July 2017.
The fleet sizes, request rates, cache sizes, the fallback configuration and the AWS load balancer mapping are ours.
R2.6 Numbers and Cost
All figures are assumptions unless marked published.
The gateway fleet for one region. Assume 500,000 edge requests/s at the peak and 5,000 requests/s per Zuul instance at our target CPU.
| Arithmetic | Instances | |
|---|---|---|
| Needed at the peak | 500,000 ÷ 5,000 | 100 |
| Survive losing an AZ | 100 ÷ 2 = 50 per AZ | 50 per AZ, 150 in all |
| Requests per instance, all AZs up | 500,000 ÷ 150 | 3,333/s |
For scale: Netflix's published total in 2018 was more than 1 million requests/s across all its Zuul clusters.
Connections: why Zuul 2. Assume 10M devices hold a connection to this region at the peak.
- Per instance: 10,000,000 ÷ 150 = 66,667 connections.
- Thread-per-connection (Zuul 1's model) would need 66,667 threads per instance. The JVM's default stack reservation on 64-bit Linux is 1 MB a thread, so that's about 65 GB of reserved stack address space before any work, plus the scheduler cost of switching among them.
- An event loop per core (Zuul 2's model): on a 16-core instance, 16 threads serve all 66,667 connections, about 4,167 each; each idle connection costs a socket and a small buffer.
Little's law in the gateway. At 3,333 requests/s per instance and a 50 ms origin, 3,333 × 0.05 = 167 requests are in flight. If an origin slows to 2 s, the same arrivals mean 3,333 × 2 = 6,667 in flight. Zuul 1 needs a thread for each; Zuul 2 needs only memory for each, which is why Zuul 2 still needs concurrency limits per origin (Netflix open-sourced "origin concurrency protection" with Zuul 2): memory runs out too, just later.
EVCache for one region. Assume 10M cache operations/s at the peak (9M reads, 1M writes), 1 KB average item, and a 30 TB hot set. Three regions at this rate make 30M/s, the same order as Netflix's published 2016 peak.
| Quantity | Arithmetic | Result |
|---|---|---|
| RAM, one full copy per AZ | 30 TB × 3 | 90 TB |
| Nodes at 100 GB usable each (assumption) | 30 TB ÷ 100 GB = 300 per AZ | 300 per AZ, 900 in all |
| Cross-AZ write bytes | 1M writes/s × 1 KB × 2 other AZs | 2 GB/s |
| Their cost | 2 GB/s × 3,600 s = 7,200 GB/h × 730 h = 5,256,000 GB × $0.02/GB | ≈ $105,000 a month |
| If reads ignored AZs | 9M/s × 1 KB × ⅔ cross an AZ = 6 GB/s = 21,600 GB/h × 730 h × $0.02 | ≈ $315,000 a month |
Data moving between AZs in one region is charged $0.01/GB in each direction, so $0.02 for every GB that crosses. Reading from our own AZ's copy is worth about $315,000 a month here; the cost of writing every copy is the price of having one.
The load balancer in front of Zuul (our AWS mapping). A Network Load Balancer with a TCP listener, passing TLS through to Zuul. NLB bills the largest of three dimensions in NLB capacity units (NLCUs) at $0.006 per NLCU-hour in us-east-1. For TCP, one NLCU is 800 new connections/s, 100,000 active connections, or 1 GB processed an hour. Assume 20,000 new connections/s and 4 KB per request and response together.
| Dimension | Arithmetic | NLCUs |
|---|---|---|
| New connections | 20,000 ÷ 800 | 25 |
| Active connections | 10,000,000 ÷ 100,000 | 100 |
| Processed bytes | 500,000/s × 4 KB = 2 GB/s = 7,200 GB/h ÷ 1 GB | 7,200 |
Bytes dominate: 7,200 × $0.006 = $43.20 an hour, × 730 hours ≈ $31,500 a month, plus $0.0225 an hour for the load balancer itself (about $16).
A second copy of the same bytes. If we also put an internal Application Load Balancer between Zuul and the device APIs, it processes the same 2 GB/s again: 7,200 LCUs (1 GB an hour per LCU for EC2 targets) × $0.008 ≈ $57.60 an hour, ≈ $42,000 a month more, plus one more box on every request. Client-side load balancing through Eureka avoids both, which is the choice Netflix made for traffic between services.
R2.7 Trade-Offs
One programmable gateway vs a service mesh
| Edge gateway (Zuul) | Service mesh (a proxy beside every service) | |
|---|---|---|
| Where the logic runs | In one fleet at the edge | In a sidecar next to every instance |
| Best for | Traffic from devices: auth, routing, canaries, shedding | Traffic between services: retries, limits, mTLS, in every language |
| Failure blast radius | A bad filter affects all device traffic | A bad proxy config affects the services that load it |
| Netflix's path | Zuul at the edge, open-sourced in June 2013, Netty-based since 2016 | Client libraries (Eureka, Ribbon, Hystrix) for service-to-service calls; since 2023, moving those features into Envoy, still using Eureka as the source of truth |
They aren't rivals: Netflix published both, for different traffic. Netflix gave the reason for the mesh: keeping the same features "in more languages, in more clients" consistent was the hard part.
Client-side vs server-side load balancing
| Client-side (chosen) | A load balancer per service | |
|---|---|---|
| Hops | Caller to callee directly | One more box on every call |
| Cost | A registry | LCUs for every byte (R2.6), per service |
| Failure | Stale registry entries (~3–5 min after a crash) | The load balancer itself can fail, as ELBs did on Christmas Eve 2012 |
| Smartness | The caller knows its own latencies and errors, and picks accordingly | Only what the load balancer measures |
The cost of chaos engineering
| Wait for failures | Cause them on purpose (chosen) | |
|---|---|---|
| When failures happen | Whenever, often at night | Business hours, engineers present |
| Member impact | Unbounded | A small, bounded slice, stopped automatically on an error budget |
| What we learn | After an incident | Before one, continuously |
R2.8 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| A blocking call inside an async filter | One core's event loop stalls; every connection on it (about 4,167 here) freezes; latency spikes on one-sixteenth of an instance at a time; health checks fail | Rule: filters never block. I/O filters are async; anything slow or CPU-heavy goes to a separate worker pool. Netflix used an instrumenting tool (Reactive-Audit) to find blocking calls hidden in libraries. Alarm on event-loop lag per instance. |
| Retry storms | An origin fails; devices and services retry; traffic triples on a sick service | Netflix's Zuul tracks error rates per origin and, when a whole service is in trouble, "throttle[s] retries from devices and disable[s] internal retries". Our clients add a retry budget (at most 10% extra requests) and exponential backoff with jitter. |
| The original outlives its retry | A retry succeeds on instance B while the first attempt, still running on slow instance A, also completes: two plays registered, two emails sent | Retry only idempotent calls. A call with side effects carries an idempotency key (for Play: the playback_id from the manifest), and the service claims the key with a conditional insert (insert-if-absent, for example a lightweight transaction at LOCAL_SERIAL or memcached add) before acting. The loser waits for the stored result or gets a 409 "in progress". A claim carries an expiry, so an attempt that crashed can be retried; a later arrival finds the completed claim and returns the first result. |
| A cache stampede on a big release | A new season launches at 00:00 UTC; millions of members open its page; its keys aren't cached; every miss goes to Cassandra at once | A design that fits: warm the new title's metadata and artwork keys before launch; collapse concurrent misses so only one request per key per instance goes to the database while others wait for its result; fill with memcached add (store only if absent), so a refill never overwrites a newer write. |
| A Eureka server outage | Registrations and fetches fail | Clients keep calling with their last copy; entries get staler. New instances can't be discovered until it's back, so deploys pause. |
| A leaked buffer in Zuul | Direct (off-heap) memory grows until the instance dies | Netflix named "ByteBuf leaks, file descriptor leaks, lost responses" as the typical async bugs. Canary every Zuul build on a small slice with leak detection on, and alarm on direct memory use. |
R2.9 Production Gotchas
| Gotcha | Why it hurts | What we do |
|---|---|---|
| Blocking I/O in an async filter | One blocked event loop stalls thousands of connections | Async I/O filters only; CPU-heavy work off the event loop; lint and runtime detection |
| Sticky DNS | Devices and resolvers keep old addresses longer than the TTL says, so removing a gateway address from DNS doesn't empty it quickly | Keep old addresses serving until their traffic drains; Round 3 needs this for evacuation |
| Unbounded cache writes | Caching every computed value (including ones nobody reads again) evicts hot keys and fills RAM | TTL on every item; cache only what's read again; for replicated caches, send an invalidation instead of the data when remote reads are rare, but not for caches on the streaming path, which must be warm after an evacuation (Round 3) |
| A fallback that's never run | It fails the one time it's needed | Chaos experiments exercise it continuously |
ThreadLocal context in async code | Request context leaks between requests sharing a thread | Pass context explicitly with the request (Netflix: "thread local variables don't work in an async non-blocking world") |
R2.10 Pillar Check
| Pillar | What Round 2 adds |
|---|---|
| Reliability | Bulkheads, breakers, fallbacks and adaptive concurrency limits for every dependency; retry budgets; idempotency keys; chaos experiments REL 5 · REL 4 · REL 12 |
| Performance Efficiency | Event-loop gateway for millions of idle connections; device-shaped endpoints; AZ-local cache reads PERF 2 · PERF 4 |
| Security | Authentication at the edge; a per-request, HMAC-protected Passport instead of raw tokens deep in the stack SEC 2 · SEC 9 |
| Cost Optimization | No load balancer between services; AZ-local reads worth about $315,000 a month; bytes, not connections, drive the NLB bill COST 8 · COST 5 |
| Operational Excellence | Canary and sticky-canary routing; per-origin error tracking and alerting; chaos experiments in the deploy pipeline OPS 6 · OPS 8 |
| Sustainability | Skipped this round: the async gateway's gain is mainly connection density, which means fewer instances for the same connections. |
R2.11 Round 2 Rubric and Follow-Ups
What a senior (L6) answer adds over L5
- Puts one programmable front door in front of the services and explains its filter model.
- Explains non-blocking honestly: connection scaling yes, CPU savings only when work per request is small.
- Uses a registry with leases and client-side load balancing, and knows how stale a crashed instance can stay (3 to 5 minutes, not the documented 90 s).
- Uses Little's law to show why a slow optional dependency exhausts a shared pool, and contains it with bulkheads, breakers, fallbacks and adaptive limits.
- Knows Hystrix is in maintenance and what replaced its role.
- Designs a per-AZ cache and prices the cross-AZ bytes.
- Proposes chaos engineering as experiments with a steady state and an automatic stop.
Follow-up questions
-
"If Zuul 2 didn't make the API cluster cheaper, why did Netflix do it?" Answer: connection scaling. Persistent connections from every device enable push and fewer chatty requests. CPU-bound work costs the same CPU either way; idle connections are where thread-per-connection loses.
-
"Eureka is down for ten minutes. What breaks?" Answer: calls between running services keep working from each caller's cached copy. New instances can't be found, crashed ones aren't removed, and deploys should pause. That's why the design treats Eureka as important but not on the request path.
-
"Why fail silent for ratings but fail fast for the playback license?" Answer: a page without stars is still a page; a Play without a license can't play. A fake license would be worse than a quick, honest error the device can retry.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "Clients call services directly" | Every service becomes public; every change touches every device. |
| "IP addresses in config" | Stale within minutes of the next deploy or scale-out. |
| "Raise the timeouts" | Slow calls hold threads longer; the shared pool empties faster. |
| "Hystrix is Netflix's current answer" | It has been in maintenance mode since 2018; adaptive limits and the mesh carry that role. |
| "Async makes everything faster" | Netflix measured no gain for CPU-heavy work. |
| "One cache cluster" | Cross-AZ latency and bills, and one AZ's loss empties it. |
| "Test more before release" | Hundreds of deploys a day outrun any test environment. |
Round 3 · Architect · "Era 3: Own the CDN, and Survive Losing a Region"
~45 min · Principal (L7) · 190 countries · more than 325M paid memberships · about 536M hours watched a day · 3 active-active AWS regions · Open Connect partnerships with more than a thousand ISPs · 3M edge requests/s worldwide (assumption) · evacuate a region in under 10 minutes
R3.0 Where We Left Off
Round 2 in 60 seconds. "Hundreds of services now sit behind one programmable front door, Zuul. Its filters authenticate each request and create a Passport, an HMAC-protected identity object that travels to every service; its routing rules run canaries; its load balancer avoids cold, failing and busy instances. Zuul moved to Netty in 2016: the win was cheap persistent connections, not CPU. Behind it, device-specific APIs build each screen in one call. Services find each other through Eureka (30-second heartbeats, a documented 90-second lease that in practice expires in 3 to 5 minutes, cached registries) and call each other directly, with bulkheads, circuit breakers, fallbacks and adaptive concurrency limits; Hystrix is in maintenance mode since 2018. EVCache keeps a copy of hot data in every AZ and reads from the local one. Chaos Monkey, chaos experiments and ChAP exercise our failure paths continuously. Per region: 500,000 requests a second, 150 Zuul instances, 10M connections. Open costs: third-party CDNs carry every video byte, and each region is a world of its own: if a region fails, its members can't play."
Architecture v2, compact
Synthesizing vector architecture diagram...
Round 2 in one picture: one front door, direct calls with containment, a cache in every zone, all inside one region.
Round 2 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | Devices talk to dozens of services | Zuul 2, device APIs, edge auth | Shared critical component |
| 2.2 | Addresses change constantly | Eureka, client-side load balancing | ~3–5 min stale entries |
| 2.3 | A slow optional service | Bulkheads, breakers, fallbacks, adaptive limits | Degraded pages |
| 2.4 | Hot data in every AZ | EVCache, a copy per AZ | Cross-AZ writes |
| 2.5 | Failures surprise us | Chaos engineering, ChAP | Bounded risk |
Open costs: someone else's CDN for every byte, far from viewers; one region per member, with nothing to fail over to.
R3.1 The Scope Raise
Interviewer: "Two things happened. On Christmas Eve 2012, a problem in AWS's load balancer service in
us-east-1stopped playback on many TV devices for hours. And we're about to launch in nearly every country, while in the evening we're already a significant share of every ISP's traffic. Deliver the video yourselves, and make the control plane survive the loss of a whole region."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| What exactly happened on Christmas Eve 2012? | Netflix published it: data in AWS's Elastic Load Balancer service in US-East was deleted by a maintenance process. A handful of Netflix's ELBs, across all AZs, stopped passing traffic; game consoles and other devices were affected for about seven hours. Netflix's service in the UK, Ireland and the Nordic countries was not affected. | A regional service can fail across every AZ at once. Multi-AZ is not enough; we need other regions (step 3.2). |
| How much video? | Netflix reports more than 97 billion hours watched in the first half of 2026. | About 536M hours a day, roughly 100 Tbps at the global peak by our assumptions (R3.6). |
| Where are the viewers? | 190 countries since January 2016; more than 325M paid memberships at the end of 2025. | Content must cross oceans once, not once per viewer (step 3.1). |
| How fast must we leave a failed region? | Under 10 minutes. Netflix published 8 minutes in 2018, down from about 50. | Capacity waiting in the survivors, and fast traffic moves (step 3.3). |
| How stale may member data be after a move? | A bookmark a second or two old is fine. Some writes, like one account per email address, must never be duplicated. | Consistency is chosen per kind of data (step 3.4). |
| How do we know it works? | We practice, in production, regularly. | Region-scale chaos exercises, judged by stream starts per second (step 3.5). |
Scope change
| Round 2 | Round 3 | |
|---|---|---|
| Video delivery | Third-party CDNs | Open Connect, our own CDN, inside ISPs |
| Regions | One per geography | Three, active-active; any region can serve any member |
| Survives | Losing an AZ | Losing a region: traffic out in under 10 minutes |
| Member data | One region | Replicated to every region |
| Countries | A few dozen | 190 |
| Edge traffic | 500,000 requests/s per region | 3M requests/s worldwide, 1M per region |
R3.2 What Breaks in the Round 2 Design
| Round 2 piece | What breaks at the new scale |
|---|---|
| Third-party CDNs | We pay per byte forever, the CDNs' servers mostly sit outside ISP networks so every evening's video crosses peering links, backbones and undersea cables, and we can't decide what is cached where. |
| One region per member | If that region's load balancers, network or control plane fail, those members can't press Play, even though the video servers are fine. |
| Member data in one region | Another region can't serve members whose bookmarks, profiles and history it doesn't have. |
| Caches per region | Even with the data replicated in the database, moving traffic finds every member's cache entries missing: a thundering herd on the database. |
| Calls between regions | Any request that calls another region fails when either region fails, and pays the distance on every call. |
| Capacity that follows the diurnal curve | Netflix: its clusters were "not overprovisioned to the point where they could absorb the additional traffic" of another region. Booting new instances took about 25 minutes. |
| DNS | Repointing DNS takes seconds; devices following it took about 5 minutes (Netflix, 2018). |
R3.3 New Requirements and API Additions
The manifest now points at our own servers. The playback service asks the steering service for appliances that have the files, are healthy, and are close to this client in network terms (our response shape):
json{ "playback_id": "pb_01J2K7Q9", "resume_position_s": 1512, "servers": [ { "rank": 1, "base_url": "https://oca-17.isp-a-sea.example.net/", "site": "embedded in ISP A" }, { "rank": 2, "base_url": "https://oca-04.ix-sea.example.net/", "site": "internet exchange, Seattle" }, { "rank": 3, "base_url": "https://oca-22.ix-pdx.example.net/", "site": "internet exchange, Portland" } ], "files": [ { "stream": "video_3000k", "path": "t_81049281/v3000.mp4" }, { "stream": "audio_en", "path": "t_81049281/a_en.mp4" } ], "url_expires_at": "2026-09-28T23:40:00Z" }
Each appliance reports to the control plane. Netflix published what an appliance reports: its health, the BGP routes it has learned from the ISP router it peers with, and which files it stores. An illustrative report:
json{ "oca_id": "oca-17", "site": "isp-a-sea", "healthy": true, "egress_utilization": 0.62, "bgp_prefixes_learned": ["198.51.100.0/22", "2001:db8:1200::/40"], "files_changed": { "added": 1834, "removed": 912 }, "reported_at": "2026-09-28T19:40:00Z" }
Evacuation controls (ours). One internal request records the decision; automation carries it out, and the call is refused unless its checks pass:
httpPOST /traffic/v1/evacuations HTTP/1.1 Content-Type: application/json { "evacuate": "us-east-1", "targets": { "us-west-2": 1.0 }, "reason": "SPS -38% in us-east-1; probes failing from us-west-2 and eu-west-1", "requested_by": "oncall:traffic", "idempotency_key": "evac-2026-09-28-01" }
httpHTTP/1.1 202 Accepted Content-Type: application/json { "evacuation_id": "evac-2026-09-28-01", "state": "CAPACITY_ATTACHING", "checks": { "targets_capacity_ready": "pending", "ec2_vcpu_quota_headroom": "ok", "residency_rules": "ok", "at_least_two_regions_on": "ok" } }
Replicate member data everywhere. Each Cassandra keyspace for member data keeps three copies in every region (CQL, a SQL-like schema language):
sqlCREATE KEYSPACE member_playback WITH replication = { 'class': 'NetworkTopologyStrategy', 'us-east-1': 3, 'us-west-2': 3, 'eu-west-1': 3 };
R3.4 Design Evolution: Our Own CDN and Three Regions
Step 3.1: Video Is Expensive and Far Away
The problem: every evening, video is a large share of each ISP's traffic, and it arrives over the ISPs' paid links from CDN servers outside their networks. In Australia, video from our US storage crosses undersea cables. We pay per byte for all of it. What would you do?
Synthesizing vector architecture diagram...
Fill runs at night in each time zone. Only the fill master pulls from S3; everything else copies from a nearby OCA, so a new episode crosses the ocean once, not once per viewer.
Loop: Design a Video Streaming Platform (multi-CDN steering and ISP caches from the buyer's side, step 3.2)
Step 3.2: A Region Failed and Members Couldn't Play
The problem: on Christmas Eve, a regional AWS service failed across all AZs in us-east-1, and members served from there couldn't start playing for hours, while the video servers were fine.
What would you do?
Primitive: Cloud Disaster Recovery and Multi-Region Active-Active
Step 3.3: Move All Traffic Out of a Region, Fast
The problem: us-east-1 is failing. We have to move its members to the other regions. Before Netflix's 2018 redesign, the whole operation took about 50 minutes.
What would you do?
Step 3.4: Evacuation Hits Replication Lag and Clock Trouble
The problem: during an evacuation, members arrive in us-west-2 a few seconds after their last writes in us-east-1. Some bookmarks haven't replicated yet. And last month one host's clock ran five minutes fast, and some members' bookmarks stopped updating for five minutes.
What would you do?
Drill: The booking that existed in Frankfurt but not in Virginia · Loop: Distributed key-value store, steps 3.1 and 3.2 (multi-region replication and conflicts from the inside)
Step 3.5: Did the Evacuation Work?
The problem: we flipped the routing controls, DNS changed, and the dashboards for us-east-1 went quiet. Is that success?
What would you do?
Round 3 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 3.1 | Video is expensive and far away | Open Connect: OCAs at IXPs and inside ISPs, filled off-peak, steered by the control plane | Hardware, operations, ISP partnerships |
| 3.2 | A region failed and members couldn't play | Active-active in three regions; stateless, local-only services; Cassandra and EVCache replicated globally | Three copies; lag; the no-cross-region-calls rule |
| 3.3 | Leave a region fast | Dark capacity (Nimble); detection from outside; gated ARC routing controls; redirect mode | Dark capacity; 9 minutes of impact |
| 3.4 | Lag and clocks | Warm global caches; a lag budget; synced clocks; SERIAL only for one-winner writes | A few seconds of bookmarks; slower signups |
| 3.5 | Did it work? | SPS; Chaos Kong; split-brain exercises; a scorecard | Practice time and bounded risk |
R3.5 Global Architecture
Synthesizing vector architecture diagram...
Two planes again, now fully separated. The control plane runs in three regions that can each serve anyone, with data replicated among them and traffic steered by controls that don't live in any one region. The data plane is Netflix's own network inside ISPs, told what to serve by the control plane. Netflix published the regions, the replication, Zuul's role and Open Connect; the ARC and Route 53 mapping, the probes and the dark-capacity labels on the diagram are our AWS rendering of the published ideas.
Sources for this round
- Netflix, Open Connect Overview (versions 2016 to 2026): OCAs at IXPs and embedded in ISPs, what OCAs report, the steering flow, fill windows.
- Netflix, Open Connect (checked September 2026): "over a thousand ISPs".
- Netflix, How Netflix Works With ISPs Around the Globe, March 2016: built from 2011, 100% of video, tens of Tbps, close to 90% direct, 8 to 90+ Gbps per server, the Australia example, 190 countries.
- Costello and Livengood, Netflix and Fill, Netflix TechBlog, August 2016: proactive caching, fill tiers, fill masters.
- Gallatin, Serving Netflix Video at 400Gb/s on FreeBSD, EuroBSDCon 2021.
- Cockcroft, A Closer Look at the Christmas Eve Outage, December 2012.
- Meshenberg, Gopalani and Kosewski, Active-Active for Multi-Regional Resiliency, December 2013: the three rules, geo-DNS and Denominator, Zuul's misrouting and shedding, the 500 ms test, Chaos Kong and split-brain.
- The EVCache team, Caching for a Global Netflix, March 2016: three regions, global replication design and latencies.
- Kosewski, Ramanujam, Behnam, Blohowiak and Probst, Project Nimble: Region Evacuation Reimagined, March 2018: the 50-minute breakdown, dark capacity, 8 minutes.
- Fisher-Ogden, Sanden and Rioux, SPS: the Pulse of Netflix Streaming, February 2015.
- Basiri, Hochstein, Thosar and Rosenthal, Chaos Engineering Upgraded, September 2015: Chaos Kong, the September 2015 event.
- Netflix, Q4 2025 shareholder letter, January 2026: more than 325M paid memberships.
- AWS, Enabling accelerated recovery for managing public DNS records: opt-in, 60-minute target.
- Netflix, What We Watched the First Half of 2026: more than 97B hours.
Netflix has not published its current traffic-steering tools, its present number of AWS regions, per-instance capacities, cache sizes or evacuation detection rules. Everything about those on this page is our design, and so are the probes, the ARC mapping, the scorecard and the consistency rules in step 3.4.
R3.6 Numbers and Cost
Figures marked published come from the sources above; everything else is our assumption or derived from one.
The data plane: how much video
| Quantity | Arithmetic | Result |
|---|---|---|
| Hours watched a day | 97B hours (published, H1 2026) ÷ 181 days | ≈ 536M hours |
| Average streams playing at once | 535.9M hours ÷ 24 h | ≈ 22.3M |
| Average video throughput | 22.33M × 3 Mbps (assumption) | ≈ 67 Tbps |
| Global evening peak | 67 Tbps × 1.5 (assumption: time zones flatten the global peak) | ≈ 100 Tbps |
| If the average bitrate were 2 or 5 Mbps | 22.33M × 2 Mbps; 22.33M × 5 Mbps | 45 or 112 Tbps on average |
| Bytes an hour of viewing | 3 Mbps × 3,600 s ÷ 8 | 1.35 GB |
| Bytes a day | 535.9M hours × 1.35 GB | ≈ 723 PB |
| Bytes a month (30 days) | 723 PB × 30 | ≈ 21.7 EB |
A sanity check against 2016. Netflix then delivered more than 125M hours a day at "tens of terabits per second" of peak. The same model gives 125M ÷ 24 = 5.2M streams × 3 Mbps = 15.6 Tbps on average, × 1.5 = 23.4 Tbps at the peak: tens of terabits. The model is in the right range.
Why own the CDN: a made-up price. Suppose a CDN charged half a cent per GB, a round number we made up for illustration, not a quote. 21.7 billion GB × $0.005 = $108.5M a month, rising with every hour watched. Money is only part of Netflix's published reasoning: the other parts are working directly with ISPs, keeping bytes off backbones and undersea cables, and caching proactively instead of on demand.
Servers are not the limit; placement is. At the 400 Gb/s per server that Netflix engineers described in 2021, 100 Tbps is 100,000 ÷ 400 = 250 servers running flat out. Open Connect has thousands of appliances because each ISP site is sized for its own local peak and storage, and most sites are small. The design problem is putting the right files in the right thousand places, not raw throughput.
The control plane in each region. Assume 3M edge requests/s worldwide at the peak (Netflix published more than 1M in 2018 with 125M members; 325 ÷ 125 = 2.6 times the members, rounded up), 1M per region, and 5,000 requests/s per Zuul instance, as in Round 2.
| Situation | Load on the region | Instances needed | Per AZ, rounded up | Fleet |
|---|---|---|---|---|
| Normal peak, sized to lose an AZ | 1M/s | 200 in 2 AZs | 100 | 300 |
eu-west-1 evacuated, split across both US regions | 1M + 0.5M = 1.5M/s | 300 in 3 AZs | 100 | 300: no extra, but no AZ headroom during the evacuation |
us-east-1 evacuated to us-west-2 (or the reverse) | 1M + 1M = 2M/s | 400 in 3 AZs | 133.3 → 134 | 402: 102 dark instances, 34 per AZ |
The same 50% headroom that protects a region from losing one of three AZs also covers an even split of one of three regions, because in both cases a third of the capacity disappears and the other two-thirds must carry it. We accept not surviving both at once. The US pair evacuates to each other at full load, so each US region keeps about a third more capacity dark (102 ÷ 300 = 34%) in every tier on the streaming path, sized by time of day as Nimble did. Netflix reported Nimble as cost neutral without saying how, so we don't guess at the saving.
Quotas. 402 Zuul instances at 16 vCPUs each (assumption) is 6,432 vCPUs, and every other tier on the streaming path grows by the same third. The target region's EC2 On-Demand vCPU quota must cover the evacuation fleet for all of them, not today's peak. Dark instances already count against the quota; the check is headroom for autoscaling beyond them. We check it monthly and before every exercise (R3.9).
Caches during an evacuation. Round 2 had 10M cache operations/s for 500,000 edge requests/s: 20 per request. At 1M/s per region that's 20M operations/s, and the target of a paired evacuation must serve 40M/s. RAM doesn't need to grow, because globally replicated caches already hold every region's members; throughput does. This holds only for fully replicated caches. Caches that replicate invalidations only are empty for moved members, so we use invalidate-only replication only where Cassandra can absorb the moved members' miss rate. Caches on the streaming path replicate data.
A cold cache would be a stampede. The evacuated region's members make about 18M cache reads/s (90% of 20M). At a 99% hit ratio (assumption), Cassandra sees 18M × 1% = 180,000 reads/s. With cold caches, all 18M/s would go to Cassandra: 100 times its normal load, at the worst moment.
Replication between regions (list price $0.02/GB for data sent from us-east-1, us-west-2 or eu-west-1 to another region; the receiver pays nothing; other regions charge more).
| Stream | Arithmetic | A month |
|---|---|---|
| Bookmarks | 22.33M streams ÷ 60 s = 372,000 writes/s on average; each write sent to 2 other regions × 300 B (assumption) = 600 B; 372,000 × 600 B = 223 MB/s = 19.3 TB a day × $0.02/GB = $386 a day | ≈ $11,600 |
| EVCache replication | More than 1M requests/s at peak (published); assume 500,000/s on average × 2 KB = 1 GB/s = 86.4 TB a day × $0.02/GB = $1,728 a day | ≈ $51,800 |
| Total | ≈ $63,400 |
Traffic control. One ARC cluster for the routing controls costs $2.50 per cluster-hour × 730 h ≈ $1,825 a month, whether or not we ever flip it.
At the peak, bookmarks run at 33.5M streams ÷ 60 s ≈ 558,000 writes/s, the figure step 3.4 used. Replication is cheap next to the fleet; what it costs in design effort is the no-cross-region-calls rule.
What a failure costs in availability. Netflix's internal goal is 99.99%, 4.32 minutes a month. A real region failure evacuated in 9 minutes hurts about a third of members for those 9 minutes: 9 × ⅓ = 3 member-weighted minutes, about 69% of the month's budget from one event. This is why cutting 50 minutes to 8 matters; Netflix's stated reason was that "even short or partial outages affect many of our customers."
R3.7 Trade-Offs
Build our own CDN vs buy
| Buy (third-party CDNs) | Build (Open Connect, chosen) | |
|---|---|---|
| Caching | On demand: the first viewer in each place causes a miss | Proactive: files placed at night before anyone asks |
| Where servers sit | Mostly outside ISPs | Inside ISPs and at exchanges; ISPs choose which customers use embedded OCAs |
| Cost shape | Per GB, forever | Hardware, operations and partnerships; bytes get cheaper as servers improve (8 to 90+ Gbps in four years) |
| Control | The CDN's software and placement | Hardware, software, placement and steering are all ours |
| Makes sense when | Traffic is modest or unpredictable | You are a large, predictable share of every ISP's evening traffic |
Active-active vs active-passive
| Active-passive | Active-active (chosen) | |
|---|---|---|
| Is the backup tested? | Rarely, with no real traffic | Every day, with real traffic |
| Caches in the other region | Cold | Warm, replicated |
| Capacity | Paid for and idle, or not there when needed | Serving, plus dark capacity for the move |
| Complexity | Failover runbooks | No cross-region calls; conflict rules for replicated data |
| Failure of the switch itself | Common, since it's rarely run | Practiced by regular exercises |
Consistency, chosen per kind of data (ours)
| Data | Rule | Why it's acceptable |
|---|---|---|
| Bookmarks | Last writer wins by synced server time; asynchronous replication | A position a second old is harmless; the player repairs it |
| Viewing history | Append by time; asynchronous replication | A late record appears late, never wrongly |
| Cached recommendations and profiles | EVCache eventual consistency; invalidate-only replication when remote reads are rare, and not on the streaming path | Netflix: "it's okay ... if Ireland and Virginia occasionally have slightly different recommendations for you" |
| One account per email address | SERIAL lightweight transaction across all regions | Must have one winner; rare enough to afford cross-region agreement |
Automatic vs approved evacuation. Fully automatic saves about a minute and risks evacuating a healthy region on a bad signal. We automate everything except the decision, and make the decision one action with the evidence in front of the on-call engineer. The capacity step starts during that minute anyway.
What changed from Round 1. Round 1 asked how to survive losing a database and answered with statelessness, replication and three AZs. Round 3 asks the same question one level up, about a whole region and about the internet itself, and answers with the same ideas at a bigger scale: nothing exists once (three regions, thousands of OCAs), state is replicated and its consistency chosen on purpose, and the control plane and data plane never share a fate. The split we made in step 1.2 to keep video off our API servers is the reason a region failure in Round 3 costs Play buttons, not video already playing.
R3.8 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| Cold caches during an evacuation (replication was broken, or a cache was never replicated) | Cassandra read load in the target jumps toward 100 times normal | Shed at the edge above a set traffic level, as Netflix's Zuul did; move traffic in steps rather than at once when the source region is still usable; collapse concurrent misses per key; concurrency limits in front of Cassandra so it degrades instead of collapsing. |
| Replication lag when a region dies | Some members resume a few seconds earlier than they stopped | Within the lag budget (about a second); the player's next update, now to the new region, repairs the bookmark. |
| A host clock runs fast | Bookmarks written through it win over later correct writes; some members' positions freeze | Clock offset over 100 ms fails the host's health check. The frozen positions heal by themselves once real time passes the bad timestamps (up to five minutes here); to heal them sooner, a repair job rewrites each affected bookmark from the player's latest update with a timestamp just above the bad one, since a write with a correct, smaller timestamp would lose again. |
| A replication batch keeps failing | Relay consumer lag grows on one partition while others drain | Split the failing batch in halves until the single bad message is found (bisect); skip it and send a delete for its key to the other regions, so they miss and reload from Cassandra instead of serving the old value. Alarm on maximum lag per partition, not the average, as Netflix does. |
| A target region was unreachable longer than the replication queue holds | Netflix: Kafka "will start dropping older messages"; the dropped updates never arrive | Treat that region's replicated caches as suspect: before sending it traffic, invalidate or re-copy them. A re-copy notes the queue position first, copies from a healthy region, then replays from the noted position; if the copy took longer than the queue's retention, start again. Keys deleted after the copy was taken are removed by the replayed deletes; every item's TTL is the backstop. |
| Evacuation on a false signal | Probes from one region fail because of its own network | The trigger needs probes from both other regions and an SPS drop; a human approves; the safety rule keeps two regions on. |
| An appliance fails | Players on it stall | Players move to the next ranked URL after a couple of failed or slow fetches; steering stops choosing it when its reports stop or turn unhealthy; ISPs with embedded OCAs also peer with Open Connect "for resiliency", so traffic can reach exchange OCAs. |
| A slow memory leak in the gateway | Direct memory creeps up; at twice the normal load during an evacuation, it runs out in hours instead of days | Alarm on direct memory; rolling restarts; leak detection on canaries (Round 2). Every Zuul build ships through a canary before it can meet evacuation load. |
| Regions lose contact with each other (split brain) | Replication queues grow; each region keeps serving | Netflix's split-brain exercise showed each region working while replication queued; SERIAL signups fail in the minority side until the link returns, which we accept. |
R3.9 Runbook and Incident Response
| Signal | Alarm | Severity | First action |
|---|---|---|---|
| SPS per region vs its expected band (computed in every region) | Below the band for 2 min | P1 | Check probes from the other regions; start the evacuation procedure if both fail |
| SPS worldwide | Below the band for 2 min | P1 | Is it one region, one device family, or a release? |
| 5xx rate at the edge, per region | > 1% for 5 min | P2 | Which origins? Zuul's per-origin error rates |
| Evacuation readiness: dark capacity attached-ready, and quota headroom, per region | Below the time-of-day target | P2 | Refill dark groups; request a quota increase |
| Replication lag, p99 and max per partition | > 1 s for 5 min | P2 | Relay health; the cross-region link; bisect a stuck batch |
| Host clock offset | > 100 ms on any host | P2 | Host out of service; check chrony |
| Share of video from OCAs inside ISPs, per country | Drops 5 points in 15 min | P2 | Which sites stopped reporting? BGP sessions? Fill failures? |
Evacuation procedure REL 13 · OPS 10
- Confirm it's a region. SPS for the region is below its band in the views of the other regions, and probes from both other regions fail. One region's probes alone, or a drop across every region, is not a regional failure.
- Start capacity at once. In each target region, raise the production groups' maximum size, then move dark instances into production (commands 1 to 3). They reach
UPin about a minute. This is safe to undo. - Check the gate. Targets' capacity
UP; vCPU quota headroom (command 4) (dark instances already count against the quota; the check is headroom for autoscaling beyond them); the target map satisfies residency rules; two regions will remain on. - Approve and flip. Turn the failed region's routing control off through any ARC cluster endpoint; the five endpoints are in five regions, so try them in turn until one answers (commands 5 and 6).
- Drain what DNS won't. If the failed region's Zuul answers, switch it to redirect mode so persistent connections close gradually and misrouted requests are sent to the targets.
- Watch SPS, not the DNS change (command 7): the worldwide SPS should return to its band within about 5 minutes of the flip, and the failed region's should fall toward zero.
- Return later, slowly. Once the region is healthy, replication has caught up and its caches are checked (R3.8), turn its routing control back on and move traffic back in steps, watching SPS all the way.
Go deeper: CLI playbook
Commands an on-call engineer runs one at a time. Replace names, IDs, ARNs and endpoints with real ones. attach-instances takes at most 20 instance IDs per call.
text# 1. Production and dark group sizes in the target region aws autoscaling describe-auto-scaling-groups --region us-west-2 --auto-scaling-group-names zuul-prod zuul-dark --query "AutoScalingGroups[].[AutoScalingGroupName,DesiredCapacity,MaxSize,length(Instances)]" # 2. Make room in the production group (attach fails if it would exceed the maximum) aws autoscaling update-auto-scaling-group --region us-west-2 --auto-scaling-group-name zuul-prod --max-size 420 # 3. Move dark instances into production, in batches of up to 20 aws autoscaling detach-instances --region us-west-2 --auto-scaling-group-name zuul-dark --instance-ids i-0abc1234def567890 i-0abc1234def567891 --should-decrement-desired-capacity aws autoscaling attach-instances --region us-west-2 --auto-scaling-group-name zuul-prod --instance-ids i-0abc1234def567890 i-0abc1234def567891 # 4. The Running On-Demand Standard instances vCPU quota in the target region aws service-quotas get-service-quota --region us-west-2 --service-code ec2 --quota-code L-1216C47A # 5. Read the routing control state through one of the cluster's five endpoints aws route53-recovery-cluster get-routing-control-state --region us-west-2 --endpoint-url https://<cluster-endpoint-us-west-2>/v1 --routing-control-arn arn:aws:route53-recovery-control::123456789012:controlpanel/0123456789abcdef/routingcontrol/abcdef0123456789 # 6. Turn off the failed region (the safety rule refuses it if fewer than two regions would stay on) aws route53-recovery-cluster update-routing-control-state --region us-west-2 --endpoint-url https://<cluster-endpoint-us-west-2>/v1 --routing-control-arn arn:aws:route53-recovery-control::123456789012:controlpanel/0123456789abcdef/routingcontrol/abcdef0123456789 --routing-control-state Off # 7. Stream starts per minute served by each region (our own metric) aws cloudwatch get-metric-statistics --region us-west-2 --namespace Streaming/Playback --metric-name StreamStarts --dimensions Name=ServingRegion,Value=us-west-2 --start-time 2026-09-28T20:00:00Z --end-time 2026-09-28T20:30:00Z --period 60 --statistics Sum
R3.10 Pillar Check
| Pillar | What Round 3 adds |
|---|---|
| Reliability | Three active-active regions with no cross-region calls; dark capacity; detection from outside the failed region; gated routing controls with a two-regions-on safety rule; regular region-scale exercises REL 10 · REL 13 · REL 12 |
| Performance Efficiency | Video from inside the viewer's ISP; steering by file availability, health and network proximity; globally replicated caches so moved traffic stays fast PERF 4 · PERF 1 |
| Security | Appliances hold no member data; signed, expiring video URLs; SERIAL agreement for identity-defining writes like one account per email; residency rules in the evacuation gate; video served over TLS SEC 7 · SEC 9 |
| Cost Optimization | Owning delivery at about 21.7 EB a month; replication at about $63,000 a month at list price; dark capacity sized by time of day instead of permanent headroom COST 5 · COST 8 · COST 9 |
| Operational Excellence | SPS as the one signal every team shares; an evacuation procedure and a return procedure; an exercise scorecard OPS 8 · OPS 10 · OPS 11 |
| Sustainability | Bytes copied across oceans once, at night, instead of once per viewer; purpose-built servers that went from 8 to 90+ Gbps (2012 to 2016); peers fill peers SUS 3 · SUS 5 |
R3.11 Round 3 Rubric and Follow-Ups
What an architect (L7) answer adds over L6
- Separates the control plane from the data plane, and designs each for its own failure and cost model.
- Explains why proactive, ISP-embedded caching beats on-demand caching for a predictable catalog, and what the ISP relationship adds.
- Chooses active-active over active-passive and states the rules that make it work: stateless services, local resources only, no cross-region calls, asynchronous replication.
- Breaks evacuation time into phases, finds the long poles (boot time, DNS), and removes them with dark capacity; counts every delay.
- Detects and moves traffic without depending on the failed region, and gates the move on capacity, quotas and residency.
- Chooses consistency per kind of data, knows last-writer-wins depends on clocks, and knows local conditional checks can't give one winner across regions.
- Proves it works with a member-facing metric and regular exercises.
- Keeps Netflix's published facts separate from its own design.
Follow-up questions
-
"Why didn't Netflix use latency-based DNS routing?" Answer: Netflix wrote that latency-based routing "could cause unpredictable traffic migration effects". With geographic routing, which users go where is decided by us, so capacity plans and evacuations are predictable.
-
"A new blockbuster launches worldwide at the same minute. What does Open Connect do differently from a normal CDN?" Answer: it doesn't wait for the first viewer. The files are copied to appliances during the preceding nights' fill windows, based on predicted popularity, peer to peer within each area, so at launch every nearby appliance already has them.
-
"What if the region you need to evacuate is the one that hosts your DNS provider's control plane?" Answer: that's exactly why the move uses ARC routing controls, whose data plane has endpoints in five regions, rather than editing records through the Route 53 API in
us-east-1(whose opt-in accelerated recovery targets 60 minutes), and why detection comes from probes in the other regions rather than health-check metrics stored inus-east-1. -
"Is Hystrix still how Netflix does resilience?" Answer: no. Hystrix is in maintenance mode; Netflix pointed to adaptive concurrency limits and resilience4j, and has been moving features like adaptive limits, circuit breaking and hedging into an Envoy-based service mesh.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "Buy more CDN capacity" | On-demand caching, per-byte cost forever, and the bytes still cross the same links. |
| "Active-passive is safer" | An untested standby with cold caches fails the one time it's used. |
| "Just change DNS" | Survivors lack capacity, new instances take many minutes, and devices follow DNS slowly. |
| "Fail over on a health check" | The checker may live in the failed region, and a false positive evacuates a healthy one. |
| "Route 53 anycast latency routing moves users" | Netflix chose geographic routing on purpose; DNS changes still take minutes to reach devices. |
| "Eventual consistency covers it" | It doesn't say which write wins, and clocks decide last-writer-wins. |
| "The DNS call succeeded, so we're done" | Only member-facing SPS tells you if members can play. |
Loop Closer: Interview Strategy for All Three Rounds
How to Run Each 60-Minute Round
| Time | Round 1 | Round 2 | Round 3 |
|---|---|---|---|
| 0–5 min | Scoping: what broke in 2008, what the service does, where video comes from | Restate Round 1 in 60 seconds | Restate Round 2 in 60 seconds |
| 5–15 min | Requirements and the manifest API | Scope raise → what breaks | Scope raise → what breaks |
| 15–40 min | Steps 1.0–1.4: one database → cloud and AZs → control vs data plane → Cassandra → autoscaling | Steps 2.1–2.5: Zuul → Eureka → breakers and limits → EVCache → chaos | Steps 3.1–3.5: Open Connect → active-active → evacuation → lag and clocks → proof |
| 40–50 min | Fleet per AZ, bookmark writes, the 214× byte ratio | Connections per instance, Little's law, cross-AZ and LCU costs | Video volume, evacuation capacity and quotas, replication cost, the timeline |
| 50–60 min | Failures and pillar check | Failures, gotchas, pillar check | Failures, runbook, pillar check |
For how to spend a single 45-minute round, see the 45-minute interview blueprint. For transcoding and adaptive bitrate, see the video streaming loop; for load shedding and limits, the rate limiter loop; for why retries must be budgeted and what a provider's 429 means, the notification system loop.
The Two Sentences That Matter Most
- Opening any round: "Netflix is two systems: a small control plane that decides what you can play and where to fetch it, and an enormous data plane that moves the video; I'll keep them apart so they never share a failure or a bill."
- When scale arrives: "I'll put one programmable front door in front of the services, contain every dependency with limits and fallbacks, run the control plane active-active in three regions with capacity waiting for an evacuation, and put the video inside the viewers' ISPs, filled before anyone asks."
Well-Architected Review Sheet
Interviewers rarely ask "which pillar is this?". They ask the pillar's question in plain words. Rehearse one sentence per row.
| Pillar | Question you'll hear | One-sentence answer | Round | Backed by |
|---|---|---|---|---|
| Reliability | "What if a database or a server dies?" (REL 10) | Stateless services in three AZs, Cassandra with a copy per AZ and quorum writes, sized to lose an AZ at the peak. | 1 | Steps 1.1, 1.3; R1.7 |
| "One service is slow. Why is the whole page down?" (REL 5) | A shared pool fills by Little's law; bulkheads, breakers, fallbacks and adaptive limits contain it. | 2 | Step 2.3 | |
| "What if a whole region fails?" (REL 13) | Active-active in three regions; dark capacity; routing controls flipped from outside the failed region; about 9 minutes. | 3 | Steps 3.2, 3.3 | |
| "How do you know failover works?" (REL 12) | Chaos Monkey, ChAP and regular Chaos Kong exercises, judged by SPS. | 2–3 | Steps 2.5, 3.5 | |
| Performance | "How do you serve millions of idle connections?" (PERF 2) | An event loop per core in Zuul 2: a connection costs a socket, not a thread. | 2 | Step 2.1, R2.6 |
| "How does video stay fast?" (PERF 4) | Appliances inside the viewer's ISP, chosen by health and network proximity, filled before demand. | 3 | Step 3.1 | |
| Security | "How do services know who the user is?" (SEC 2) | Zuul authenticates at the edge and passes an HMAC-protected Passport to every service. | 2 | Step 2.1 |
| "Where may member data live during an evacuation?" (SEC 7) | Wherever the replication choice put it; the evacuation gate checks a pre-approved target map. | 3 | Step 3.3 | |
| Cost | "Why build your own CDN?" (COST 5) | About 21.7 EB a month of predictable demand: Netflix's published reasons are efficiency (proactive caching, bytes kept off backbones) and working directly with ISPs. | 3 | Step 3.1, R3.6 |
| "Where do data transfer charges hide?" (COST 8) | Cross-AZ reads ($315,000 a month avoided by local reads), load balancer bytes, and a second load balancer's copy. | 2 | R2.6 | |
| Operations | "How do you know members are OK?" (OPS 8) | Stream starts per second against its expected band, per region and worldwide. | 3 | Step 3.5, R3.9 |
| Sustainability | "How do you avoid waste?" (SUS 3) | Autoscaling with the evening, and video copied across oceans once, at night. | 1–3 | Step 1.4, step 3.1 |
Rubric Across Levels
| Dimension | L5 (Round 1) | L6 (Round 2) | L7 (Round 3) |
|---|---|---|---|
| Failure scope | A server, a database, an AZ | A slow dependency, a bad deploy, a retry storm | A whole region, a CDN site, a clock |
| Edge and traffic | One load balancer, stateless API | Zuul filters, device APIs, Passport, Netty for connections | Geographic routing, gated routing controls, redirect mode, dark capacity |
| Data | Cassandra keyed by profile, quorum writes | EVCache per AZ, cache-aside, TTLs | Global replication, lag budgets, clocks, SERIAL only where one winner is required |
| Video | Third-party CDNs through signed manifests | (unchanged) | Open Connect: embedded OCAs, fill windows, steering |
| Proof | Health checks and alarms | Chaos experiments with a steady state | Region evacuations judged by SPS |
| Honesty | Labels assumptions | Knows Hystrix's status and async's real gains | Separates what Netflix published from its own design |
Sources
All Netflix sources used on this page, oldest first.
- Netflix, Q4 2008 results press release, January 2009.
- Ciancutti, Four Reasons We Choose Amazon's Cloud as Our Computing Platform, December 2010.
- The Netflix Simian Army, July 2011.
- Schmaus, Making the Netflix API More Resilient, December 2011.
- Jacobson, Embracing the Differences: Inside the Netflix API Redesign, July 2012.
- Cockcroft, A Closer Look at the Christmas Eve Outage, December 2012.
- Meshenberg, Gopalani and Kosewski, Active-Active for Multi-Regional Resiliency, December 2013.
- Fisher-Ogden, Sanden and Rioux, SPS: the Pulse of Netflix Streaming, February 2015.
- Basiri, Hochstein, Thosar and Rosenthal, Chaos Engineering Upgraded, September 2015.
- Izrailevsky, Vlaovic and Meshenberg, Completing the Netflix Cloud Migration, February 2016.
- The EVCache team, Caching for a Global Netflix, March 2016.
- Netflix, How Netflix Works With ISPs Around the Globe, March 2016.
- Costello and Livengood, Netflix and Fill, August 2016.
- The Cloud Gateway team, Zuul 2: The Netflix Journey to Asynchronous, Non-Blocking Systems, September 2016.
- ChAP: Chaos Automation Platform, July 2017.
- Duvedi, Li, Garg and Fisher-Ogden, Scaling Time Series Data Storage, Part I, January 2018.
- Kosewski, Ramanujam, Behnam, Blohowiak and Probst, Project Nimble: Region Evacuation Reimagined, March 2018.
- Landau, Thurston and Bozarth, Performance Under Load, March 2018.
- Gonigberg et al., Open Sourcing Zuul 2, May 2018.
- Edge Authentication and Token-Agnostic Identity Propagation, February 2021.
- Gallatin, Serving Netflix Video at 400Gb/s on FreeBSD, EuroBSDCon 2021.
- Vroom, Mulcahy, Yuan and Gulewich, Zero Configuration Service Mesh with On-Demand Cluster Discovery, August 2023.
- Netflix, Q4 2025 shareholder letter, January 2026.
- Netflix, What We Watched the First Half of 2026, 2026.
- Open-source documentation: Netflix/zuul, Netflix/Hystrix, the EVCache wiki, the Eureka wiki, and Netflix's Open Connect Overview and Open Connect site, checked September 2026.