A bidder fetches features from eight online-store shards in parallel per impression - why is the fetch stage's p99 far worse than one shard's p99?
answer
- the slowest lookup ends the stage
- a maximum, not an average
- eight lookups, eight chances to be slow
- 0.99 to the eighth is 92%
- each shard must hold p99.87
basics
~20 sBecause the stage ends when the slowest lookup answers. With eight independent lookups, roughly 8% of impressions contain at least one shard past its own p99, so the stage's p99 sits far out in a single shard's tail.
solid answer
~50 sThe stage is a barrier: it finishes when the last of the eight lookups returns, so its latency is the maximum, not the average. If each shard independently exceeds its p99 on 1% of calls, all eight stay inside on 0.99^8 of impressions, about 92.3% - so roughly 7.7% of impressions wait past a shard's p99. Turned around, holding the stage at p99 requires each shard to be inside the line at about its p99.87, because 0.99 raised to the one-eighth power is 0.99874. That is why per-shard averages are useless for budgeting this line. The controls that matter are the width of the fan-out, a hedged duplicate for a straggler, and making the stage deadline-driven so it can stop waiting. Measure the stage's own p99 rather than deriving it, because real shards are not independent.
go deeper
Remember that a parallel fetch finishes when its slowest call returns, so the stage takes the maximum of the lookups rather than their average, no matter how well the others did.
Do the arithmetic: with eight independent lookups, about 8% of requests see at least one past its p99, and holding the stage at p99 needs each shard inside the line at roughly p99.87.
Pick the lever and defend it - narrow the fan-out, hedge the straggler, or make the stage stop waiting - and say why you would measure the stage rather than derive it from per-shard numbers.
Treat fan-out width as a budget commitment across teams: adding features widens the fetch and silently takes milliseconds from the scorer, so the trade needs numbers and an owner.
## The stage is a barrier, so its latency is a maximum A per-impression feature fetch that touches eight partitions of an online store is not eight small latencies; it is one latency, the **maximum** of eight draws, because the scorer cannot start until the feature vector is complete. Every intuition built on averages fails here. Adding a ninth lookup does not add one-ninth of anything - it adds one more chance for the whole stage to be slow. ## The arithmetic of tail amplification Assume for a moment that the shards are independent and each exceeds its own p99 on 1% of calls. The chance that **all** of them stay inside is 0.99 raised to the width of the fan-out. | fan-out width | all inside the per-shard p99 | at least one past it | per-shard percentile needed to hold the stage at p99 | |---|---|---|---| | 1 | 99.0% | 1.0% | p99 | | 2 | 98.0% | 2.0% | p99.5 | | 4 | 96.1% | 3.9% | p99.75 | | 8 | 92.3% | 7.7% | p99.87 | | 16 | 85.1% | 14.9% | p99.94 | The last column is the same statement read the other way: to keep the stage's own p99 inside its line, each shard must hold that line at the percentile whose eighth power is 0.99, which is 0.99874 - the shard's p99.87. Very few stores are characterised that far out, which is one reason fan-out tails are so often a surprise. ## What independence assumes, and why the real number must be measured The table is the independent case, and independence is the **worst** case for pure amplification. If the shards were perfectly correlated - slow together, fast together - the maximum would be no worse than a single draw and the amplification would vanish. Real systems sit between: shards share hosts, racks, network paths and sometimes a noisy co-tenant, so their slow moments partly coincide. The trap is to conclude that correlation makes the problem worse. It does not, mechanically; what makes it worse is that the **shared cause which creates the correlation usually degrades each shard's own distribution at the same time**, and a heavier marginal beats a friendlier dependence structure. The practical instruction is short: derive the table to understand the mechanism, then **measure the stage's p99 directly** and budget from the measurement. ## What to do about it 1. **Narrow the fan-out.** Co-locating the keys one impression needs, so the request makes two lookups instead of eight, moves the requirement from p99.87 per shard to p99.5 per shard - by far the largest single lever. Check afterwards that concentrating those keys has not made the surviving partitions hotter, because a heavier marginal can give back what the narrower fan-out won. 2. **Hedge the straggler.** Once a lookup has been outstanding longer than the per-shard p95, send a duplicate to another replica and take whichever answers first. This converts a rare long wait into a short one at the cost of a few percent extra read load. It pays when the tail comes from one unlucky host; it makes things worse when the store is globally saturated, because the duplicates add load to the exact resource that is short. 3. **Make the stage deadline-driven.** The fetch has a line, and when the line is spent the stage should stop waiting rather than run to completion, handing on whatever it has. 4. **Never budget this line from an average.** A stage whose latency is a maximum has to be budgeted from the distribution of that maximum. ## The budget consequence Fan-out width is a **budget decision**, not only a data-layout one. Every millisecond the fetch stage takes past its line is a millisecond the scorer does not get, and the checkpoint that guards the deadline will spend it by degrading the response rather than by running late. A team that widens the fan-out to add features is therefore trading model input for scoring time, and that trade should be made explicitly with the numbers on the table rather than discovered later as a rise in the timeout rate. The same reasoning explains why the fetch stage is so often the one that breaches: it is usually the widest fan-out on the path, and width is the one property that turns an ordinary dependency tail into the stage that decides the budget.
- You cut the fan-out from eight lookups to two by co-locating the keys. What happens to the stage p99?The amplification requirement drops from roughly the per-shard p99.87 to p99.5, so the stage's measured p99 moves much closer to a single lookup's tail. The catch is that the same keys now sit on fewer partitions, so check whether those partitions got hotter: a heavier per-shard distribution can give back the win the narrower fan-out bought.
- When does hedging a slow lookup pay, and what does it cost?Send a duplicate to another replica once the original passes the per-shard p95, and take the first answer. It pays when the tail is caused by one unlucky host, since the second copy is very unlikely to be unlucky too, and the extra load is only a few percent. It is harmful when the store is already saturated, because the duplicates add load to the resource that is short.
- Should the fetch stage be budgeted from the measured stage p99 or from the per-shard numbers?From the measured stage p99. The independence arithmetic is a way to understand the mechanism and to predict how a width change will behave, but real shards share hosts and network paths, so the true distribution of the maximum has to be observed rather than derived.
saying these in an interview costs you the question
- Says parallel fan-out is free because the calls overlap
- Sizes the stage from a shard's average lookup time
- Believes adding shards lowers the stage p99 by spreading load
- Expects hedging to be free and ignores the extra read load
- Derives the stage p99 by arithmetic instead of measuring it
- Treats fan-out width as a storage layout choice with no latency budget effect