OperationsGuide 8 of 11
Resilience in distributed systems
How do you keep a partial failure from spreading through the whole system? Resilience mechanisms, when to use them, and where their limits are.
Updated 10 min read
// on this page
A distributed system changes the failure model completely. Inside a local process, a function either returns a result or throws an exception. In a remote interaction:
flowchart LR A["Service A"] --> N["Network"] --> B["Service B"]
the request may never arrive, it may arrive twice, Service B may process the operation and lose the response, Service B may take too long, or only part of the system may be failing.
A partial failure doesn’t stay inside the component that failed. If Payments waits indefinitely on a downed provider, it stops serving checkout; if it retries a charge that already went through, it charges twice. The mechanisms in this article don’t stop things from failing: they limit the damage, make retrying safe, and decide what degrades when not everything can be served.
Timeout
A timeout keeps a dependency from holding resources indefinitely:
sequenceDiagram participant P as Payment participant Pr as Provider P->>Pr: request Note over P,Pr: ... P--xPr: timeout
But:
A timeout doesn’t mean the operation didn’t happen.
This detail is critical:
sequenceDiagram participant C as Client participant Pr as Provider C->>Pr: charge() Note over Pr: charges the card Pr--xC: response lost
The client sees a timeout. The card has already been charged.
Retry
In the face of transient failures, a retry sends the same operation again:
flowchart LR R["Request"] --> F1["Failure"] --> Rt1["Retry"] --> F2["Failure"] --> Rt2["Retry"]
A robust retry usually accounts for exponential backoff, jitter, a maximum number of attempts, and a time budget. Without that ceiling, a slow provider turns into load amplification: every client retries and the provider gets even more traffic.
Retry + non-idempotent operation
This is dangerous:
sequenceDiagram participant C as Client participant Pr as Provider C->>Pr: ChargeCard() Note over C,Pr: timeout C->>Pr: retry ChargeCard() Note over Pr: charged twice
That’s why timeouts, retries, and idempotency have to be designed together.
Circuit breaker
Lets you stop sending requests to a dependency that is failing:
stateDiagram-v2 [*] --> CLOSED CLOSED --> OPEN: failure threshold exceeded OPEN --> HALF_OPEN: probe after cooldown HALF_OPEN --> CLOSED: probe succeeds HALF_OPEN --> OPEN: probe fails
While it’s OPEN, the system can fail fast instead of piling up timeouts and continuing to hammer the dependency.
Idempotency
Take this:
POST /payments
Idempotency-Key: abc123
On the first request, the provider creates the payment but the response is lost. On the second request, with the same Idempotency-Key, the system recognizes that the operation has already been processed and returns the existing payment instead of creating a new one:
flowchart LR K["Idempotency-Key: abc123"] --> E["Existing payment"]
It’s essential for payments, orders, webhooks, and message consumers.
At-least-once delivery
Many messaging systems deliver each message at least once:
sequenceDiagram participant Broker participant Consumer Broker->>Consumer: message Note over Consumer: processes Note over Consumer: fails before the ack Broker->>Consumer: message (redelivered) Note over Consumer: duplicate message
That’s why one practical strategy is to combine at-least-once delivery with an idempotent consumer.
Bulkhead
The name comes from the bulkheads of a ship: one flooded compartment shouldn’t sink the rest. In software, you isolate resource pools — connections, threads, semaphores — so a failure in one area doesn’t take another down with it:
flowchart TD App["Application"] --> PP["Payment Pool"] App --> OP["Order Pool"] App --> NP["Notification Pool"]
If Notifications degrades, it shouldn’t consume every resource available to Payments.
Rate limiting
Controls the incoming flow before the system saturates:
flowchart LR C["Client"] --> RL["Rate Limiter"] RL -->|allowed| A["Backend"] RL -->|rejected| R["429"]
It protects against abusive clients, traffic spikes, and clients that generate extra load without meaning to. Among the best-known implementations are token bucket and leaky bucket. Unlike load shedding, rate limiting acts at the edge: it decides how much to accept, not what to drop once you’re already saturated.
Backpressure
Say a producer runs at 20k msg/s and a consumer at 5k msg/s. The gap builds up:
flowchart LR P["Producer<br/>20k msg/s"] --> Q["Queue<br/>(growing)"] --> C["Consumer<br/>5k msg/s"]
Backpressure means the consumer or the queue slows the producer down when they can’t keep up, instead of accepting work without limit and letting the queue grow unbounded.
Load shedding
When the system can’t process everything:
flowchart TD L["Incoming load"] --> O["Saturated system"] O -->|Critical| Pr["Process"] O -->|Optional| Re["Reject"]
It’s better to reject some requests than to let the whole system collapse.
Dead letter queue (DLQ)
A DLQ is the queue messages go to when a consumer couldn’t process them:
flowchart LR Q["Queue"] --> C["Consumer"] -->|failure + retries exhausted| DLQ["DLQ"]
The DLQ lets you inspect messages, fix problems, reprocess them later, and keep one bad message from blocking the queue permanently.
Saga
When a business operation spans several services, there’s no single ACID transaction. We can implement a Saga:
flowchart LR O["Create Order"] --> P["Authorize Payment"] --> I["Reserve Inventory"] --> S["Create Shipment"]
If Inventory fails, it triggers a chain of compensating actions:
flowchart LR F["Reserve Inventory ✘"] --> RP["Release Payment"] --> CO["Cancel Order"]
Graceful degradation
Not every component should carry the same level of criticality:
flowchart TD Ch["Checkout"] --> Pay["Payment — critical"] Ch --> Ord["Order — critical"] Ch --> Rec["Recommendation — optional"]
If Recommendations is down, Checkout stays available and only the Recommendations feature is lost. The architecture keeps what’s critical up.
How to analyze a failure
One practical approach is to walk through a series of questions:
flowchart TD D["The dependency fails"] --> Q1["Does the request time out?"] Q1 --> Q2["Should we retry?"] Q2 --> Q3["Is retrying safe?"] Q3 --> Q4["Should the circuit open?"] Q4 --> Q5["Can there be duplicate processing?"] Q5 --> Q6["What happens to queued messages?"] Q6 --> Q7["Is compensation needed?"] Q7 --> Q8["What does the user experience?"]
Tools
Depending on the platform: Resilience4j, Envoy, Istio, AWS API Gateway, AWS SQS, Kafka, RabbitMQ, Kubernetes health probes. Tools implement mechanisms. They don’t replace designing for failure modes.
None of these mechanisms is worth anything if you don’t verify that they work. A circuit breaker that has never opened in a controlled environment is a hypothesis, not a protection: it has to be exercised with performance testing under load and with failures injected on purpose.