Sven Erik Matzen

Software Architect | Cloud & Security Expert | AI-enabled Solutions

The Long Tail of Latency: Tail Latency and Why Averages Lie in the Cloud

🎧 Listen to this article

Cloud Computing · 2026-09-27

EU label: fully AI-generated content Fully AI-generated article (no prior review).

The Hook: Every Part Is Fast, the Whole Is Slow

Imagine you run a service spread across 100 machines. A search query goes out to all 100, each machine searches its slice of the data, and only once the last answer has arrived can the result go back to the user. You have measured carefully: each individual machine responds within 10 milliseconds in 99 out of 100 cases. Only in one case out of 100 does it take a full second — some background process, a garbage collection run, an unfortunately timed disk access. One percent outliers. That sounds like a system you can put into production with a clear conscience.

Do the arithmetic. The probability that all 100 machines are fast is 0.99 to the power of 100 — roughly 0.366. Which means: in 63 percent of all user requests, at least one machine takes a full second, and because you must wait for the last answer, the user waits that second too. A system whose parts are lightning fast 99 percent of the time is sluggish as a whole in almost two out of three cases.

This calculation opens one of the most influential papers of the cloud era: "The Tail at Scale" by Jeffrey Dean and Luiz André Barroso, published in Communications of the ACM in 2013. In it, the two Google engineers articulated an insight that shifted how we think about distributed systems: as soon as a system consists of many components that all have to answer, it is not the typical response time but the rare bad one that determines the user experience. And the larger the system grows, the stronger this effect becomes. Scaling does not merely improve throughput — it systematically amplifies the influence of outliers.

The term for this is tail latency — the latency at the right-hand edge, in the "tail" of the distribution. And perhaps the most uncomfortable consequence: almost every dashboard developers look at daily shows precisely the number that says nothing here — the average.


Part 1: What Latency Actually Is — and Why the Average Answers the Wrong Question

Latency Is a Distribution, Not a Number

The first error is linguistic. We say "the service's latency is 40 milliseconds," as though latency were a property like mass or length. It is not. Latency is a random variable with a distribution, and in practically every real system that distribution is right-skewed and heavy-tailed: there is a lower bound (you cannot go faster than the sum of the physically necessary steps), but no upper one. The path downward is short; upward it is open.

For such a distribution, the mean is one of the least informative statistics you can compute. It mixes two entirely different populations — the dense main body of fast responses and the thin, far-stretched banner of slow ones — into a single number that describes neither. A service with a mean of 40 ms may mean "almost all requests take 38 to 42 ms" or "95 percent take 10 ms and 5 percent take 600 ms." These are completely different systems that behave completely differently in production — and they look identical on the average-latency dashboard.

Worse still: the mean is systematically too optimistic relative to what users experience. If 1 percent of requests take 100 times as long as the rest, that one percent contributes only about half of the mean — so the mean merely doubles, while for the affected users the world stands still. The arithmetic mean is built to smooth outliers away. That is exactly what you must not want here.

Percentiles: The Right Question, Properly Posed

The usable alternative is percentiles. The 99th percentile (customarily abbreviated p99) is the value below which 99 percent of all measured response times lie. "p99 = 250 ms" means: one request in a hundred takes longer than 250 milliseconds. That is a statement about users, not about statistics.

One point that often gets lost in practice matters here: which percentile you look at is not a matter of taste but follows from the number of interactions per user session. If you build a web application in which a single page view generates 100 backend calls, then p99 is not a rare exception for you but the normal case: with 100 calls, practically every page view hits at least one p99 event. Dean and Barroso state the consequence starkly: in high fan-out systems, the 99.9th percentile (p999) is the relevant metric, because a building block's 99th percentile can already become the median of the composed service.

A second, subtler point: percentiles do not add up. If service A has a p99 of 50 ms and service B has a p99 of 50 ms, then the chain A→B does not have a p99 of 100 ms. The probability that both simultaneously exhibit their bad case is small (0.0001 under independence), but the probability that at least one is slow nearly doubles: 1 − 0.99² ≈ 1.99 percent. The chain's p99 therefore lies closer to 50 ms than to 100 ms, while its p98 degrades. Percentiles are not quantities you can add budgets with — you have to reason in probabilities.

The Measurement Error That Spoils Almost Every Benchmark: Coordinated Omission

Before thinking about taming the tail, you have to be able to measure it correctly. And here lies a systematic error that the Java performance engineer Gil Tene, in his much-cited talk How NOT to Measure Latency, named coordinated omission: most load generators report high percentiles that are too good — not because they compute wrongly, but because they never take the worst cases into their sample in the first place.

The mechanism is simple, and that is what makes it so treacherous. A typical load generator works in a loop: send request, wait for response, record time, send next request. If the system now stalls for two seconds, that thread stays blocked for two seconds — and sends no further requests during that time. At a target rate of 1,000 requests per second, it should have dispatched 2,000 requests in those two seconds, all of which would have been slow. Instead it records one measurement of 2,000 ms. The 1,999 bad measurements that should have existed are missing. The load generator has tacitly "cooperated" with the system under test and stopped measuring precisely during the outage.

The result is not a small inaccuracy but a distortion of orders of magnitude in exactly the region you care about. A test rig reporting p999 = 20 ms may in truth have p999 = 2,000 ms. Tene therefore built HdrHistogram, a tool that records latencies across many orders of magnitude at constant relative precision and allows corrections for coordinated omission. The practical rule that follows: measure against a schedule, not against a loop. A request that should have gone out at 10:00:00.000 and is answered only at 10:00:01.500 has 1,500 ms of latency — even if the server first laid eyes on it at 10:00:01.498. That is the latency the user experiences.


Part 2: The Mathematics of Fan-Out Amplification

The Basic Law

The formula behind the opening example is trivially simple and brutal in its consequences. Let p be the probability that a single building block answers slowly, and n the number of building blocks whose answers must be awaited. Then the probability that the overall request becomes slow is:

P(slow) = 1 − (1 − p)ⁿ

For small p and moderate n, this is well approximated by P ≈ n · p. To a first approximation, the probability of a slow overall response therefore grows linearly with the number of components involved. That is the real message: tail latency is not a problem that also appears with size — it is a problem that size creates.

Dean and Barroso give a second example that matches the scale of modern systems: a service spread across 2,000 machines in which only one request in 10,000 per machine is slow still delivers a slow response on nearly one in five user requests (1 − 0.9999²⁰⁰⁰ ≈ 18 percent). You improve the reliability of the individual component by a factor of 100 — and the situation is still bad, because you have simultaneously scaled out by a factor of 20.

What This Does to Real Numbers

Dean and Barroso present measurements from a Google-internal service operating over a BigTable dataset. The numbers are instructive because they make the amplification visible in a concrete system:

Measurement (taken at the root service) 99th percentile
A single, randomly chosen sub-request 10 ms
95 percent of sub-requests completed 70 ms
All sub-requests completed 140 ms

Between "one sub-request" and "all sub-requests" lies a factor of 14. And the jump from "95 percent done" to "100 percent done" — the price of the last five percent — doubles the latency once more. Anyone who has internalized this table understands why the question "do we really need all the partial answers?" is an architectural question of the first rank in distributed systems.

Why Microservice Architectures Build the Problem In

The two dangerous patterns can be named precisely.

The first is fan-out: one service queries many services in parallel and waits for all of them. Here the formula above applies in full force, because the outlier probabilities accumulate.

The second is depth: a request traverses a chain of services sequentially. Here the latencies do add up, but the same principle holds: the probability that an outlier sits somewhere along the way grows with the chain's length. With ten hops of p99 = 20 ms each, the chance that at least one hop exhibits its bad behavior is already 1 − 0.99¹⁰ ≈ 9.6 percent.

In practice, real architectures combine both — a call graph with fan-out at every level and several levels of depth. The outlier probability multiplies along the whole graph. I am of the opinion that this is the most underappreciated argument against unnecessarily fine-grained service decomposition: every additional network boundary in the critical path is not merely an additional number of milliseconds, but an additional draw from a heavy-tailed distribution.


Part 3: Where the Slow Responses Come From

If you want to shorten the tail, it helps to understand where it comes from. Dean and Barroso call the sources variability — and they show that in shared infrastructure this variability does not arise randomly but structurally.

Shared Resources and the Noisy Neighbor

A modern cloud machine does not run one process but many — containers belonging to different tenants, background services, monitoring agents. They share CPU cores, level-3 cache, memory bandwidth, network interface, and disk queues. Each of these resources is a coupling through which a neighbor's load becomes your process's latency. The term noisy neighbor understates the phenomenon: it is not about noise but about queues in which someone else's work stands ahead of yours.

The coupling through shared caches is particularly insidious. When a neighboring process fills the level-3 cache with its data, memory access time for your process rises from nanoseconds to hundreds of nanoseconds — an effect visible in no application metric, which accumulates across millions of memory accesses into milliseconds. It is precisely this problem that makes isolation technologies such as Betting on a Stranger's Code: Firecracker microVMs and the End of the Container-versus-VM Dilemma not only a security topic but a latency topic as well.

Background Maintenance

The second major source is periodic background activity that serves operations but hurts in the foreground:

  • Garbage collection. A collection pause in a managed runtime (JVM, .NET, Go) stops exactly those threads that are currently serving requests. Even modern, largely concurrent collectors have phases with stop-the-world character.
  • Compaction. Write-optimized storage engines must reorganize their data periodically. How this mechanism works and why it inevitably produces load spikes is described at length in Write First, Sort Later – Log-Structured Merge-Trees and the Inversion of the Database: the price of very fast writes is background merges that consume disk I/O and CPU in irregular bursts.
  • Log rotation, index rebuilds, metric aggregation, certificate renewal. All small operations — but each one a candidate for a 100 ms outlier if it lands badly.
  • Power and clock management. Processors with aggressive power-saving modes need microseconds to milliseconds to move from a deep sleep state to full-load operation. Thermal throttling works in the same direction.

The decisive point: all of these activities are legitimate. You cannot switch them off, only shift, smooth, or — and this is the architectural lever — synchronize them. Dean and Barroso recommend running maintenance across a replica group simultaneously rather than independently: if all replicas are briefly slow at the same time, there is a short, predictable gap; if each replica becomes slow independently and at random, there is permanently always at least one slow replica — and with fan-out you always hit it.

Queueing: Why Utilization Is Punished Exponentially

The deepest cause is mathematical and has nothing to do with software at all. It sits in queueing theory.

Consider the simplest case: one server, requests arriving randomly (Poisson-distributed), service times exponentially distributed — the classic M/M/1 model. Let ρ be the utilization (arrival rate divided by service rate). Then the mean waiting time in the queue, excluding the actual service, is:

W_q = (1 / μ) · ρ / (1 − ρ)

The factor ρ/(1 − ρ) is the decisive term. Plug in numbers:

Utilization ρ Factor ρ/(1−ρ) Waiting time relative to service time
50 % 1.0 1×
70 % 2.3 2.3×
80 % 4.0 4×
90 % 9.0 9×
95 % 19.0 19×
99 % 99.0 99×

The curve is a hyperbola with a pole at ρ = 1. The practical consequence is fundamental: between 50 and 70 percent utilization, additional load costs almost nothing; between 90 and 99 percent it costs everything. And because the curve is nonlinear, the same holds for the spread: at high utilization the distribution is not merely shifted to the right but becomes dramatically wider. The tail therefore does not arise from errors in the first instance — it arises from utilization.

Here lies the real reason why "cost efficiency through high utilization" and "low tail latency" are mutually exclusive goals. Anyone running their machines at 95 percent CPU has chosen bad high percentiles — not through a poor implementation, but through the choice of operating point. On this reading, every capacity reserve you hold is not waste but purchased latency stability.

Complementing this, Little's Law (L = λ · W) provides the bridge to observability: the mean number of requests simultaneously in the system equals the arrival rate multiplied by the mean time in the system. Anyone measuring the number of in-flight requests — a metric almost every framework provides for free — therefore has an indirect and very early indicator of latency problems, often earlier than the latency histogram itself.

The Microsecond Gap

A fourth, comparatively new factor: the hardware landscape has shifted. In a follow-up paper, "Attack of the Killer Microseconds" (Barroso, Marty, Patterson, and Ranganathan, Communications of the ACM, 2017), the authors argue that our tools are well built for two time scales — nanoseconds (hardware parallelism, out-of-order execution, prefetching) and milliseconds (context switching, asynchronous I/O) — but offer almost nothing serviceable for the microsecond.

The figures from the paper make the problem tangible: DRAM access sits at tens to hundreds of nanoseconds, a classic hard disk at a few milliseconds, flash storage at tens of microseconds, and traversing a datacenter across 200 to 300 meters of cable costs about one microsecond. It is exactly in this band — flash, RDMA, emerging non-volatile memory — that modern cloud infrastructure operates. And here both familiar strategies fail: hardware parallelism cannot hide microseconds, and a software context switch often costs more than the wait it is supposed to bridge. This is why techniques such as The Programmable Kernel: eBPF and the Sandbox at the Heart of the Operating System — processing inside the kernel, without the detour through user space — have become so relevant to latency.


Part 4: Tail-Tolerant Design — The Techniques That Shorten the Tail

The conceptual shift Dean and Barroso propose is the real punchline of their paper: you should not try to make every component uniformly fast — with shared infrastructure, legitimate background work, and queueing physics, that is a hopeless program. Instead, you should build systems that are tail-tolerant: systems that deliver a stable overall latency despite variable components. This is the same intellectual move by which reliability engineering went from "failure-free components" to "fault-tolerant systems."

Hedged Requests: Ask the Same Question Twice

The simplest technique: send the request to one replica. If after a short wait — say the duration of the 95th percentile — no answer has arrived, send the same request to a second replica and take whichever comes back first. The second request is cancelled as soon as an answer is in hand.

The effect in numbers, measured on a BigTable benchmark retrieving 1,000 values from 100 servers: a second request issued after a 10 ms delay lowered the 99.9th percentile for collecting all 1,000 values from 1,800 ms to 74 ms — at only 2 percent additional requests. That is a factor of 24 for two percent of extra work, and it is why this technique became so popular.

Why does it work so well? Because the delay kicks in only past the 95th percentile, the second request affects only the 5 percent of cases that are slow anyway — and in those cases the probability that both replicas have an outlier simultaneously is the product of two small numbers. You buy good high percentiles with a few percent of extra load. The price is real load, and that is the pitfall: a threshold set too early doubles the traffic and thereby pushes ρ upward — which, per the table above, achieves precisely the opposite of what was intended.

Tied Requests: Let Both Replicas Inform Each Other

Hedged requests have a window of waste: the time between the first and second request. Tied requests eliminate it. You send the request to two replicas simultaneously, but you tell each one who the other is. As soon as one replica dequeues the request and begins working on it, it sends the partner replica a cancellation message. The partner discards the request, provided it is still sitting in the queue.

You are thus not choosing the replica with the shortest queue, but the one that actually starts working first — a subtle but important difference, because queue lengths say nothing about the service time of the waiting tasks. Google's measurements: in an otherwise unloaded cluster, a 1 ms delay reduced median latency by 16 percent and the 99.9th percentile by nearly 40 percent, at under 1 percent additional disk utilization.

The technique has a prerequisite that must not be overlooked: the cancellation message must reach the partner replica faster than the request is processed. In a datacenter with microsecond round-trips that holds; over a wide-area link it does not.

Micro-Partitioning and Selective Replication

The third family attacks not the individual request but the data distribution.

Micro-partitioning means: split the data into considerably more partitions than there are machines — factors of 10 to 100 are common — and assign partitions to machines dynamically. This has two effects. First, load balancing becomes a fine-grained operation: an overloaded machine hands off some of its partitions instead of requiring a recomputation of the entire assignment. Second, recovery becomes faster, because a failed machine's partitions can be spread across many healthy machines. How to design such an assignment so that a change in the set of machines moves only a minimal fraction of the partitions is the subject of The Ring That Shares the Load: Consistent Hashing and the Art of Moving Gracefully.

Selective replication goes one step further: detect especially hot partitions and create additional copies of exactly those. This is the answer to the fact that real access distributions are almost never uniform but follow a power law — a few keys account for most of the traffic.

Latency-induced probation is the corresponding operational discipline: a system watches its backends' response times and temporarily removes a conspicuously slow replica from rotation while continuing to send it shadow requests to check when it has recovered. The decisive part is that this is a latency-based decision, not an error-based one: the replica does answer, just too slowly. A classic health check based on "responds / does not respond" would declare it healthy. How a cluster arrives at a shared view of which nodes are in which state at all is treated in The Whisper of Machines: Gossip Protocols, SWIM, and How a Cluster Learns Who Is Still Alive.

Canary Requests and "Good Enough" Instead of "Complete"

Two further techniques from the paper deserve mention because they cost so little.

Canary requests protect against the case where a request triggers a latent bug in all replicas at once — a fan-out to 10,000 servers that crashes all 10,000 is a real incident type. The countermeasure: send the request to one or two servers first, and only once they answer successfully, proceed with the full fan-out. The cost is an additional round-trip; the benefit is avoiding a total outage.

"Good enough" responses are perhaps the most effective and least used technique. It says: define for your service what an incomplete but usable answer is, and deliver it when the deadline expires. In a search across 100 shards, a result from 98 shards is almost always as good as one from 100 — and per the table in Part 2 it saves a factor of two. The prerequisite is a product decision, not a technical one: you must be willing to trade completeness for predictability.

Choosing the Target: The Power of Two Choices

Where two or more replicas are candidates, the question arises how to choose. The answer is one of the most elegant results in computer science. Michael Mitzenmacher showed in The Power of Two Choices in Randomized Load Balancing (IEEE Transactions on Parallel and Distributed Systems 12(10), pp. 1094–1104, 2001) that having one versus two options to choose from makes a qualitative difference.

If you distribute n tasks purely at random across n servers, the maximum queue length grows on the order of log n / log log n. If instead you draw two servers at random for each task and take the one with the shorter queue, the maximum length falls to roughly log log n / log 2 — an exponential gain. The difference is qualitative, not gradual: log n grows without bound with system size, log log n practically does not — with a million servers, log₂ log₂ n is about four. And the asymptotics say something else important: the jump from one choice to two delivers almost the entire attainable gain, the jump from two to three very little more.

For the practitioner this is extraordinarily good news. A load balancer that draws two candidates at random and picks the less loaded one ("power of two random choices," often called P2C) needs no global state, no coordination, and no complete view of the cluster — and still comes very close to ideal balance. Schemes that poll all servers are not only more expensive, they are also prone to herd behavior, because many clients discover the same "best" server simultaneously.

Shuffle Sharding: Isolation as a Latency Tool

One final, often underestimated technique comes from the operational practice of large cloud providers and is described in the AWS Builders' Library by Colm MacCárthaigh. The basic idea: if a single tenant slows down its backends with a pathological request, as few other tenants as possible should be affected.

With classic sharding you divide eight workers into four pairs; a problem hits a quarter of all tenants. With shuffle sharding you assign each tenant a randomly drawn combination of two of the eight workers. There are C(8,2) = 28 such combinations — the scope of impact shrinks to one twenty-eighth, that is, to one seventh of what classic sharding delivers, without having added a single machine.

The effect scales combinatorially in a breathtaking way. Amazon Route 53 arranges its capacity into 2,048 virtual name servers and assigns each customer domain a combination of four of them. That is roughly 730 billion possible combinations — with the consequence that no customer domain shares more than two virtual name servers with any other. A tenant whose traffic brings one backend to its knees therefore cannot practically take down any other tenant entirely.

Shuffle sharding is thus the structural answer to part of the tail problem: it does not limit variability itself, but its blast radius.

A Map of the Techniques

Technique Attacks Additional load Prerequisite
Hedged requests an individual slow request a few percent if threshold > p95 idempotent reads, ≥ 2 replicas
Tied requests queueing time before service < 1 % fast cancellation message between replicas
Micro-partitioning load imbalance, recovery metadata overhead dynamic partition assignment
Selective replication hot partitions storage for extra copies hotspot detection
Latency-induced probation persistently slow replica shadow requests per-backend latency measurement
Canary requests request-induced total outages one round-trip two-stage fan-out
"Good enough" responses the last few percent of fan-out none product decision about completeness
Power of two choices target selection negligible coarse load info on 2 candidates
Shuffle sharding blast radius of a problem none controllable tenant-to-capacity mapping
Capacity reserve (lower ρ) the root of the problem infrastructure cost willingness not to run at 95 %

Part 5: What Hurts in Practice

Retry Storms and Metastable Failure

The most obvious reaction to slow responses — "then just try again" — is also the most dangerous. If a service becomes slow because it is overloaded, and all clients thereupon repeat their requests, then load rises at exactly the moment the system can least absorb it. Per the queueing table in Part 3, this pushes ρ upward, latency explodes, more timeouts fire, more retries are created. The feedback is positive — in the mathematical, not the colloquial sense.

Bronson, Aghayev, Charapko, and Zhu gave this phenomenon a precise name in 2021: metastable failure. The core of their analysis: a system can enter a state in which it remains permanently overloaded even though the triggering disturbance is long over. The overload sustains itself, because the work it generates (retries, cache misses after a cache flush, aborted and restarted transactions) is greater than the work it completes. The practical consequence is unpleasant: such states do not leave a system on their own. You have to actively take load away — through load shedding, throttling, or, in the worst case, a controlled restart.

The remedies are well known and described in the AWS Builders' Library (Marc Brooker, Timeouts, retries, and backoff with jitter): exponential backoff, jitter — a random component in the wait time so clients do not retry in lockstep — retry budgets that cap retries at a small percentage of total traffic, and circuit breakers that temporarily shut off a path recognized as unhealthy. Especially important is the rule not to stack retries across multiple layers: three layers with three attempts each yield, in the worst case, 27 requests for one user action.

Incast

A purely networking effect that strikes regularly with fan-out: when a client queries 100 servers simultaneously and all 100 answer simultaneously, 100 responses arrive practically at once on the same network interface. The switch buffer in front of it overflows, packets are dropped, TCP interprets this as congestion and waits — and a microsecond response becomes a response with a TCP retransmission timeout. Rajesh Nishtala's team describes in Scaling Memcache at Facebook (NSDI 2013) how they fought this effect with a sliding window over outstanding requests: the client does not send all requests at once but observes a limit on in-flight requests. Chosen too small, latency rises through unnecessary serialization; chosen too large, incast returns.

Timeouts Are Architectural Decisions

A timeout is not a safety measure but an assertion about the latency distribution. Set it to the mean plus a little buffer and you abort a substantial share of requests that were on the verge of succeeding — producing load without benefit. Set it too high and threads and connections block long enough for overload to pile up. The sensible starting point is a measured high percentile of the downstream service, and the sensible discipline is passing along a remaining budget (deadline propagation): the topmost service decides how much time the whole request may take, and every downstream call is told the remaining budget instead of inventing its own timeout.

SLO Arithmetic

Anyone formulating service level objectives should keep the fan-out formula in mind. A latency SLO such as "99 percent of all requests under 300 ms" is tenable for a service with twenty internal calls only if each of those calls meets a considerably stricter target — not 300/20 = 15 ms in the sense of addition, but a percentile target high enough that 1 − (1 − p)²⁰ stays below one percent. That requires p < 0.0005, that is, a p9995 per call. Such numbers are uncomfortable, but they are the honest arithmetic.


Part 6: The Overarching Design Principle

Stepping back, tail latency is a special case of a much more general pattern: in systems composed of many parts whose overall result is determined by the worst part, the relevant quantity is not the average but the spread. A system whose components are slow on average but uniformly so behaves predictably. A system whose components are fast on average but occasionally catastrophic behaves unpredictably — and unpredictability is what users experience as unreliability.

The same pattern recurs in many well-designed distributed systems: Clocks That Know Their Own Uncertainty: Google Spanner, TrueTime, and Mastering Time in the Cloud replaces a timestamp that tacitly claims to be exact with an interval whose uncertainty is carried explicitly — the spread becomes a first-class datum instead of an unspoken risk. Eleven Nines: Erasure Coding, Reed-Solomon, and How the Cloud Makes Data Practically Unlosable builds durability not from durable disks but from computational redundancy over unreliable ones. How Machines Come to Agree: Distributed Consensus from FLP to Paxos to Raft achieves progress not by having all nodes answer but by having a majority suffice — quorum logic is at heart a tail-tolerance technique.

The Central Takeaway

The practical lesson is this: stop looking at averages and start looking at a distribution — then compute how many draws from that distribution a single user request requires. Concretely, that is four steps you can take in any project. First: collect latency histograms instead of means, and check whether your measurement tooling is armored against coordinated omission — a dashboard that shows p999 as too good is worse than none. Second: count the fan-out and the depth of your critical path and compute 1 − (1 − p)ⁿ; the number that comes out is the honest description of your user experience. Third: look at utilization — if your services run above 80 percent, capacity rather than code is the cheapest latency improvement you can buy. Fourth: deliberately pick a tail-tolerant technique from the table in Part 4 instead of hoping to make every component uniformly fast — with shared infrastructure, that goal is physically unattainable.

A Question to Ponder

Every technique in this article buys latency stability with some resource: with additional load (hedged requests), with additional storage (selective replication), with unused capacity (low ρ), or with completeness ("good enough" responses). None is free. The question that follows is not a technical one but one about values: how much completeness would you be willing to give up in your own system to gain predictability — and who in your organization is actually allowed to make that decision? Because in experience it is rarely made explicitly. It is made implicitly, in the moment someone writes a timeout value into a configuration file.


Cross-References in the Vault


Sources

← All articles