Why does p99 latency explode long before 100 % CPU?
Short answer: a request waits whenever every core is busy at the moment it arrives, and that wait grows with 1 / (1 − utilization), not with utilization. On one core, a request that needs 10 ms of work has a p99 of 92 ms at 50 % CPU, 230 ms at 80 % and 921 ms at 95 %, while the CPU graph still shows room.
More cores move the bend to the right and make it sharper; bursty traffic moves it to the left. We saw it on a real Postgres: p99 was 3.5 ms at 60 % CPU, 12.6 ms at 88 % and 21 ms at 94 %. Plan small machines for 60–70 %, and count a fleet of small machines as small machines.
A CPU graph is an average over a minute. A request does not meet the average: it meets the machine in the one millisecond in which it arrives. At 80 % CPU the machine is idle a fifth of the time and busy the rest, and a request that finds it busy waits for the work in front of it. Throughput does not show this (everything is still served), the mean barely shows it, and p99, the latency of the unluckiest request in a hundred, shows it first.
The arithmetic for one core
Utilization is arrival rate times the work per request: 80 requests a second at 10 ms each keep one core busy 80 % of the time. If requests arrive at random moments (independent users do) and their work varies around its mean the way an exponential distribution does, queueing theory gives two short formulas for one core. With S the mean work per request and ρ the utilization:
mean response time = S / (1 − ρ)
p99 response time = ln(100) × S / (1 − ρ) ≈ 4.6 × S / (1 − ρ)
Nothing in them is special about 100 %. What counts is the idle share, 1 − ρ, and every halving of the idle share doubles the latency: from 50 % to 75 % CPU, from 75 % to 87.5 %, from 87.5 % to 93.75 %. The first doubling takes 25 points of CPU, the third takes six. On a dashboard that looks like a wall; in the formula it is the same curve all along.
More cores move the bend, they do not remove it
With several cores a request waits only if all of them are busy, which is rarer. At 70 % load that happens to 70 % of the requests on one core, to 58 % on two cores and to 13 % on sixteen (Erlang's formula for a queue with several servers). The table is the p99 of the same 10 ms request:
| CPU utilization | 1 core | 2 cores | 16 cores |
|---|---|---|---|
| 50 % | 92 ms | 57 ms | 46 ms |
| 70 % | 154 ms | 83 ms | 46 ms |
| 80 % | 230 ms | 119 ms | 47 ms |
| 90 % | 461 ms | 233 ms | 53 ms |
| 95 % | 921 ms | 463 ms | 72 ms |
Two things to read from it. A large machine can run much hotter than a small one before the tail moves: sixteen cores at 90 % are better off than one core at 50 %. And the large machine gives less warning: its p99 looks perfect at 85 %, and the ten points from there to 95 % add more latency than the first 85 did.
Bursts count as much as the average
The formulas above assume random arrivals and work that varies. Kingman's approximation for one core says what happens when either is different:
mean wait ≈ ρ / (1 − ρ) × (Ca² + Cs²) / 2 × S
Ca and Cs are the variability of the gaps between arrivals and of the work per request (standard deviation divided by mean; 1 for random arrivals, 0 for a metronome). Variability multiplies the wait. Evenly spaced arrivals halve it; arrivals in clumps (a cron job on the hour, a client that retries, a page that fires twenty calls at once) have a Ca far above 1 and multiply it. The same holds for the work: one slow endpoint among fast ones raises Cs for everybody who queues behind it.
A measurement: Postgres on 2 vCPU
We measured one Postgres on a 2-vCPU cloud VM under an open-loop load in steps of 60 seconds, once with random arrivals and once with evenly spaced ones. The full setup and all steps are in How many requests per second can one Postgres handle?; these are seven of its ten steps:
| Offered req/s | DB CPU (random) | p99 random | DB CPU (even) | p99 even |
|---|---|---|---|---|
| 640 | 9 % | 1.4 ms | 8 % | 1.1 ms |
| 3,200 | 49 % | 2.5 ms | 35 % | 1.1 ms |
| 4,480 | 70 % | 5.7 ms | 58 % | 1.6 ms |
| 5,120 | 80 % | 7.3 ms | 73 % | 3.3 ms |
| 5,440 | 88 % | 12.6 ms | 83 % | 4.2 ms |
| 5,760 | 94 % | 21.0 ms | 91 % | 6.0 ms |
| 6,080 | 99 % | 1,966 ms | 99 % | 1,572 ms |
With random arrivals p99 was four times its idle value at 70 % CPU and fifteen times at 94 %, with every request still served. From half load on, the same requests per second with evenly spaced arrivals gave a p99 two to four times lower: the direction Kingman's formula gives. Past the limit the queue no longer drains and the latency is whatever the length of the overload makes it.
The same curve in the live model
Four designs you can open in Stackrig, each one service whose requests need 10 ms of CPU on
average: one small server with 2 cores; eight of those small servers with nothing in front, so
that each request reaches a server picked at random; the same eight behind a load balancer that
deals the requests in turn; and one large server with 16 cores (AWS c7g.large and
c7g.4xlarge, where one vCPU is a full core). Eight small servers have the same
sixteen cores as the large one at the same price; the load balancer costs extra. Each design
runs at 50 % CPU at 1× traffic; the traffic slider at 1.4×, 1.6×, 1.8× and 1.9× gives the other
rows.
| CPU utilization | 1 small server | 8 small servers, picked at random | 8 small servers behind a load balancer | 1 large server |
|---|---|---|---|---|
| 50 % | about 60 ms | about 60 ms | about 55 ms | about 50 ms |
| 70 % | about 85 ms | about 85 ms | about 65 ms | about 50 ms |
| 80 % | about 125 ms | about 125 ms | about 85 ms | about 50 ms |
| 90 % | about 230 ms | about 230 ms | about 150 ms | about 50 ms |
| 95 % | about 460 ms | about 460 ms | about 270 ms | about 75 ms |
The two middle columns are the ones to remember. A request that reaches one of eight small servers at random finds what it finds at one small server: each server has its own queue, and a request that waits at a busy server cannot use an idle core on the server next to it. A load balancer that deals the requests in turn spaces each server's arrivals more evenly, and the tail is about a third lower at 90 % CPU. It is still about three times the tail of one large server with the same sixteen cores. A fleet of small servers behaves like small servers, not like one large server.
These columns are the model's. Our own measurement of fleets on real machines is the subject of a coming guide.
Have an invite? Open the four designs:
- One small server (2 cores)
- Eight small servers, picked at random
- Eight small servers behind a load balancer
- One large server (16 cores)
What to do about it
- Plan by utilization, per tier. For a tier on machines with few cores, 60–70 % at the daily peak is where p99 is still close to its idle value. Large machines can run hotter, but leave room: their bend is sharper.
- Measure your own curve. Raise the load in steps and write down CPU and p99 at each step, with a load generator that keeps sending on schedule when the system slows down (open loop). A generator that waits for each answer slows down with the system and hides the bend.
- Smooth the arrivals. Spread scheduled jobs over the minute instead of the full hour, put a queue in front of batch work, and give retries a backoff with jitter: see Retry storms: why retries take your service down.
- Pool where you can. One queue in front of many workers beats many queues with one worker each. A connection pool in front of a database is such a shared queue: see How big should your database connection pool be? Fewer, larger machines pool better than many small ones, at the price of losing more when one fails.
- Keep slow work away from fast work. A report that takes two seconds in the same pool as a lookup that takes two milliseconds makes the lookup's p99 the report's problem.
What this page does not tell you
- The formulas are a model. They assume random arrivals, exponentially distributed work and a queue that never gives up. Real work is often less variable (lower latencies than the table) and real arrivals often burstier (higher).
- Only the wait for a core. Garbage collection, locks, disk and network waits, a full connection pool and noisy neighbours add their own tails on top.
- vCPUs are not always cores. Two hyperthreads of one core do less than two cores, so on many machine types the bend comes earlier than the vCPU count suggests.
- The live model is not a measurement. It follows the arithmetic. In the one comparison with a real machine that we have published, the $12 server, the model found the knee within 5 %, and past the knee its p99 was 1.7–2× too low. How accurate is Stackrig? answers a different question: how closely the live model agrees with an exact simulation.
- One measurement. The Postgres numbers are one workload on one machine type on one day.
Sources
- J. F. C. Kingman, The single server queue in heavy traffic, Mathematical Proceedings of the Cambridge Philosophical Society 57 (1961): the approximation for the wait with general arrivals and work.
- A. K. Erlang, Solution of some Problems in the Theory of Probabilities of Significance in Automatic Telephone Exchanges, Elektroteknikeren 13 (1917): the share of arrivals that must wait when several servers share a queue.
- Mor Harchol-Balter, Performance Modeling and Design of Computer Systems, Cambridge University Press (2013): the derivations of the response-time formulas used here.
- Jeffrey Dean and Luiz André Barroso, The Tail at Scale, Communications of the ACM 56 (2013): why the slowest hundredth decides how a service with many parts feels.
Open “Web app with database” in the playground
The playground is invite-only during the private preview: join the waitlist to get an invite.