# How Distributed Systems Fail — Scalability, Availability & Reliability

Source: https://www.skillbyai.com/en/scalability/r-failures

> Recognise crash, slow, partial and cascading failures.

## Failures are normal, and often partial

At scale something is always broken: a disk, a node, a network link, a dependency. The dangerous failures are rarely clean crashes. **Slow failures** (a dependency answering in 30 seconds instead of 30 ms) tie up threads and connections and spread upstream. **Partial and gray failures** affect only some requests, some users or some nodes, and pass health checks. **Cascading failures** happen when one overloaded component sheds load onto others, which then overload too: a node dies, its traffic moves to the remaining nodes, they tip over, and so on. **Retry storms** multiply load exactly when a system is struggling: three layers each retrying three times turn one request into 27. A **thundering herd** happens when many clients act at once, for example reconnecting after an outage. **Metastable failures** are states where the system stays broken even after the trigger disappears, because the recovery load (retries, cold caches) keeps it overloaded.

## A cascading failure

One overloaded node pushes its load onto neighbours until they fail too.

![A row of five boxes where the first is dark and cracked, arrows push weight to the next, which is darker, and so on, with the last boxes still light.](assets/figures/scalability/section-6-map.svg) — Figure 6.1 — Overload spreading across a cluster.

## How retries amplify load

Retries at every layer multiply when a deep dependency slows down.

```text
client       -> retries 3x
  gateway    -> retries 3x
    service  -> retries 3x
      db     (slow)

worst case: 3 x 3 x 3 = 27 database calls for one user action

fixes:
  retry at one layer only (usually the closest to the failure)
  use exponential backoff with jitter and a retry budget (e.g. retries <= 10% of requests)
  fail fast with timeouts shorter than the caller's timeout
```

## Slow is worse than down

A dead dependency fails fast; a slow one quietly exhausts your thread pools and connection pools. Every network call needs a timeout, chosen deliberately.

**Quiz:** Why can a slow dependency be more dangerous than one that is completely down?

- [ ] Slow dependencies return wrong data
- [ ] Slow dependencies cannot be monitored
- [x] Waiting requests pile up and exhaust threads and connections in every caller
- [ ] Load balancers route more traffic to slow services

*Answer:* Waiting requests pile up and exhaust threads and connections in every caller. Hung calls hold resources, so the slowness spreads upstream.
