Partial failure
Partial Failure
Understand why distributed systems fail in fragments rather than entirely, why you cannot distinguish a slow component from a dead one, and how to design around uncertainty you can never remove.
Prerequisites
Recommended Knowledge
- Networking fundamentals and request lifecycle
- Processes, threads, and memory basics
- Service-to-service communication patterns
- Timeouts, retries, and idempotency
- Queues, delivery semantics, and consumer behavior
- Replication and replica lag
- Load balancing and health checks
- Latency percentiles and capacity reasoning
The Defining Difference
In a single process, a failure is total. The program either completes an operation or crashes, and if it crashes, nothing downstream of that point executed. This binary property makes local reasoning tractable.
A distributed system has no such property. Some components succeed while others fail, and crucially, the surviving components cannot reliably determine which is which.
Distributed Failure Some, Unknown Which
Simple Analogy
You post a letter and receive no reply. Was it lost in transit, delivered but unread, read but unanswered, or answered with a reply that was itself lost? The silence is identical in every case, yet each demands a different response.
The Uncertainty of a Timeout
When a request times out, the caller knows only that no response arrived within the allotted time. Every one of the following remains possible.
| Actual Situation | Did the Work Happen? | Correct Response |
|---|---|---|
| Request lost before arrival | No | Safe to retry |
| Server crashed before processing | No | Safe to retry |
| Server crashed mid-processing | Partially | Depends on atomicity |
| Processed, response lost | Yes | Retry would duplicate |
| Processed, response delayed | Yes | Retry would duplicate |
| Still processing slowly | In progress | Retry may run concurrently |
The Two Generals Problem
Two generals must attack simultaneously to win. They communicate only by messenger across enemy territory, where any messenger may be captured. Can they reach certain agreement?
The first general sends a plan but cannot know it arrived. An acknowledgement helps, but the second general cannot know the acknowledgement arrived. Each confirmation requires its own confirmation, without end.
| Guarantee | What Actually Happens | Requirement Placed on You |
|---|---|---|
| At-most-once | Never retried, may be lost | Tolerate missing operations |
| At-least-once | Retried until acknowledged | Handle duplicates safely |
| Exactly-once effect | At-least-once plus deduplication | Idempotency keys and state tracking |
Failure Detection Is Always a Guess
Since a slow node is indistinguishable from a dead one, every failure detector must choose a threshold. That choice trades two types of error against each other.
Short Timeout
- Detects genuine failures quickly
- Declares healthy slow nodes dead
- Removes capacity during load spikes
- Can trigger unnecessary failovers
Long Timeout
- Avoids false accusations
- Sends traffic to dead nodes for longer
- Extends user-visible outage
- Delays recovery actions
Health Checks Lie
A health check that returns success proves only that the health endpoint responded. It says nothing about whether the service can perform real work.
| Check Type | Verifies | Risk |
|---|---|---|
| Liveness | The process is running | Passes while wholly unable to serve |
| Readiness | Dependencies are reachable | One slow dependency removes every instance |
| Deep check | A real operation succeeds | Expensive and can itself cause load |
| Synthetic transaction | End-to-end path works | Slower to detect, costly to run |
Gray Failure
The most troublesome failures are not clean crashes. A component that is intermittently slow, partially correct, or failing only for certain requests evades every binary detector.
| Gray Failure Mode | Observed Behavior | Why Detection Fails |
|---|---|---|
| Degraded latency | Responds, but far slower | Health checks still succeed |
| Partial packet loss | Some requests fail intermittently | Probes may land on working paths |
| One bad replica | Stale data for a subset of reads | Responses are well-formed |
| Resource exhaustion | Fails only under specific load | Idle probes always succeed |
| Corrupted cache entry | Wrong answer returned quickly | Looks like a normal success |
| Asymmetric partition | A reaches B, B cannot reach A | Each side has a different view |
Network Partitions
A partition splits the system so that groups of nodes can communicate internally but not across the divide. Each side observes the other as failed.
| Partition Type | Characteristic | Particular Danger |
|---|---|---|
| Complete split | Two isolated groups | Both may accept conflicting writes |
| Asymmetric | Traffic flows one direction only | Contradictory failure views |
| Intermittent | Connectivity flaps repeatedly | Repeated leader changes |
| Single node isolated | One node cut off | Isolated node may act as if in charge |
Cascading Failure
Partial failure becomes total failure when the system's response to trouble generates more trouble. The most common amplifier is retrying.
With three layers each retrying three times, one client request can produce twenty-seven requests at the innermost service. A dependency that is merely slow receives an order of magnitude more load precisely when it can least handle it.
Initial Slowdown
A dependency degrades, and response times rise beyond the caller's expectation.
Resource Accumulation
Callers hold threads and connections while waiting. Available capacity shrinks even though nothing has failed outright.
Retry Storm
Timeouts trigger retries at every layer, multiplying load on the struggling dependency.
Spread
Callers exhaust their own resources and begin failing for unrelated requests that never touched the original dependency.
Recovery Blocked
Even after the root cause is fixed, accumulated retries and queued work prevent the system from stabilizing.
Bulkheads
Named after ship compartments, bulkheads isolate resource pools so exhaustion in one area cannot consume capacity needed elsewhere.
class Bulkhead {
constructor(name, maxConcurrent, maxQueued) {
this.name = name;
this.maxConcurrent = maxConcurrent;
this.maxQueued = maxQueued;
this.active = 0;
this.queued = 0;
}
async execute(operation) {
if (this.active >= this.maxConcurrent) {
if (this.queued >= this.maxQueued) {
throw new BulkheadRejection(
`${this.name} at capacity`
);
}
this.queued += 1;
await this.waitForSlot();
this.queued -= 1;
}
this.active += 1;
try {
return await operation();
} finally {
this.active -= 1;
this.releaseSlot();
}
}
}
const bulkheads = {
payments: new Bulkhead("payments", 40, 20),
recommendations: new Bulkhead("recommendations", 10, 0),
search: new Bulkhead("search", 30, 15)
};
Without isolation, a slow recommendation service can occupy every available thread, taking down checkout alongside it. With isolation, recommendations fail while payments continue.
Circuit Breakers
A circuit breaker stops calling a failing dependency, failing fast instead of accumulating waiting requests. It converts slow failure into immediate failure, which is far less damaging.
| State | Behavior | Transition Condition |
|---|---|---|
| Closed | Requests pass through normally | Opens when failure rate exceeds threshold |
| Open | Requests rejected immediately | Half-opens after a cooldown period |
| Half-open | Limited probe requests permitted | Closes on success, reopens on failure |
class CircuitBreaker {
constructor(config) {
this.state = "closed";
this.failureCount = 0;
this.successCount = 0;
this.openedAt = 0;
this.config = config;
}
async call(operation, fallback) {
if (this.state === "open") {
const elapsed = Date.now() - this.openedAt;
if (elapsed < this.config.cooldownMs) {
return fallback
? await fallback()
: Promise.reject(new CircuitOpenError());
}
this.state = "half-open";
this.successCount = 0;
}
try {
const result = await operation();
this.recordSuccess();
return result;
} catch (error) {
this.recordFailure();
if (fallback) {
return await fallback();
}
throw error;
}
}
recordSuccess() {
if (this.state === "half-open") {
this.successCount += 1;
if (this.successCount >= this.config.probesToClose) {
this.state = "closed";
this.failureCount = 0;
}
return;
}
this.failureCount = 0;
}
recordFailure() {
this.failureCount += 1;
if (this.state === "half-open") {
this.state = "open";
this.openedAt = Date.now();
return;
}
if (this.failureCount >= this.config.failureThreshold) {
this.state = "open";
this.openedAt = Date.now();
}
}
}
Timeout Budgets
Timeouts configured independently at each layer produce incoherent behaviour. An inner timeout longer than the outer one means the caller has already abandoned the request while the inner service continues consuming resources.
async function callWithDeadline(request, operation) {
const remainingMs = request.deadline - Date.now();
if (remainingMs <= MINIMUM_USEFUL_MS) {
throw new DeadlineExceededError(
"Insufficient time remaining to attempt call"
);
}
const budgetMs = Math.min(
remainingMs - RESERVE_FOR_RESPONSE_MS,
MAX_SINGLE_CALL_MS
);
return await operation({
timeoutMs: budgetMs,
deadline: request.deadline
});
}
Idempotency as the Foundation
Because a caller can never know whether an operation completed, safe retry requires that repeating an operation produces the same outcome as performing it once.
CREATE TABLE idempotency_records (
idempotency_key VARCHAR(80) PRIMARY KEY,
operation_type VARCHAR(60) NOT NULL,
request_hash VARCHAR(64) NOT NULL,
status VARCHAR(20) NOT NULL,
response_body JSONB NULL,
created_at TIMESTAMP NOT NULL,
completed_at TIMESTAMP NULL,
expires_at TIMESTAMP NOT NULL
);
-- Claim the key atomically; a conflict means a retry
INSERT INTO idempotency_records (
idempotency_key, operation_type, request_hash,
status, created_at, expires_at
)
VALUES (
:key, :operation, :hash,
'in_progress', CURRENT_TIMESTAMP,
CURRENT_TIMESTAMP + INTERVAL '24 hours'
)
ON CONFLICT (idempotency_key) DO NOTHING;
| Retry Encounters | Correct Behavior |
|---|---|
| No existing record | Proceed with the operation |
| Record marked completed | Return the stored response |
| Record still in progress | Return a conflict, do not run concurrently |
| Same key, different request body | Reject as a key reuse violation |
Graceful Degradation
Not every dependency is essential to every request. Designing which features may be sacrificed converts partial failure into reduced functionality rather than an error page.
| Dependency | Criticality | Behavior When Unavailable |
|---|---|---|
| Primary datastore | Essential | Fail the request explicitly |
| Authentication service | Essential | Reject, never bypass |
| Cache | Performance only | Fall through to origin |
| Recommendation engine | Enhancement | Serve a static popular list |
| Search service | Important | Offer category browsing instead |
| Analytics collector | Non-essential | Drop events silently |
| Inventory check | Business decision | Accept order with later verification |
Failure Modes and Responses
| Failure Mode | Symptom | Design Response |
|---|---|---|
| Crash stop | Node halts and stays down | Replication and failover |
| Crash recovery | Node returns with stale state | Durable logs and state reconciliation |
| Omission | Some messages silently dropped | Acknowledgements and retry |
| Timing | Responses arrive too late to use | Deadlines and cancellation |
| Byzantine | Component returns wrong answers | Checksums, validation, cross-checking |
| Correlated | Many components fail together | Diverse zones and staggered deploys |
| Metastable | System stays broken after cause is removed | Load shedding and queue draining |
Testing for Partial Failure
Partial failure behaviour cannot be verified by testing the happy path. Failure must be injected deliberately.
Conditions Worth Injecting
- Dependency returning errors at a set rate
- Dependency responding just under the timeout
- Dependency accepting connections but never replying
- Network partition between service groups
- Asymmetric connectivity in one direction only
- Duplicate message delivery
- Out-of-order message arrival
- Instance termination mid-operation
- Clock skew between nodes
- Replica lag exceeding normal bounds
Observing Partial Failure
Signals Worth Tracking
- Success rate per dependency, not just overall
- Latency percentiles rather than averages
- Timeout rate distinguished from error rate
- Retry volume and amplification ratio
- Circuit breaker state transitions
- Bulkhead rejections by pool
- Thread and connection pool saturation
- Queue depth and oldest item age
- Requests abandoned past deadline
- Per-instance divergence in behavior
- Client-observed versus server-reported success
Common Design Mistakes
Weak Design
- Treating remote calls like local functions
- Assuming a timeout means nothing happened
- Retrying independently at every layer
- Omitting timeouts entirely
- Sharing one thread pool across all dependencies
- Health checks that query shared dependencies
- Treating all dependencies as essential
- Monitoring averages instead of percentiles
- Never testing degraded conditions
Strong Design
- Treats every remote call as fallible
- Makes operations idempotent by default
- Assigns retry ownership to one layer
- Propagates absolute deadlines
- Isolates dependencies with bulkheads
- Uses independent, shallow health checks
- Classifies dependencies by criticality
- Tracks tail latency per dependency
- Injects failure routinely
System Design Interview Discussion
| Question | What Your Answer Should Cover |
|---|---|
| What happens if this call times out? | Ambiguity of outcome and idempotent retry |
| How do you detect a failed node? | Threshold trade-off and false positives |
| What if a dependency is slow, not down? | Gray failure, bulkheads, circuit breakers |
| How do you prevent cascading failure? | Retry budgets, isolation, load shedding |
| Which features can be degraded? | Dependency criticality classification |
| How are timeouts coordinated? | Deadline propagation across hops |
| What happens during a partition? | Quorum and split-brain prevention |
| How is this verified? | Fault injection and failure drills |
Design Checklist
Production Checklist
- Set an explicit timeout on every remote call
- Propagate absolute deadlines across hops
- Make write operations idempotent
- Accept and honour idempotency keys
- Assign retry ownership to a single layer
- Apply exponential backoff with jitter
- Enforce a retry budget per dependency
- Isolate dependencies into separate pools
- Add circuit breakers on external calls
- Classify every dependency by criticality
- Define fallback behavior per non-essential dependency
- Keep security controls failing closed
- Avoid shared dependencies in readiness checks
- Use quorum to prevent split-brain writes
- Monitor tail latency per dependency
- Alert on retry amplification and breaker state
- Inject failure regularly in a controlled way
- Document expected behavior for each dependency outage
Knowledge Check
Why is partial failure harder than total failure?
Some components succeed while others fail, and the survivors cannot determine which occurred, so the system must act under permanent uncertainty.
What does a timeout actually tell you?
Only that no response arrived in time. The operation may have fully succeeded, partially completed, or never started.
Why is exactly-once delivery impossible?
The two generals problem shows acknowledgements require their own acknowledgements indefinitely. Systems achieve exactly-once effects through idempotency instead.
Why do retries cause cascading failure?
Independent retries at each layer multiply, so a struggling dependency receives several times its normal load exactly when it is least able to cope.
Why is gray failure worse than a crash?
A crash is detected and routed around. A component that is slow or intermittently wrong passes health checks while continuing to receive and damage traffic.
Summary
Partial failure is the defining characteristic of distributed systems. Unlike a single process, where failure is total, distributed components fail independently, and the surviving components cannot determine which failed or whether requested work was completed.
A timeout conveys nothing about remote state. The two generals problem establishes that certain agreement over an unreliable network is impossible, which is why exactly-once delivery cannot exist and idempotency becomes foundational rather than optional.
Failure detection is always a guess balancing false positives against slow recovery. Gray failures, where components are slow or intermittently wrong, evade binary detectors and frequently cause more damage than clean crashes.
Partial failure becomes total failure through cascades, usually amplified by uncoordinated retries. Containment comes from bulkheads isolating resources, circuit breakers failing fast, deadline propagation bounding work, and deliberate degradation that sacrifices non-essential functionality rather than the whole service.
Key Takeaway
Design for uncertainty you cannot eliminate. Assume every remote call may have succeeded, failed, or be ongoing. Make operations idempotent, own retries at one layer, propagate deadlines, isolate dependencies, fail fast when a dependency is unhealthy, and decide in advance which functionality you are willing to lose.