Table of Contents

    Partial failure

    SYSTEM DESIGN • CHAPTER 14.1

    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.

    Learning objective: By the end of this article, you will understand the defining difference between local and distributed failure, the two generals problem and why exactly-once delivery is impossible, failure detection limits, gray failure, cascading collapse, and the patterns that contain damage rather than pretending to prevent it.

    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.

    THE CORE PROPERTY
    Local Failure All or Nothing

    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 hardest problem in distributed systems is not that components fail. It is that you cannot tell the difference between a component that failed and one that is merely slow.

    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 FUNDAMENTAL LIMIT
    A timeout tells you about your own waiting, never about the remote system's state. No timeout value makes this distinction possible.
    Dangerous Assumption Treating a timeout as evidence that the operation did not occur. Retrying a payment on this basis charges the customer twice.
    Correct Design Assume the operation may have completed. Make the operation idempotent so a retry converges on the same result regardless of which scenario actually occurred.

    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 INFINITE REGRESS
    Send Plan Need Acknowledgement Need Ack of Ack Never Terminates

    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.

    Why this matters practically: This is the formal reason exactly-once delivery across an unreliable network is impossible. Systems claiming it achieve at-least-once delivery combined with idempotent processing, which produces exactly-once effects.
    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
    DETECTION TRADE-OFF
    \[ T_{\text{detect}} \uparrow \;\Rightarrow\; P_{\text{false positive}} \downarrow, \; T_{\text{recovery}} \uparrow \]

    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
    Correlated Health Failure If every instance's readiness check queries the same database, one slow database marks the entire fleet unhealthy. The orchestrator then removes all capacity, converting degradation into a complete outage.

    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
    The observability gap: Gray failures are often visible to clients long before they appear in the system's own health metrics. Monitor from the caller's perspective, not only from the component's self-report.

    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.

    THE SPLIT-BRAIN RISK
    Partition Occurs Both Sides Assume Leadership Divergent State
    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
    Quorum as Protection Requiring a majority before accepting writes ensures only one side of a partition can proceed, since two disjoint groups cannot both hold a majority of the same cluster.

    Cascading Failure

    Partial failure becomes total failure when the system's response to trouble generates more trouble. The most common amplifier is retrying.

    RETRY AMPLIFICATION
    \[ L_{\text{observed}} = L_{\text{original}} \times \prod_{i=1}^{n} r_i \]

    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.

    1

    Initial Slowdown

    A dependency degrades, and response times rise beyond the caller's expectation.

    2

    Resource Accumulation

    Callers hold threads and connections while waiting. Available capacity shrinks even though nothing has failed outright.

    3

    Retry Storm

    Timeouts trigger retries at every layer, multiplying load on the struggling dependency.

    4

    Spread

    Callers exhaust their own resources and begin failing for unrelated requests that never touched the original dependency.

    5

    Recovery Blocked

    Even after the root cause is fixed, accumulated retries and queued work prevent the system from stabilizing.

    RETRY BUDGET RULE
    Only one layer should own retries. Independent retries at every hop multiply rather than add, and nobody is accountable for the total.

    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();
            }
        }
    }
    The counterintuitive benefit: Refusing to send requests gives the struggling dependency room to recover. Continuing to call it guarantees it never will.

    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.

    BUDGET PROPAGATION
    \[ T_{\text{remaining}} = T_{\text{deadline}} - T_{\text{elapsed}} \]
    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
        });
    }
    Deadline Propagation Pass an absolute deadline rather than a relative timeout. Each hop computes its own budget from the remaining time, and work stops everywhere once the deadline passes.

    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.

    IDEMPOTENCY PROPERTY
    \[ f(f(x)) = f(x) \]
    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
    Degradation Anti-Pattern Bypassing authorization when the permission service is unreachable. Security controls must fail closed, regardless of availability pressure.

    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
    Metastable failure deserves attention: A system can remain in a broken state indefinitely after the trigger disappears, because retries and queued work sustain the overload. Recovery may require shedding traffic deliberately.

    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
    A system that has never been tested against a slow dependency has an untested response to the most common real failure.

    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
    AVERAGES CONCEAL
    Partial failure affects a subset of requests. An average latency that looks acceptable can hide a tail where a meaningful fraction of users are failing entirely.

    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

    1

    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.

    2

    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.

    3

    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.

    4

    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.

    5

    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.