Lesson 16 / 25

How Distributed Systems Fail

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.
Figure 6.1 — Overload spreading across a cluster.

How retries amplify load

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

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.

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

  • Slow dependencies return wrong data
  • Slow dependencies cannot be monitored
  • 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.