Part 2 · 10 chapters · ~55 min

Performance: latency, throughput, and the tail

Two different problems that share a word. A transfer is a latency problem and the accrual batch is a throughput one, and the same optimisation can help one while hurting the other. This part covers the arithmetic that makes a capacity conversation concrete in ten seconds, why the tail rather than the mean is what customers experience, and the four things that actually make software slow.

18

Latency and throughput are different problems

worked numbers
latency    how long ONE operation takes
throughput how many operations per unit time

they are related and not inversely:

  adding servers       → more throughput, same latency
  batching             → more throughput, worse latency
  caching              → better both
  more threads         → more throughput until contention, then worse both

"make it faster" is ambiguous, and the two answers differ.
why it matters for a bank specifically
A transfer is a latency problem; the accrual batch is a throughput problem. The transfer needs a 300 ms p99 and does 464 per second; the batch needs to finish 2 million loans in a window and nobody cares about any individual one. The same optimisation can help one and hurt the other, which is why Part 5 batched the accrual and Part 1 did not batch the posting.
19

Little’s Law, and using it in a room

worked numbers
L = λW

  L = items in the system   (concurrency)
  λ = arrival rate          (throughput)
  W = time in the system   (latency)

it holds for any stable system, with no assumptions about
distribution. that is what makes it usable in an interview.
three ways to use it
  1. Sizing a pool. 464 transfers per second at 80 ms each needs L = 464 × 0.08 = 37 concurrent. So around 40 connections, not 200, and not 5.
  2. Predicting the effect of a slowdown. If the ledger call goes from 80 ms to 400 ms at constant arrival rate, required concurrency goes from 37 to 186. If the pool is 50, you now have a queue and the latency is worse than 400 ms.
  3. Sanity-checking a claim. "We handle 10,000 requests per second with 20 threads at 50 ms" implies L = 500 against 20 available. The numbers are wrong and now you know which one to question.
why this is the most useful thing here
It turns a vague performance conversation into arithmetic in ten seconds, it requires no tooling, and it is the fastest way to find out whether someone has thought about capacity or is reciting. The second use is the one that predicts incidents: a dependency slowing down silently multiplies your concurrency requirement, and that is the mechanism behind the cloud module's autoscaling feedback loop.
20

Why the tail matters more than the mean

worked numbers
a page that makes 20 backend calls, each p99 = 1 second:

  P(all 20 fast) = 0.99²⁰ = 0.82
  → 18% of page loads hit at least one slow call

with 100 calls: 63%

the p99 of one service becomes the p50 of a page that fans out.
this is why tail latency is amplified by architecture.
what the tail actually is
  1. The mean tells you almost nothing. A bimodal distribution of 10 ms cache hits and 2 s cache misses has a comfortable mean and a terrible experience.
  2. p99 is one customer in a hundred. At 18,000 balance reads a second, that is 180 people per second having a bad time.
  3. p99.9 is where the interesting bugs live: GC pauses, lock contention, a cold cache, a retry, a failover.
  4. Your heaviest users hit the tail most. They make more requests, so their probability of encountering it approaches one.
the practice
Alert on p99, budget on p99, and look at p99.9 when investigating. And never average percentiles across instances or time windows, which is mathematically meaningless and extremely common in dashboards. Part 13 of the CBA module specified distributions rather than averages for exactly this reason.
21

Where latency actually goes: a budget, itemised

Part 8 of the CBA module budgeted card authorisation. The same exercise for a transfer, which is the more instructive one because most of it is not compute.

Stepp50p99Note
TLS and gateway3 ms12 msConnection reuse matters enormously here
Auth and validation2 ms6 ms
Limit check (Redis)1 ms4 msOne round trip, one Lua script
Fraud score18 ms150 msHard timeout, fails open
Shard routing<1 ms1 msIn-memory map
Posting transaction22 ms95 msDominated by fsync and replica ack
Outbox insert2 ms5 msSame transaction
Response1 ms3 ms
Total~50 ms~276 msAgainst a 300 ms budget
the two observations worth making
The durability cost is visible and was chosen. Part 1 accepted 5 to 15 ms for fsync plus replica acknowledgement, and here it is, roughly a third of the p99. And the fraud timeout dominates the tail while contributing little to p50, which is exactly why it fails open: it is the single largest thing you can remove from the worst case.
22

Queueing theory, the useful ten percent

worked numbers
the one result worth memorising:

  W = S / (1 − ρ)     where ρ = utilisation

  ρ = 0.5 → wait = 2× service time
  ρ = 0.8 → wait = 5×
  ρ = 0.9 → wait = 10×
  ρ = 0.95 → wait = 20×
  ρ = 0.99 → wait = 100×

latency is non-linear in utilisation, and the knee is around 70-80%.
what follows from that curve
  1. Target 60 to 70% utilisation on latency-sensitive paths. Running at 90% to "use the capacity efficiently" means ten times the queueing delay, and no headroom for a spike.
  2. Headroom is a feature. It is what absorbs the 4x peak from Part 3 and the backlog after a failover.
  3. Variance makes it worse. The formula assumes exponential service times; a bimodal distribution queues harder at the same utilisation.
  4. This is why autoscaling on CPU is dangerous. By the time CPU shows saturation, you are past the knee and latency has already collapsed.
23

The four things that make software slow

CauseSignatureTypical fix
Doing too muchLinear in input size, CPU-boundDo less. Better algorithm, less data, precompute
WaitingLow CPU, high latencyParallelise, batch, or remove the round trip
ContentionLatency rises with concurrencyReduce the critical section, shard the hot thing
CoordinationLatency floors at network RTTFewer participants, or relax the guarantee
worked numbers
for our bank, in order of how often it is the answer:

  1. waiting           N+1 queries, serial external calls
  2. contention      the hot row, the shared lock, the pool
  3. coordination    fsync, replica ack, cross-shard saga
  4. doing too much  rarely, and usually a missing index

most "slow service" investigations end at waiting or contention,
and almost none end at the CPU.
the diagnostic shortcut
Look at CPU first, not to optimise it but to classify the problem. High CPU means doing too much. Low CPU with high latency means waiting or contention, which is the common case and a completely different investigation. Teams that jump straight to profiling CPU on a waiting problem find nothing and conclude the profiler is broken.
24

Measuring before optimising, properly

the sequence
  1. Reproduce it. A performance problem you cannot reproduce is a performance problem you cannot verify you fixed.
  2. Measure end to end first. Where does the time actually go? Part 13 tracing exists for this, and the answer is frequently not where anyone guessed.
  3. Form a hypothesis with a number. "I think the N+1 in the beneficiary lookup costs 40 ms of the p99." Falsifiable.
  4. Change one thing. Two changes at once and you learn nothing about either.
  5. Measure again, the same way. Same load, same environment, same percentile.
  6. Keep it or revert it. An optimisation that did not measurably help is complexity you added for nothing. Revert it.
the failure this prevents
Optimising the part you understand rather than the part that is slow. Everyone has rewritten a function for 2 ms while a serial external call cost 200. The measurement is not a formality before the real work; it is the work, and the fix is usually obvious once the measurement is honest.
25

Load testing that resembles production

what makes a load test worthless
  1. Uniform data. Every request hitting a different account means no contention and no hot keys, which is the opposite of production where the fee account is on every transaction.
  2. A warm cache and no cold start. Production has cache misses, failovers and deploys.
  3. No think time. Hammering as fast as possible measures a different system from one with realistic arrival patterns.
  4. A single operation. Production is a mix, and the mix is what causes contention between paths.
  5. Too little data. A table with 10,000 rows behaves nothing like one with 10 billion, per Part 15.
  6. Testing to failure only. Useful once. What you need is the shape of the latency curve up to and past your target, so you can see where the knee is.
the most valuable test
Replay real production traffic, at a multiple, against a realistic dataset. Shadow traffic to a staging environment gives you the real distribution of accounts, amounts and operation mix, including the skew that synthetic tests never reproduce. The hot key you did not think of is the one that takes you down, and only real traffic contains it.
26

Capacity planning, and headroom as policy

worked numbers
the inputs:
  current peak                      464 transfers/s
  growth rate                      ~8% per month
  lead time to add capacity    2 weeks
  target utilisation               65%
  failure headroom               survive losing one AZ of three

required capacity =
  464 × 1.08^3        (3 months of growth)
  ÷ 0.65                 (utilisation target)
  × 1.5                  (survive losing 1 of 3 AZs)
  = ~1,350 transfers/s of provisioned capacity

nearly 3x current peak, and every multiplier is justified.
headroom as a policy rather than a number
Write down the utilisation target and the failure-survival requirement, because otherwise every capacity conversation restarts from scratch and someone always proposes running hotter to save money. "We run at 65% and survive one AZ" is a policy; "we have enough servers" is a hope.
27

Performance as a regression test

what to put in CI, and what not to
  1. In: microbenchmarks on hot paths, with a tolerance. Catches an accidental O(n²) before it ships.
  2. In: query plan checks. Assert that the balance query uses the expected index. A plan change is a silent 100x regression and it happens when data grows.
  3. In: a smoke load test on a realistic dataset, checking p99 against the budget.
  4. Out: full-scale load tests. Too slow and too noisy for every commit. Run them nightly or before a release.
  5. Out: absolute thresholds on shared CI. Noisy neighbours make them flaky, and a flaky performance test gets disabled within a month.
the one most teams skip
Asserting the query plan. Part 15 showed that index-only scans depend on vacuum state and that a composite index only helps when the predicate matches its prefix. A query that silently stops using an index is the most common serious performance regression in a database-backed system, it does not show up in a unit test, and asserting the plan in CI costs almost nothing.