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
- 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.
- 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.
- 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
- 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.
- p99 is one customer in a hundred. At 18,000 balance reads a second, that is 180 people per second having a bad time.
- p99.9 is where the interesting bugs live: GC pauses, lock contention, a cold cache, a retry, a failover.
- 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.
| Step | p50 | p99 | Note |
|---|---|---|---|
| TLS and gateway | 3 ms | 12 ms | Connection reuse matters enormously here |
| Auth and validation | 2 ms | 6 ms | |
| Limit check (Redis) | 1 ms | 4 ms | One round trip, one Lua script |
| Fraud score | 18 ms | 150 ms | Hard timeout, fails open |
| Shard routing | <1 ms | 1 ms | In-memory map |
| Posting transaction | 22 ms | 95 ms | Dominated by fsync and replica ack |
| Outbox insert | 2 ms | 5 ms | Same transaction |
| Response | 1 ms | 3 ms | |
| Total | ~50 ms | ~276 ms | Against 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
- 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.
- Headroom is a feature. It is what absorbs the 4x peak from Part 3 and the backlog after a failover.
- Variance makes it worse. The formula assumes exponential service times; a bimodal distribution queues harder at the same utilisation.
- 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
| Cause | Signature | Typical fix |
|---|---|---|
| Doing too much | Linear in input size, CPU-bound | Do less. Better algorithm, less data, precompute |
| Waiting | Low CPU, high latency | Parallelise, batch, or remove the round trip |
| Contention | Latency rises with concurrency | Reduce the critical section, shard the hot thing |
| Coordination | Latency floors at network RTT | Fewer 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
- Reproduce it. A performance problem you cannot reproduce is a performance problem you cannot verify you fixed.
- 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.
- Form a hypothesis with a number. "I think the N+1 in the beneficiary lookup costs 40 ms of the p99." Falsifiable.
- Change one thing. Two changes at once and you learn nothing about either.
- Measure again, the same way. Same load, same environment, same percentile.
- 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
- 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.
- A warm cache and no cold start. Production has cache misses, failovers and deploys.
- No think time. Hammering as fast as possible measures a different system from one with realistic arrival patterns.
- A single operation. Production is a mix, and the mix is what causes contention between paths.
- Too little data. A table with 10,000 rows behaves nothing like one with 10 billion, per Part 15.
- 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
- In: microbenchmarks on hot paths, with a tolerance. Catches an accidental O(n²) before it ships.
- 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.
- In: a smoke load test on a realistic dataset, checking p99 against the budget.
- Out: full-scale load tests. Too slow and too noisy for every commit. Run them nightly or before a release.
- 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.