Tutorials › Core Cloud Architecture › Reliability in Distributed Systems

Core Cloud Architecture · Part 9 of 13

Reliability in Distributed Systems

One slow part of the system is a performance problem. Left unchecked, it becomes everyone's problem.

Every piece described so far (a load balancer, a database, a cache, a queue) will eventually be slow, unavailable, or wrong. A distributed system's reliability comes from what happens to everything else when one of them is, since preventing it altogether isn't achievable at scale.

Timeouts, retries, backoff, and jitter

A timeout is a limit on how long a caller will wait for a response before giving up. Without one, a single slow dependency can hold a caller's resources (a thread, a connection) indefinitely, which is the mechanism behind the cascading failure further down. A retry tries the same call again after a failure, on the assumption that some failures are transient. Retrying immediately and repeatedly can make things worse, though: it adds load onto a dependency that's already struggling, at the moment it can least handle it. Exponential backoff waits longer between each successive retry (1 second, then 2, then 4, and so on), giving a struggling dependency room to recover instead of being retried into the ground. Jitter adds randomness to that wait time so that many callers retrying after the same failure don't all retry at exactly the same moment, which would otherwise recreate the original spike in near-perfect synchrony.

Circuit breakers

A circuit breaker tracks a dependency's recent failure rate and, once it crosses a threshold, stops sending it requests entirely for a cooldown period, failing fast instead of continuing to wait on (and add load to) something that's clearly not responding. After the cooldown, it lets a small number of requests through to test whether the dependency has recovered, and only resumes normal traffic if they succeed. This protects both sides: the caller stops wasting time waiting on a doomed request, and the struggling dependency stops receiving load it can't handle while it's trying to recover.

Idempotency, graceful degradation, backpressure, and load shedding

Idempotency, from the messaging patterns, matters here too: if a caller times out and retries a request that succeeded on the far end, an idempotent operation makes that retry harmless instead of a duplicate charge or a duplicate order. Graceful degradation means serving a reduced but useful experience when a non-critical dependency is down: showing a product page without personalized recommendations if the recommendation service is unavailable, instead of failing the whole page. Backpressure is a system signaling upstream that it can't keep up, so the sender slows down instead of piling on load the receiver has no way to process. Load shedding is the blunter version of the same idea: rejecting some incoming requests (often the least important ones, by some priority) once the system is over capacity, so the requests that are accepted can be served.

A concrete cascading failure

These patterns exist because of a specific, common failure shape: one slow dependency, left unmanaged, taking down services that have nothing wrong with them on their own.

flowchart TD
  A[Payment provider
becomes slow] --> B[Checkout calls it
with no timeout] B --> C[Each call holds a
thread/connection] C --> D[Checkout's thread
pool fills up] D --> E[Checkout can't accept
new requests] E --> F[Storefront & mobile app
callers time out too] F --> G[Whole app looks down;
one slow dependency]

The payment provider never failed outright. Slow was enough to take down the chain. A timeout on the call to the payment provider would have freed those threads instead of letting them queue up indefinitely; a circuit breaker would have stopped sending it new requests once its failure rate crossed a threshold, instead of continuing to feed the backlog; graceful degradation might have let checkout queue the order and confirm payment asynchronously instead of blocking the whole request on it. Any one of these patterns, applied at the point where checkout calls the payment provider, breaks the chain before it reaches the storefront.

RTO, RPO, high availability, and disaster recovery

High availability is about surviving component failure with little to no interruption: the multi-AZ pattern is a high-availability design. Disaster recovery (DR) is about recovering from a much larger failure, such as a whole region gone or a catastrophic data corruption, and is measured by two numbers. RTO (recovery time objective) is how long the system is allowed to be down before it's back up. RPO (recovery point objective) is how much data the organization can afford to lose, measured as time: if backups run every 6 hours, the RPO is up to 6 hours of data. Pick both numbers from what the business can tolerate before building an architecture to hit them. A tighter RTO and RPO are achievable, and each notch tighter costs more in infrastructure and operational complexity.

These patterns keep one failure as one failure, contained to the part of the system that owns it, instead of a chain reaction that takes down parts that were never broken. A distributed system with enough moving parts will have something failing somewhere at any given time.