KnowSys

Reliability Engineering

Follow one click on a Pay button through a handful of services: how often it works when every service is very good, how to turn that into a target you can measure and spend, and the four defences (timeouts, retries with jitter, circuit breakers and load shedding) that stop one slow service from taking the rest down.

⏱ 40 min read◆ BeginnerAssumes: a terminal and Python; RPCs and timeouts help
Start reading

You click Pay in an online shop, and half a second later the page says "Order placed". Behind that click, the code in your browser called the shop's frontend, the frontend called the orders backend, and the backend asked a handful of other programs for help: who is this (auth), what's in the cart (cart), is it in stock (inventory), take the money (payments), and write the order down (the orders database). Each of these is a service, a program that answers requests from other programs over the network, and a service that another service calls on is its dependency. The click works only if all five dependencies answer.

Suppose each of the five is excellent and answers correctly 999 times out of 1,000. It's natural to expect the click to be about that reliable. It won't be. It will fail more often than any single dependency, and when conditions get bad, the services' own attempts to recover (asking again, queueing more work) can turn a small hiccup into an outage that lasts hours after the original problem is gone.

This chapter follows that one click and asks: how often will it work, how do we decide how often it should, and what stops one sick service from taking the whole shop with it? We'll start with the arithmetic, turn it into a target we can measure and spend, and then add the defences one at a time, each one because the one before it leaves a gap: timeouts, retries, circuit breakers and load shedding.

01How often the click works

1.1Counting in nines

We need a way to say how often something works. The usual one is availability, the share of time it works, written as a percentage. A service that works 999 minutes out of every 1,000 is 99.9% available. People call that "three nines", and 99% is "two nines", 99.99% is "four nines".

Each extra nine divides the allowed failure by ten, which makes the allowed downtime shrink fast. Over a 365-day year, 99.9% allows about 526 minutes of downtime, which is under nine hours, while 99.99% allows about 53 minutes. (Later in the chapter we'll count requests instead of minutes, which is how reliability is measured in practice, but the arithmetic is the same.)

1.2Why five good services make a worse one

Now back to the click. It needs auth, cart, inventory, payments and the orders database, and if any one of them fails to answer, the click fails. Suppose each answers 999 times in 1,000, and that their failures are unrelated to each other. The click works only when all five work, so the chance of that is 0.999 × 0.999 × 0.999 × 0.999 × 0.999, which is 0.999 to the fifth power.

Think of a relay team where every handoff can drop the baton. Each extra handoff multiplies the chance of a clean race by a number just under one, so a long chain of very good handoffs still loses more often than any one of them. Dependencies work the same way. You can put numbers on it with a dozen lines of Python. Save this as nines.py and run it with python3 nines.py. The first loop prints, for four availability levels, how many minutes of downtime that level allows per year and per 30 days. The second loop raises 0.999 to the power of 1, 3, 5 and 10 and converts what's lost into hours per year. (The :>13.3% pieces are format codes that right-align a number and print it as a percentage.)

Print the downtime allowed by 2 to 5 nines, then the availability of 1 to 10 dependencies at 99.9% each
python
Python
year_min = 365 * 24 * 60
print(f"{'availability':>13} {'downtime / year':>17} {'downtime / 30 days':>19}")
for a in (0.99, 0.999, 0.9999, 0.99999):
    print(f"{a:>13.3%} {(1 - a) * year_min:>13.1f} min {(1 - a) * 30 * 24 * 60:>15.1f} min")
 
print()
for n in (1, 3, 5, 10):
    print(f"{n:>2} dependencies, each 99.9%, all needed: {0.999 ** n:.3%}  ({(1 - 0.999 ** n) * year_min / 60:.1f} hours down a year)")
output
C++
 availability   downtime / year  downtime / 30 days
      99.000%        5256.0 min           432.0 min
      99.900%         525.6 min            43.2 min
      99.990%          52.6 min             4.3 min
      99.999%           5.3 min             0.4 min
 
 1 dependencies, each 99.9%, all needed: 99.900%  (8.8 hours down a year)
 3 dependencies, each 99.9%, all needed: 99.700%  (26.3 hours down a year)
 5 dependencies, each 99.9%, all needed: 99.501%  (43.7 hours down a year)
10 dependencies, each 99.9%, all needed: 99.004%  (87.2 hours down a year)

In the top table you can see the factor of ten: 99.9% allows 525.6 minutes a year and 99.99% allows 52.6. The bottom lines are the click. One dependency at 99.9% costs 8.8 hours of downtime a year, and our five-service click costs 43.7 hours, about five times as much, because each dependency's failures add to the total. Ten dependencies cost about 87 hours, and the click is down to roughly 99.0%, two nines.

A runner holding out a relay baton to a teammate who is reaching back for it
A relay handoff. Every handoff is a chance to drop the baton, so the more handoffs a race has, the lower its odds of a clean finish, however good each runner is. A click that passes through five services takes that chance five times.Photo: Patrick Bell, CC BY 2.0, via Wikimedia Commons

1.3What you can do about it

Here are the numbers so far side by side, with one extra row at the bottom that we'll come to in a moment.

One dependency at 99.9%0.99999.9%
Five in series0.999⁵99.50%
Ten in series0.999¹⁰99.00%
Two replicas at 99%, either will do1 − 0.01²99.99%
ten 99.9% dependencies leave you attwo nines

If the shop wants to promise 99.9% and needs five dependencies that each deliver 99.9%, it has already broken the promise on paper, before any bug of its own. There are three moves: make the dependencies better, need fewer of them at once, or let the click survive when one of them fails.

Of the three, the last is probably the most useful, and the bottom row of the table shows one way to do it. Run a replica, a second copy of a service, so the click needs either copy and not a particular one. Each copy at 99% fails 1 time in 100, so both fail together only 1 time in 10,000 (0.01 × 0.01), and the pair is 99.99% available. When the click needs every dependency, each one's chance of failing adds to the total. When it needs only one of two copies, the chances of failing multiply together and become tiny.

That works only if the failures are independent. Two replicas that both read the same bad configuration push, or that sit in the same data centre during a power event, fail together, and the pair is no better than one copy.

So the arithmetic tells us what the shop can achieve. It doesn't tell us what the shop should promise, or how anyone would know whether the promise was kept. For that we need a target made of numbers that someone can measure.

02Putting a number on reliability

"The shop should be reliable" can't be tested, alerted on or traded off against anything. The first job is to turn it into a number.

2.1SLI, SLO and SLA

Start by deciding what counts as a click working. Say a request is good if the shop answers it without a server error (an HTTP status in the 500s, written 5xx) in under 300 milliseconds. Count the good requests and the valid ones, which are the requests that count towards the measurement. The ratio of good to valid is called the SLI, the service level indicator: it's the measurement itself, such as "the share of requests answered well".

Next comes a target for that measurement over a window of time. Say 99.9% of requests good, over a rolling 30 days (a window that always covers the most recent 30 days). That target is the SLO, the service level objective. If the target is written into a contract that has consequences when it's missed, that contract is the SLA, the service level agreement.

TermWhat it isExample
SLI (indicator)A measurement of good events over valid eventsShare of HTTP requests that returned non-5xx in under 300 ms
SLO (objective)A target for the SLI over a window99.9% of requests good, over a rolling 30 days
SLA (agreement)A contract with consequences if the SLO is missedService credits below 99.5% in a calendar month

Set the SLA looser than the SLO, and the SLO looser than what the service usually achieves. The gap between them is warning time: the team should learn that it's missing the internal target long before a contract notices.

?Why measure requests instead of uptime?

Because "the server was up" isn't what users experience. A service can answer health checks all day while 5% of real requests time out. An SLI counted from real requests, ideally at the load balancer (the program in front of the servers that spreads requests across them) or in the client, sees what users see.

An SLI should also be a ratio of events and not an average or a percentile. "p99 latency under 300 ms" (the time that 99 of every 100 requests beat) is hard to budget against. "99% of requests completed in under 300 ms" is a count of good events that you can add up across servers and over time.

2.2What each nine costs

To pick the number, we need to know what each level asks of the team. Here are the allowed bad minutes at each level:

SLOBad time per yearPer 30 daysPer week
99%3.65 days7.2 hours1.68 hours
99.9%8.77 hours43.2 minutes10.1 minutes
99.95%4.38 hours21.6 minutes5.04 minutes
99.99%52.6 minutes4.32 minutes1.01 minutes
99.999%5.26 minutes25.9 seconds6.05 seconds

?Why don't teams just aim for five nines?

Because of what the last row implies about people. Twenty-six seconds a month is less time than it takes a person to notice an alert, unlock a laptop and log in. At four or five nines, detection and recovery have to be automatic, rollouts have to be gradual enough that a bad one hurts a sliver of traffic, and every dependency has to be better still. Each extra nine is a different engineering organisation and costs far more than a tweak to the last one.

2.3Checking the promise against the dependencies

Section 1 gave us the check: the shop's SLO can't be higher than what its critical dependencies allow. Google's rule of thumb, from The Calculus of Service Availability (Treynor et al., 2017), is the rule of the extra 9: a critical dependency should offer one more nine than the service that depends on it. When a dependency can't, we change it from critical to something the click can survive without, using a cache, a fallback or graceful degradation (section 6.3 lists the options).

A 99.9% SLO says something else that turns out to be more useful than "be reliable": 0.1% of requests are allowed to fail. The next section is about what to do with that allowance.

03Error budgets and burn rates

3.1A budget meant to be spent

The allowance of failures that the SLO permits is the error budget. Say the shop's API takes 1,000 requests a second. Over 30 days that's about 2.6 billion requests, and a 99.9% SLO allows about 2.6 million of them to be bad. That's the budget, and it's meant to be spent: on deploys, experiments, migrations and the occasional bad day.

?Why is a budget better than "as reliable as possible"?

Because it turns an argument into a measurement. When the budget is healthy, the team ships faster. When it's nearly gone, the team slows releases and spends time on reliability. Nobody has to win a debate about whether the service is "reliable enough"; the counter says.

A budget also caps how reliable the service should be. A service at 99.99% against a 99.9% SLO is leaving budget unspent, which usually means it's shipping more slowly than it could.

3.2Burn rate

A budget raises a new question: how fast is it being used up? The burn rate compares the current speed of spending to the speed that would use the budget up exactly at the end of the window. A burn rate of 1 spends a 30-day budget in 30 days. A burn rate of 10 spends it in 3.

For a 99.9% SLO the budget is 0.1% of requests, so the burn rate is just the current error rate divided by 0.1%:

Error rateBurn rateBudget gone in
0.1%130 days
0.6%65 days
1.44%14.450 hours
10%1007.2 hours
100% (full outage)1,00043 minutes
Predict before you read on

Your service has a 99.9% monthly SLO. A bad deploy starts failing 1.44% of requests. What share of the month's error budget is gone after one hour?

3.3Alerts on burn rate

Now we can decide when to wake someone. An alert that wakes the on-call engineer is called a page. The obvious first rule is a fixed threshold on the error rate, something like "5xx above 1% for five minutes". It goes wrong in both directions. A blip that crosses 1% for five minutes burns at a rate of 10 for five minutes, which is only about 0.1% of the month's budget, yet it pages someone at 3 a.m. Meanwhile a slow leak at 0.5% errors never crosses 1%, so it never pages, and in three days it spends half the budget (burn rate 5 for 72 hours). The alert should follow how fast the budget is going, and not how the graph looks.

So we choose how much budget loss should wake someone. The SRE workbook's chapter on alerting on SLOs walks through several schemes and ends on this one, where each alert fires only if the burn rate is high over both a long window and a short one:

SeverityLong windowShort windowBurn rateBudget consumed when it fires
Page1 hour5 minutes14.42%
Page6 hours30 minutes65%
Ticket3 days6 hours110%

A ticket is a task for working hours and wakes nobody. Each budget figure comes from one formula, burn rate × window ÷ SLO period. For the first row, 14.4 × 1 h ÷ 720 h = 2%.

?Why two windows per alert?

A long window stops the alert firing on a blip: 14.4× burn over a whole hour is real. The short window, a twelfth of the long one by the workbook's guideline, makes the alert stop firing soon after the problem is fixed, instead of paging for the rest of the hour while old errors age out of the long window.

Here's how the three alerts behave during one incident:

A bad deploy, as the burn-rate alerts see it
●
◉
Healthy
burn ≈ 0.2
⚡
Fast burn
1 h and 5 m windows
⏱
Slow burn
6 h and 30 m windows
▤
Ticket
3 d and 6 h windows
✓
Resolved
short windows clear
Step 1. The service runs at 0.02% errors against a 0.1% budget: burn rate 0.2. No alert is close to firing.
1 / 5

In Prometheus, a monitoring system that stores counters over time, the fast-burn page is two ratio queries joined with and. In each one, rate(...[1h]) is the per-second average over the last hour, sum adds up all the servers, and dividing the 5xx count by the total gives the error ratio. The threshold is the burn rate times the 0.1% budget, 14.4 × 0.001:

promql
(
  sum(rate(http_requests_total{job="api",code=~"5.."}[1h]))
  / sum(rate(http_requests_total{job="api"}[1h])) > (14.4 * 0.001)
)
and
(
  sum(rate(http_requests_total{job="api",code=~"5.."}[5m]))
  / sum(rate(http_requests_total{job="api"}[5m])) > (14.4 * 0.001)
)

The alerts now tell us when the budget is draining. Next we need to know what drains it. Besides services that go down, the most common cause is a service that gets slow, and a slow dependency is worse for its callers than a dead one.

04Timeouts and deadlines

4.1A call with no timeout

A dead service refuses the call at once, and the caller can fail quickly. A slow service takes the call and doesn't answer, so the caller waits. To see what that does, we need to know how a backend handles clicks. Each request being handled occupies a thread, a worker inside the backend program, and a backend has a fixed number of them, its thread pool. When the orders database gets sick and takes ten seconds to answer, every thread that calls it is stuck for ten seconds. Here's that on a pool shrunk to four threads so we can draw it:

A slow database holds every thread; a timeout frees them
Waiting clicksno free thread yetBackend thread pool4 threads, to keep the picture smallOrders databaseclick 1asking dbidleidleidleorders db~20 ms a callclick 5needs a threadclick 6needs a threadclick 7needs a threadhealth checkneeds no dbquery
Step 1. A normal day. A click takes a thread, asks the database, gets the answer in about 20 ms and gives the thread back. Most threads are idle.
1 / 6

A call with no timeout turns a dependency's latency into the caller's resource usage. We can count how much. The number of requests in flight at once equals the arrival rate times the time each takes, a relationship known as Little's Law (chapter 16 covers it). At 500 requests a second, a database that goes from 20 ms to 10 seconds takes the in-flight count from 10 to 5,000. If the pool has 200 threads, the backend is down.

A common extra guard borrows its name from ships, whose hulls are divided into watertight compartments so one leak floods only one compartment. Give each dependency its own separate thread pool, and a sick database can fill only the database pool while calls to cart or auth still find free threads. That separate pool is called a bulkhead, and a bulkhead would also have kept the health check in the scene above from waiting.

?Why not just set every timeout to something generous, like 30 seconds?

Because a generous timeout is almost no timeout. It lets a stuck dependency hold resources for 30 seconds each, which is exactly the pile-up above. A timeout should be close to the latency a healthy dependency has.

A workable rule, roughly: start from the dependency's p99.9 under normal load (the time that 999 of every 1,000 calls beat), add a margin, and check what share of calls would have timed out last week. If a healthy dependency trips the timeout often, it's too tight. If a stuck one can hold a thread for many times its normal latency, it's too loose.

4.2Deadlines across several hops

Our click isn't one call. The frontend calls the backend, which calls the database, and if each hop picks its own timeout, the timeouts don't add up to anything sensible. The frontend gives up after 1 second, but the backend it called is still waiting 2 seconds on the database, doing work nobody will read.

One fix is a deadline: an absolute time by which the whole request must be done, passed down with each call, so that every hop knows how much time is left. The SRE book's chapter on cascading failures gives the arithmetic for a remote procedure call (RPC), a call from one service to another: if server A picks a 30-second deadline and spends 7 seconds before calling B, "the RPC from A to B will have a 23-second deadline."

gRPC, a popular RPC framework, does this for you: a client's deadline travels with the request in a grpc-timeout header, and a server that hands the incoming request's deadline to its own outgoing calls passes on whatever time is left. Over plain HTTP you have to carry the deadline yourself, in a header, and check it before starting expensive work.

A timeout turns a stuck call into a failed one. But many failures last only a moment: a dropped packet, a server restarting, a database replica taking over from a copy that failed. The cheapest answer to a momentary failure is to ask again.

05Retries, backoff and jitter

A retry hides a transient failure at the cost of one more request, which makes it the cheapest reliability feature there is. It's also the most common way a small problem becomes a big one, and the rest of this section is about how.

5.1What's safe to retry

Before how to retry comes whether. Think about the Pay click. The backend asks payments to charge the card, payments charges it, and the response is lost on the way back. The backend sees a timeout and retries, and the customer is charged twice. A retry is only safe if doing the operation twice leaves things the same as doing it once, a property called idempotent.

OperationSafe to retry?Why
GET, readsYesNo side effects
PUT of a full value, DELETEYes, by designIdempotent: the second call leaves the same state
POST that creates or chargesOnly with an idempotency keyThe first call may have succeeded before the response was lost
Call that timed outOnly if idempotentA timeout means "unknown", not "failed"

For the third row, the standard fix is an idempotency key: the client generates a unique ID per logical operation and sends it on every attempt. Stripe's design (Brandur Leach, 2017) has clients send it in an Idempotency-Key header, and when a retry arrives for a key that already succeeded, "the server replies with a cached result of the successful operation". The retried charge then returns the first charge's receipt, and the customer pays once.

5.2Retries multiply

Safe retries still have a problem in a layered system: each layer that retries multiplies the load on everything below it. Our click passes through three callers, the browser's code, the frontend and the backend, before it reaches the database. To draw it, let every layer try at most twice, one try and one retry, and watch what reaches the database when it's overloaded and failing:

One click, three layers that each retry once
Browser codeFrontendBackendOrders databaseoverloaded, so it keeps failing1 clicktry 1try 2try 1try 2try 3try 4try 1try 2try 3try 4try 5try 6try 7try 8
Step 1. The user clicks once. The browser's code will try again if the frontend fails. To keep the picture small, every layer makes at most 2 attempts.
1 / 6

These numbers come from the SRE book's own example: with backend, frontend and JavaScript layers each issuing 3 retries, "a single user action may create 64 attempts (4^3) on the database." Its advice is to "avoid amplifying retries by issuing retries at multiple levels."

?Why not just retry at the top?

Retrying at the top is safest for load, but it's a bit slow: the whole call tree is redone for a failure in one leaf. The usual compromise is to retry at exactly one layer, the one closest to the failure that knows the operation is idempotent, and to have every other layer pass errors straight up.

5.3Backoff, and why it needs jitter

Even with one retrying layer, a client that retries immediately adds load at the worst moment. So retries wait, and wait longer after each failure. This is called backoff. Marc Brooker's Exponential Backoff And Jitter (AWS Architecture Blog, 2015) is the standard reference, and it defines the variants below. In the formulas, base is the first wait, cap is the longest wait allowed, and n is the number of the attempt. Some schemes make the wait partly random, and that randomness is called jitter:

SchemeSleep before attempt nBehaviour
Exponentialmin(cap, base × 2^n)Spreads retries out in time, but every client waits the same amount
Full jitterrandom(0, min(cap, base × 2^n))Same ceiling, uniformly spread below it
Equal jitterhalf the exponential value, plus random over the other halfNever retries very early
Decorrelated jittermin(cap, random(base, previous × 3))Each client's delay depends on its own last delay

?Why isn't plain exponential backoff enough?

Because clients that failed together retry together. If a thousand clients hit an error at the same instant, exponential backoff moves the whole thousand to 1 second later, then 2, then 4. The spike is postponed, and each wave meets the same capacity limit as the first.

A simulation makes that concrete. A thousand clients send at once to a server that accepts 10 requests per millisecond and rejects the rest. Rejected clients back off with base 1 ms and cap 1 second, and retry until they succeed. Besides the four schemes in the table, it runs "none" (retry after 1 ms, with no backoff) and "grpc" (exponential, with each wait multiplied by a random number between 0.8 and 1.2, the small jitter gRPC uses). The code block shows how each policy picks its sleep, and the table shows how many calls the thousand clients made in total and when the last one finished:

1,000 simultaneous clients, a server taking 10 per ms, six retry policies
python
Python
for c, attempt, last in arrivals_this_ms:     # shuffled
    if accepted < CAP: accepted += 1; continue # CAP = 10 per ms
    exp = min(cap, base * 2 ** attempt)        # base 1 ms, cap 1000 ms
    sleep = {
        "none":   base,
        "exp":    exp,
        "grpc":   exp * uniform(0.8, 1.2),
        "equal":  exp / 2 + uniform(0, exp / 2),
        "full":   uniform(0, exp),
        "decorr": min(cap, uniform(base, last * 3)),
    }[policy]
    schedule(c, attempt + 1, at=now + max(1, round(sleep)))
output
Output
policy   total calls  all done after
none          50,500          99 ms
exp           50,500      90,023 ms
grpc           7,410         531 ms
equal          7,029         230 ms
full           7,226         311 ms
decorr         7,341         189 ms

Ideally, 1,000 calls finish at 100 ms, because that's how fast the server can accept a thousand requests. Retrying every millisecond with no backoff finishes on time but sends 50 times the necessary calls. Unjittered exponential backoff is the worst of both: the same 50,500 calls, and the last client waits 90 seconds, because all thousand stay synchronised and keep colliding. Every jittered variant cuts the work to about 7,000 calls and finishes in under a second.

Brooker's own simulation, of clients competing to update the same record, reaches the same verdict: jittered backoff "should be considered a standard approach for remote clients," with full jitter using the least work and decorrelated jitter finishing slightly sooner.

5.4Retry budgets

Backoff controls when each client retries. It doesn't limit how much retrying a fleet of clients does in total, and under a real overload that total is what matters. A retry budget caps retries as a share of normal traffic. Three real designs are in the table below. One of them belongs to Envoy, a proxy: a program that sits between callers and a service and forwards their requests, so it can count and limit them on the way through. We'll meet Envoy again in section 6.

SystemRuleDefault
Google (SRE book, Handling Overload)Up to 3 attempts per request, and retry only while retries are under 10% of a client's requests10% per client
Envoy retry_budgetRetries in flight ≤ a percentage of the requests in flight plus those waiting20%, minimum 3 concurrent
gRPC retry throttlingA counter of tokens: failures cost 1, successes earn tokenRatio; no retries at or below half fullSet per service config

gRPC's version is small enough to read whole. Its counter is a token bucket: a supply of tokens that failures spend and successes refill. Each failed attempt calls throttle(), which removes one token and answers "throttle this retry?" with yes whenever the tokens are at or below the threshold, half the maximum. Each success calls successfulRPC(), which adds ratio tokens back, up to the maximum.

Go
type retryThrottler struct {
	max    float64
	thresh float64   // set to max / 2 from the service config
	ratio  float64
 
	mu     sync.Mutex
	tokens float64
}
 
// throttle subtracts a retry token from the pool and returns whether a retry
// should be throttled (disallowed) based upon the retry throttling policy in
// the service config.
func (rt *retryThrottler) throttle() bool {
	/* ... nil check ... */
	rt.mu.Lock()
	defer rt.mu.Unlock()
	rt.tokens--
	if rt.tokens < 0 {
		rt.tokens = 0
	}
	return rt.tokens <= rt.thresh
}
 
func (rt *retryThrottler) successfulRPC() {
	/* ... nil check ... */
	rt.mu.Lock()
	defer rt.mu.Unlock()
	rt.tokens += rt.ratio
	if rt.tokens > rt.max {
		rt.tokens = rt.max
	}
}

?Why does a budget beat a per-request retry limit?

Because a per-request limit of 3 still allows 4× load when everything is failing, which is exactly when the backend can least afford it. A budget scales retries with success: when most calls succeed, the odd failure gets retried. When most fail, the tokens drain and retries stop, and the backend sees roughly its normal load.

A budget stops us from making a failing dependency worse. It doesn't spare us the cost of calling one that is clearly down. While payments is dead, every click still sends its first call there and waits out a timeout, and the backend's threads keep filling up. We need something that notices the dependency is down and stops calling it.

06Circuit breakers

6.1Three states

A circuit breaker works like the one in a house: when a dependency is clearly failing, it trips, stops calling it for a while and fails fast instead. The pattern comes from Michael Nygard's Release It! and Martin Fowler's CircuitBreaker write-up. A breaker is always in one of three states. In the closed state, the circuit is complete and calls flow through as normal. In the open state it has tripped and calls fail at once. In the half-open state it lets a few trial calls through to find out whether the dependency has recovered.

A row of miniature circuit breakers clipped onto a rail inside an open electrical cabinet, each with a toggle and a numbered label
Circuit breakers in a distribution board. Each one cuts its own circuit when the current stays too high and stays off until someone resets it, so one faulty appliance can't overheat the wiring for the whole building. The software version keeps the idea and replaces the person with a timer and a few trial calls.Photo: Santeri Viinamäki, CC BY-SA 4.0, via Wikimedia Commons

Resilience4j, a common circuit-breaker library for programs on the Java platform, has these three states plus some special ones, and its documented defaults make the state machine concrete. Here's a breaker on default settings sitting between the backend and the payments service:

A Resilience4j circuit breaker on default settings
BackendCircuit breakerdefault settingsPaymentsWhat the breaker has seenCLOSEDcalls flowpaymentshealthylast 100 calls2 failedcall
Step 1. CLOSED: calls flow through. The breaker records outcomes in a count-based window of the last 100 calls. It won't judge anything until it has seen at least 100 (minimumNumberOfCalls).
1 / 6

?Why fail fast instead of waiting for the timeout?

Two reasons. For the backend, a failed call that returns in microseconds frees the thread, where a call that waits out a 2-second timeout holds it. For the dependency, zero traffic is often what it needs to recover: an overloaded database with its queue drained can come back, and one still receiving retries can't.

6.2Envoy's circuit breakers are concurrency limits

Envoy uses the same name for something different. Envoy calls the service it forwards requests to the upstream, and the group of servers that together provide one upstream service a cluster. Its "circuit breakers" are limits on how many connections and requests can be in flight to one cluster at once. They have no states to move through: each request that would push a count over its limit fails, and the next request that fits goes through. Here are the limits and their defaults, starting with the retry budget from section 5.4:

api/envoy/config/cluster/v3/circuit_breaker.proto
envoyproxy/envoy @ v1.31.0 ↗
protobuf
    message RetryBudget {
      // Specifies the limit on concurrent retries as a percentage of the sum of active requests and
      // active pending requests. [...]
      // This parameter is optional. Defaults to 20%.
      type.v3.Percent budget_percent = 1;
 
      // [...] This parameter is optional. Defaults to 3.
      google.protobuf.UInt32Value min_retry_concurrency = 2;
    }
    /* ... */
    // The maximum number of connections that Envoy will make to the upstream
    // cluster. If not specified, the default is 1024.
    google.protobuf.UInt32Value max_connections = 2;
 
    // The maximum number of pending requests that Envoy will allow to the
    // upstream cluster. If not specified, the default is 1024.
    google.protobuf.UInt32Value max_pending_requests = 3;
 
    // The maximum number of parallel requests that Envoy will make to the
    // upstream cluster. If not specified, the default is 1024.
    google.protobuf.UInt32Value max_requests = 4;
 
    // The maximum number of parallel retries that Envoy will allow to the
    // upstream cluster. If not specified, the default is 3.
    google.protobuf.UInt32Value max_retries = 5;

When a limit is hit, Envoy fails the request at once and increments a counter such as upstream_rq_pending_overflow or upstream_rq_retry_overflow. For HTTP requests the router also sets the x-envoy-overloaded header, per its circuit breaking docs.

Envoy puts the state-machine behaviour in a separate feature, outlier detection, which ejects (stops sending traffic to) individual hosts, never the whole cluster. Its defaults in the same release: eject a host after 5 consecutive 5xx, sweep every 10 s, eject for 30 s multiplied by the number of times it's been ejected, and never eject more than 10% of the cluster (outlier_detection.proto). Side by side:

Resilience4j breakerEnvoy circuit breakersEnvoy outlier detection
UnitOne dependencyOne upstream clusterOne host in a cluster
Trips onFailure or slow-call rateConcurrency above a capConsecutive errors, or success-rate outliers
While trippedAll calls fail fastOnly calls over the cap failThat host gets no traffic
RecoversHalf-open probesAs soon as load dropsAfter the ejection time

6.3What to do when the breaker is open

A breaker that opens and returns an error has only made the failure faster. What the caller does instead is the real design question, and for the Pay click the answer depends on which dependency tripped:

FallbackGood forWatch out for
Serve stale data from a cacheProfiles, catalogues, configHow stale is acceptable, and does anyone know it's stale
Degrade the featureRecommendations, ratings, non-essential widgetsHiding a failure nobody notices for weeks
Fail open: let the request through as if the check had passedRate limiters, feature-flag lookupsSecurity checks must never fail open
Queue for laterWrites that can be delayed: emails, analyticsThe queue needs its own limits
Return an errorPayments, anything that must be correctThis is the honest default for critical paths

Everything so far protects a caller from something it depends on. There's a second direction. The Pay click itself arrives at a server, and sometimes more clicks arrive than the server can handle.

07Overload and load shedding

Timeouts, retries and breakers protect a client from a bad dependency. Load shedding protects a server from its clients: when more work arrives than the server can do, something has to be refused. The only choice is whether you refuse it deliberately and cheaply, or accidentally and expensively.

7.1Why a busy server can serve nobody

Say the backend can handle 100 requests a second and a burst of clicks arrives at 150 a second. The extra work waits in a queue, and the server takes requests from it. In a FIFO queue (first in, first out) the server always takes the oldest. The problem is that clients don't wait forever. A client that gave up after a second has gone, but the server doesn't know that, and the request is still in the queue. Here's a small queue after a burst:

A queue of clicks, FIFO and then with a deadline check
Queueoldest at the frontServerone request at a timeWhat the clients gotclick 1client goneclick 2client goneclick 3client waitingclick 4client waiting
Step 1. The burst has left four clicks in the queue. The clients of clicks 1 and 2 waited past their 1-second timeout and have already left. The server can't tell.
1 / 7

Because the server was busy the whole time, its processor sat at 100% and yet it served nobody. That's why we need two measures. Throughput is the number of requests the server completes per second, whether or not anyone is still waiting for the answer. Goodput is the number it completes within the client's deadline, which is the number that matters.

We can measure the effect with a simulation. One server handles 100 requests a second (each takes a random time averaging 10 ms). Normal load is 80 a second. From t=60 s to t=90 s a burst pushes it to 150 a second, then it returns to 80. Clients give up after 1 second. We try four ways of running the queue and, for each, clients that never retry or that retry up to 3 times when they time out:

  • fifo is the plain queue from the scene above.
  • deadline is FIFO plus the deadline check: a request already past its deadline is dropped when it's dequeued.
  • bounded also rejects a request on arrival when 50 requests are already queued.
  • lifo serves the newest request first whenever the queue has been busy for 100 ms (section 7.3 explains why).
A 30-second overload, four queue policies, with and without 3 retries on timeout
python
Python
def start(t):                       # called whenever the server is free
    while q:
        if policy == "lifo" and t - last_empty > 0.1:
            req = q.pop()           # queue busy for 100 ms: newest first
        else:
            req = q.popleft()
        if policy != "fifo" and t > req.deadline:
            continue                # expired: drop it without doing the work
        serve(req, t)               # takes expovariate(1 / 0.010) seconds
        return
# "bounded" also rejects on arrival when 50 requests are already queued.
# Retrying clients resend immediately when their 1 s timeout fires, up to 3 times.
output
Output
goodput, requests/s        0-30  30-60 60-90 90-120 120-150 150-180 ... 270-300
                                        burst
fifo      no retries          80    81    11     0      0      76  ...    79
fifo      3 retries           80    81    11     0      0       0  ...     0
deadline  no retries          80    81    73    82     79      78  ...    79
deadline  3 retries           80    81    31    48     79      82  ...    79
bounded   3 retries           80    81   101    80     79      79  ...    80
lifo      3 retries           80    81   101    84     79      78  ...    79
 
fifo, 3 retries: offered 535/s during the burst, then ~320/s for the rest of
the run; ~100/s of completed work wasted every window after t=60

Each row is a policy, each column is a 30-second window, and each number is goodput in requests per second. Before the burst every policy delivers about 80, all the offered load. Read the rows one by one:

  • Plain FIFO keeps the server 100% busy through the burst and serves almost nobody for over a minute. The burst leaves about 1,500 requests queued (50 a second over capacity, for 30 seconds), and they drain at only 20 a second of spare capacity: 75 seconds.
  • FIFO with retries never recovers. Every timed-out request comes back as a retry, so offered load, the rate at which requests arrive including retries, settles around 320 a second against a capacity of 100, long after the burst has gone. Goodput is zero for the remaining 210 seconds.
  • One deadline check before starting each request fixes most of it. The server stops working on requests nobody is waiting for.
  • A bounded queue or adaptive LIFO keeps goodput at the server's full capacity during the burst. Rejecting early is cheaper than timing out late.

Look again at the second row, because it's the alarming one: the burst ended after 30 seconds, and the outage lasted for the rest of the run. That behaviour has a name.

7.2Metastable failure

A burst that ends should leave a system that recovers. The FIFO-with-retries row didn't. Bronson, Aghayev, Charapko and Zhu's Metastable Failures in Distributed Systems (HotOS 2021) calls this a metastable failure: "a trigger causes the system to enter a bad state that persists even when the trigger is removed," held there by "a sustaining effect—often involving work amplification". Leaving it "requires a strong corrective push, such as rebooting the system or dramatically reducing the load." In our simulation the trigger was the 30-second burst, and the sustaining effect was the retries.

The life of a metastable failure, from the simulation
●
◉
Stable
80% busy
◌
Vulnerable
retries configured
⚡
Trigger
30 s burst
↻
Sustained
retries > capacity
⏏
Push
shed or restart
Step 1. The server runs at 80 requests a second against a capacity of 100. Latency is fine and every request succeeds.
1 / 6

Perhaps the best-known real case is the AWS us-east-1 outage of 7 December 2021, which had this shape. Per AWS's summary, an automated scaling activity "triggered an unexpected behavior from a large number of clients," the surge overwhelmed networking devices between two internal networks, and the delays caused "even more connection attempts and retries," leading to "persistent congestion." The clients did have well-tested back-off behaviour, but "a latent issue prevented these clients from adequately backing off during this event." Recovery took about seven hours.

Now we know what to prevent. The remaining question is where in the path of a click to refuse the extra work, since the later it's refused, the more it has already cost.

7.3Shed early, shed cheaply

A request is cheapest to refuse at the first point it can be refused. Here are the places, from the client inwards:

Where an overloaded request can be refused, cheapest first
●
◉
Client
throttles itself
⇆
Proxy
Envoy limits
▣
Admission
concurrency limit
≡
Queue
bounded wait
⚙
Handler
deadline check
Step 1. Client-side throttling. The client counts how many requests it sent and how many the backend accepted, and starts refusing its own requests once it sends more than about twice what gets accepted. The backend never sees this traffic.
1 / 5

The first and fourth steps need a closer look. Client-side throttling comes from Google's Handling Overload. Each client keeps two counts over the last couple of minutes, requests (attempted) and accepts (accepted by the backend), and rejects a new request locally with probability

C++
max(0, (requests − K × accepts) / (requests + 1))

K is a multiplier the operator picks, and Google says it "generally prefer[s] the 2x multiplier". When the backend accepts everything, requests is close to accepts, the top of the fraction is negative, and the probability is zero. When the backend starts rejecting, accepts falls behind and the probability rises, so clients cut their own traffic and the backend spends less of its capacity saying no.

The two queue rules from the fourth step come from Ben Maurer's Fail at Scale (ACM Queue, 2015), on how Facebook runs its services. Its CoDel variant (short for "controlled delay") says: if the queue hasn't been empty in the last N milliseconds, limit time in queue to M milliseconds. Facebook found M = 5 ms and N = 100 ms "tends to work well across a wide set of use cases." Adaptive LIFO serves FIFO normally and switches to newest-first (LIFO, last in first out) once a queue forms, because the oldest request's user has probably given up. In the simulation, the lifo policy uses the same 100 ms rule.

?Why LIFO? Isn't it unfair?

It's unfair to the requests at the back of a queue that's already too long, but they were probably going to time out anyway. Under FIFO everyone waits and most fail. Under LIFO the newest requests, whose clients are still waiting, succeed. Fairness in a queue that's overflowing means spreading failure evenly, and nobody benefits from that.

7.4Not all traffic is equal

Shedding blindly drops the checkout request as readily as the prefetch. Google's services tag every request with a criticality, one of CRITICAL_PLUS, CRITICAL, SHEDDABLE_PLUS and SHEDDABLE (Handling Overload), and shed from the bottom up. Even two classes help:

ClassExamplesUnder overload
User-facing, criticalLogin, checkout, the page itselfShed last
User-facing, optionalRecommendations, counts, previewsShed early; degrade the UI
Batch and backgroundReindexing, reports, prefetchShed first; retry later
RetriesAny retried requestShed before first attempts of the same class

Priority has to travel with the request the same way the deadline does. A background job that calls an API that calls a database should arrive at the database marked as background, or the database can't tell it from a user waiting on a page.

A paper triage tag with tear-off coloured strips at the bottom labelled Morgue, Immediate, Delayed and Minor
A mass-casualty triage tag. The first responder tears off strips until the bottom one shows the patient's class (Immediate, Delayed or Minor), and the tag stays tied to the patient, so everyone downstream treats people in that order without examining them all again. A request's criticality works the same way: decided once, carried with it, read at every hop.Photo: Aleichem, CC BY-SA 3.0, via Wikimedia Commons

With the click's defences in place, from timeouts to shedding, what's left is to apply them to a real service in a sensible order.

08Applying it to your service

8.1A checklist, in order

Each step follows from a section above, in the order the chapter built them:

  1. Write down the SLI and SLO for your most important user journey, counted at the load balancer. Check it against the product of your critical dependencies (section 2.3).
  2. Add the two paging burn-rate alerts (14.4× over 1 h and 5 m; 6× over 6 h and 30 m) and delete the fixed-threshold error alert (section 3.3).
  3. Put a timeout on every outbound call, based on the dependency's p99.9, and propagate a deadline where your framework supports it (section 4).
  4. Retry at one layer, with full or decorrelated jitter, a cap, and a budget. Make every retried write idempotent (section 5).
  5. Bound every queue, and check the deadline when you dequeue (section 7).
  6. Size proxy concurrency limits from arrival rate × latency, not defaults (section 6.2).
  7. Run the recovery test from section 7.2: overload for a minute, drop back, and confirm the service recovers without help.

8.2What you trade for what

You getYou payWhen the bill arrives
A target the whole team can read off a counterA tighter SLO means a different engineering organisation for each nineWhen the budget runs out and releases slow down
Timeouts that free threadsSome calls fail that would have succeeded slowlyWhen a timeout is set too tight and healthy calls trip it
Retries that hide brief failuresExtra load on the dependency, multiplied by every retrying layerIn an overload, as 64× the traffic
A breaker that lets a dependency recoverA fallback you have to design and keep honestWhen a degraded feature goes unnoticed for weeks
Load shedding that protects the serverSome users get errors on purposeWhen the shed class turns out to be one that mattered

8.3Symptom, cause, fix

SymptomLikely causeFix
Stays down after the trigger is gone; CPU pinned, goodput near zeroMetastable: retries or queued expired work sustain the overloadShed load or drain the queue; then add retry budgets and deadline checks
One slow dependency exhausts all threadsMissing or overlong timeoutTimeout from p99.9, bulkhead the pool per dependency
Database load spikes to many times normal during an incidentRetries at several layersRetry at one layer, add a retry budget
Clients all retry in wavesBackoff without jitterFull or decorrelated jitter
503s with x-envoy-overloaded on a normal dayEnvoy circuit-breaker defaults too lowSize max_requests and max_pending_requests from Little's Law
Duplicate charges or orders after timeoutsNon-idempotent write retriedIdempotency keys
Pages at night for blips; slow leaks never pageFixed-threshold error alertsMultiwindow burn-rate alerts
Breaker flaps open and closedWindow too small, or half-open lets in too much trafficLarger window and minimum calls; fewer half-open probes

09Summary

  1. Dependencies multiply. A click that needs five 99.9% services works about 99.5% of the time, and ten dependencies leave it near 99%. Critical dependencies need an extra nine, or a fallback.
  2. Each nine is ten times harder. 99.99% leaves 4.3 minutes a month, which has to be handled by automation, not people.
  3. An SLO is a count of good events over a window. Measure it where users feel it, and set the SLA looser than the SLO.
  4. The error budget is meant to be spent. It decides when to ship and when to fix, without an argument.
  5. Alert on burn rate over two windows. 14.4× over an hour spends 2% of a month's budget and deserves a page; a fixed error threshold doesn't.
  6. Every call needs a timeout, and call trees need a deadline. A call with none turns a dependency's latency into your threads. Check the deadline before doing work, not only while waiting.
  7. Retries multiply across layers. Four attempts at three layers is 64 attempts at the bottom. Retry at one layer, within a budget, and only operations that are idempotent.
  8. Backoff without jitter postpones the spike. In the simulation it took the last client 90 seconds; full jitter took 0.3.
  9. Circuit breakers fail fast so the dependency can recover. Envoy's "circuit breakers" are concurrency caps; its outlier detection is the state machine.
  10. Overload is survived by refusing work early. A bounded queue, a deadline check or adaptive LIFO kept goodput at capacity where plain FIFO with retries stayed at zero.
  11. A failure can outlive its trigger. Retries that sustain an overload make it metastable, so test recovery and not only normal load.

10Build this

A metastability test harness.

  • Write a small HTTP service with a fixed capacity (a worker pool of 4 and a handler that sleeps 10 ms) and a load generator that sends at a fixed rate, times out at 1 second and optionally retries.
  • Reproduce section 7.1: normal load at 80% of capacity, a 30-second burst at 150%, then back to 80%. Plot goodput per second. Confirm that FIFO plus retries never recovers.
  • Add the fixes one at a time: a deadline check at dequeue, a bounded queue, adaptive LIFO, a client retry budget. Measure how long each takes to recover.
  • Put Envoy in front with default circuit breakers, then with limits sized from Little's Law, and compare the overflow counters.

11Interview questions

beginnerWhat's the difference between an SLI, an SLO and an SLA?›

The SLI is the measurement: good events over valid events, such as the share of requests that succeed in under 300 ms. The SLO is the internal target for it over a window, such as 99.9% over 30 days. The SLA is an external contract, usually looser than the SLO, with penalties if it's missed.

beginnerWhy add jitter to exponential backoff?›

Because clients that fail together retry together. Without jitter, a thousand clients that failed at the same moment all retry at 1 s, then 2 s, then 4 s, and each wave hits the same limit. Jitter spreads them out. In a simulation of 1,000 clients against a server taking 10 per millisecond, unjittered backoff took 90 seconds to serve everyone; full jitter took 0.3 seconds with a seventh of the calls.

intermediateWhat is a burn rate, and why alert on it?›

The ratio of your current error rate to the rate that would use up the error budget exactly at the end of the SLO window. For a 99.9% SLO, a 1.44% error rate is a burn rate of 14.4, which spends 2% of a 30-day budget in an hour. Alerting on burn rate ties the page to the thing you care about, how fast you're using up reliability, instead of to an arbitrary error threshold.

intermediateA request goes through browser, frontend, backend and database, and each of the first three retries 3 times. What's the worst case at the database?›

Each layer makes up to 4 attempts, and the attempts multiply: 4 × 4 × 4 = 64 attempts for one user action. That's the SRE book's own example. The fix is to retry at only one layer and cap retries with a budget, so a failing database sees close to its normal load instead of 64 times it.

intermediateWhat does a circuit breaker do in each of its states?›

Closed: calls pass, and outcomes are recorded in a sliding window. When the failure rate crosses a threshold (50% by Resilience4j's default, after at least 100 calls), it opens. Open: calls fail immediately without reaching the dependency, for a wait period (60 s by default). Half-open: a small number of trial calls (10) are allowed; if they succeed it closes, otherwise it reopens.

deepWhat's a metastable failure, and how do you prevent one?›

A failure where a temporary trigger pushes the system into a bad state that persists after the trigger is gone, because of a sustaining effect such as retries or work on already-expired requests. Load that the system handled before the trigger is now more than it can handle. Prevention is removing the sustaining effect in advance: retry budgets, jitter, deadline checks at dequeue, bounded queues and load shedding. And test recovery by overloading and then returning to normal load.

deepWhy would a server process its queue LIFO under load?›

Once a queue has formed, the oldest requests have waited longest, and their clients have probably timed out. Serving them first spends capacity on work nobody will read, so everyone fails. Serving newest first gives recent requests, whose clients are still waiting, a good chance of success. Facebook uses adaptive LIFO (FIFO normally, LIFO when a queue forms) together with a CoDel-style bound on queueing time.

deepEnvoy returns 503s with x-envoy-overloaded on a normal day. What's happening?›

A cluster circuit breaker is overflowing, most likely max_requests or max_pending_requests, both 1,024 by default. By Little's Law a service at 20,000 requests a second and 60 ms latency has about 1,200 requests in flight, so the default is too small. Check the upstream_rq_pending_overflow and upstream_rq_active_overflow counters, then size the limits from arrival rate × latency with headroom.

12Go deeper

check yourself
Your SLO is 99.95% per 30 days. How much full outage can you afford?›

21.6 minutes: 0.05% of 43,200 minutes.

A client library retries with min(cap, base × 2^n) and no randomness. What goes wrong?›

Clients that failed together stay synchronised and retry in waves, each wave meeting the same capacity limit. Add full or decorrelated jitter.

gRPC retry throttling: when does a client stop retrying?›

When its token count falls to half of maxTokens or below. Each failure costs one token and each success earns tokenRatio.

Why check the deadline when dequeuing a request, not only when it arrives?›

Because it may have expired while it waited. Processing it spends capacity on a response nobody will read, which is how an overload sustains itself.

Operating Systems: Three Easy Pieces, Distributed Systems and Sun's NFS

Distributed systems from the ground up: unreliable networks, acknowledgements, timeout and retry, and why NFS makes its operations idempotent. Free online

Google SRE Workbook: Alerting on SLOs

Walks through six alerting schemes, from a plain error threshold to multiwindow multi-burn-rate, with the trade-offs of each. sre.google

The source of the burn-rate table in section 3.3.
Google SRE Book: Handling Overload and Addressing Cascading Failures

Client-side throttling, criticality, retry budgets, deadline propagation and the 64× retry example. Handling Overload, Cascading Failures.

Bronson et al., Metastable Failures in Distributed Systems (HotOS 2021)

The vocabulary for failures that outlive their trigger: vulnerable state, trigger, sustaining effect. PDF

Read it before your next load test.
Marc Brooker, Exponential Backoff And Jitter (2015)

The simulation that made full jitter the default recommendation. AWS Architecture Blog

Ben Maurer, Fail at Scale (ACM Queue, 2015)

How Facebook runs services under overload: controlled delay, adaptive LIFO and concurrency control. CACM

Treynor, Dahlin, Rau and Beyer, The Calculus of Service Availability (2017)

The rule of the extra nine and how to reason about critical dependencies. CACM

envoy: circuit_breaker.proto and outlier_detection.proto

Every limit and default in section 6.2, documented in the API definitions. v1.31.0

gRFC A6: gRPC client retries

Retry policy, hedging policy and retry throttling, as specified. proposal

Contention, Queueing & Tail Latency

The queueing arithmetic under every section here: the utilisation curve, Little's Law, hedged requests and coordinated omission. Chapter 16.

Observability

How the SLIs in this chapter get measured: counters, histograms, the Prometheus data model and traces. Chapter 39.

The System Design Method

Where SLOs and capacity numbers enter a design, before the boxes are drawn. Chapter 38.

The Linux Networking Stack

The accept queue and SYN retries: load shedding the kernel does for you, silently. Chapter 10.