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:

Calculated, not measured: p99 response time of a request with 10 ms mean work, random (Poisson) arrivals, exponentially distributed work, one shared queue in front of the cores (the M/M/c queue). With no waiting at all the p99 is 46 ms, the slowest hundredth of the work itself.
CPU utilization1 core2 cores16 cores
50 %92 ms57 ms46 ms
70 %154 ms83 ms46 ms
80 %230 ms119 ms47 ms
90 %461 ms233 ms53 ms
95 %921 ms463 ms72 ms
Calculated p99 response time against CPU utilization for a 10 ms request: on one core it rises from 92 ms at 50 % to 230 ms at 80 % and passes 500 ms at 91 %; on sixteen cores it stays near 46 ms up to 85 % and then rises steeply, to about 300 ms at 99 %. 0 50 % 80 % 100 % 0 250 500 ms 1 core 16 cores
p99 against CPU utilization: the same formulas as the table, drawn. Sixteen cores stay flat for longer and then rise faster: at 99 % they are at about 300 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:

Measured 2026-09-30 on Google Cloud n2-standard-2 (2 vCPU = two hyperthreads of one core). p99 of the whole HTTP request: the application, two network hops and one query. At the last step the offered load exceeded what the database can serve.
Offered req/sDB CPU (random)p99 randomDB CPU (even)p99 even
6409 %1.4 ms8 %1.1 ms
3,20049 %2.5 ms35 %1.1 ms
4,48070 %5.7 ms58 %1.6 ms
5,12080 %7.3 ms73 %3.3 ms
5,44088 %12.6 ms83 %4.2 ms
5,76094 %21.0 ms91 %6.0 ms
6,08099 %1,966 ms99 %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.

From Stackrig's live model, not from real machines: p99 of the median second out of 20 after 100 seconds at each load, rounded. In the columns with random arrivals (all but the load balancer's) the model stays within 15 % of the formula above for 2 and for 16 cores.
CPU utilization1 small server8 small servers, picked at random8 small servers behind a load balancer1 large server
50 %about 60 msabout 60 msabout 55 msabout 50 ms
70 %about 85 msabout 85 msabout 65 msabout 50 ms
80 %about 125 msabout 125 msabout 85 msabout 50 ms
90 %about 230 msabout 230 msabout 150 msabout 50 ms
95 %about 460 msabout 460 msabout 270 msabout 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:

What to do about it

  1. 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.
  2. 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.
  3. 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.
  4. 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.
  5. 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

Sources

Open “Web app with database” in the playground

The playground is invite-only during the private preview: join the waitlist to get an invite.

More