AI System Design and Architecture

Best Practices for Scalability and Fault Tolerance


At 02:41 a Redis node fails over. It is only a cache — the inference service can run without it, just slower. Instead, the entire API goes down for nineteen minutes.

Here is the chain. The Redis client was created without a socket timeout, so cache.get() blocked indefinitely. The service runs 200 worker threads; within four seconds all 200 were blocked on a dead socket. The health check endpoint needed a worker thread to answer, so it stopped answering. The load balancer marked every instance unhealthy and removed all of them. The autoscaler saw zero healthy targets and started replacing nodes, each of which came up, connected to the same dead Redis, and hung in exactly the same way.

Not one component in that story was designed to fail like that. The failure was emergent: an optional dependency became mandatory because nobody bounded how long the system would wait for it, and a health check reported the wrong thing. Fault tolerance is the practice of making sure this specific class of surprise cannot happen — and almost all of it comes down to bounding waits, isolating resources, and being honest about what "healthy" means.

How a cache failure took the whole API downRedis nodefails overCache callhas no timeoutEvery requestthread blocksPool exhausted,readiness failsAll replicaspulled fromthe balancerA cache is only an optimisation if failing to reach it is faster than not using it at all.
Nineteen minutes of outage built from one missing timeout and one thread pool shared with the critical path.

The arithmetic of redundancy

Two components in series — a request must pass through both — multiply their availabilities. Two in parallel — either one suffices — multiply their unavailabilities.

ArrangementFormulaWith 99.9% componentsDowntime per 30-day month
One componentAA99.900%43.2 min
Three in seriesA3A^399.700%129.6 min
Five in seriesA5A^599.501%215.5 min
Two in parallel1−(1−A)21-(1-A)^299.9999%0.043 min
Three in parallel1−(1−A)31-(1-A)^399.9999999%0.00004 min

Redundancy is spectacularly effective and dependency chains are spectacularly corrosive. Every synchronous call you add to the request path is a series term. A request that touches an auth service, a feature store, a model service and a logging sink — all synchronous, all 99.9% — is 0.9994=99.6%0.999^4 = 99.6\% available, or 172 minutes of downtime a month, before any of your own code has a bug.

Removing a synchronous dependency from the request path improves availability more than making that dependency more reliable ever will.

Sizing redundancy properly (N+1 is usually wrong)

Peak traffic 600 req/s; each replica sustains 200 req/s. The naive answer is 3 replicas, and the slightly-less-naive answer is 4 ("N+1"). Both are wrong, because they ignore utilisation.

With 4 replicas at peak, each carries 150 req/s — 75% utilisation, already at the edge of where queue delay bites. Lose one replica and the remaining 3 carry 200 req/s each: 100% utilisation, unbounded queue growth, cascading timeouts. The redundancy you paid for buys you an outage that arrives ten seconds later than it would have anyway.

Redundancy that leaves the survivors at 100% utilisation is not redundancy — it is an outage with a ten-second delay.

Size for the degraded state instead. If your maximum acceptable utilisation is 75%, then after losing one replica you need (N−1)×200×0.75≥600(N-1) \times 200 \times 0.75 \ge 600, so N−1≥4N - 1 \ge 4 and N=5N = 5. Five replicas: 120 req/s each at 60% in normal operation, 150 req/s each at 75% after a loss. The cost is 5/3 = 1.67× the bare minimum, and that multiple is your availability.

Spread them across failure domains too. Five replicas in one availability zone survive one machine failure and not one zone failure. Across three zones, losing a whole zone removes at most two, leaving three replicas carrying 200 req/s each — back at 100% utilisation. That is exactly why the calculation should use "largest failure domain", not "one machine": to survive a zone loss at 75% you need four survivors, so six replicas, two per zone.

Bounding waits: timeouts, retries, breakers

Timeouts on everything, always

The opening outage had exactly one root cause: an unbounded wait. Every network call — HTTP, database, cache, message broker, DNS — needs an explicit timeout set to something derived from measured p99, not left at the library default, which is frequently infinite.

Python
import redisfrom redis.backoff import NoBackofffrom redis.retry import Retry# The line that would have prevented the outage.cache = redis.Redis(    host="cache.internal",    socket_connect_timeout=0.2,   # p99 connect is 8 ms; 200 ms is generous    socket_timeout=0.05,          # p99 GET is 0.9 ms; 50 ms means "it is broken"    retry=Retry(NoBackoff(), 0),  # a cache is optional; do not retry it    health_check_interval=30,)def cached_predict(key, compute):    try:        hit = cache.get(key)        if hit is not None:            return decode(hit)    except redis.RedisError:        CACHE_ERRORS.inc()        # degrade, never fail the request    value = compute()    try:        cache.setex(key, 3600, encode(value))    except redis.RedisError:        pass    return value

A 50 ms timeout on a call whose p99 is 0.9 ms is not aggressive — it is 55× the observed tail. Note the retry argument: since redis-py 6, the client retries timeouts by default (the old retry_on_timeout=False flag is deprecated and no longer turns this off), so without it a 200 ms connect timeout against a dead host can stretch to several seconds per call. With it in place, the Redis failover would have cost each request 50 ms and the service would have stayed up at a degraded 45 ms → 95 ms latency.

Retries, and the multiplication problem

Retries help exactly when failures are independent and transient. They hurt when the downstream is overloaded, because a retry is more load applied precisely when load is the problem.

Two rules make retries safe. First, retry at one layer only. Three layers each retrying three times means one client request can become 3×3×3=273 \times 3 \times 3 = 27 downstream calls, so a service running at 70% utilisation is pushed to 18.9× its capacity by a brief blip. Second, use exponential backoff with full jitter.

The jitter point deserves numbers. Suppose 10,000 clients see an error at the same instant and all retry after exactly 1 second. At t=1 s the service receives a 10,000-request spike, fails again, and at t=3 s receives another. The retries are perfectly synchronised, forever. With full jitter — waiting a uniformly random time in [0,2n][0, 2^n] seconds — the first retry wave spreads over 2 seconds at an average of 5,000 req/s, the second over 4 seconds at 2,500 req/s, and the herd disperses instead of resonating.

Python
import random, timedef call_with_retry(fn, attempts=3, base=0.25, cap=4.0, budget_s=8.0):    deadline = time.monotonic() + budget_s    for n in range(attempts):        try:            return fn(timeout=min(2.0, deadline - time.monotonic()))        except Permanent:            # 4xx: retrying cannot help            raise        except Transient:            if n == attempts - 1 or time.monotonic() >= deadline:                raise            time.sleep(random.uniform(0, min(cap, base * 2 ** n)))  # full jitter

Note the budget_s deadline. Retrying past your caller's timeout is pure waste: the caller has already given up, so the work is discarded whatever the outcome.

Circuit breakers

A circuit breaker stops sending requests to a dependency that is clearly broken, so you fail in one millisecond instead of waiting for a timeout.

Text
          failures ≥ 50% of last 20 calls   CLOSED ──────────────────────────────────► OPEN     ▲                                          │     │  probe succeeds                          │ after 30 s     │                                          ▼     └──────────────── HALF-OPEN ◄──────────────┘          (allow 1 request through and watch)                  probe fails → back to OPEN

Quantify what it saves. A dependency is down; requests arrive at 200 req/s; the timeout is 5 seconds. Without a breaker, in-flight requests accumulate to 200×5=1,000200 \times 5 = 1{,}000 — every thread and connection in the pool consumed by calls that are certain to fail. With a breaker open, each request fails in about 1 ms, so in-flight work is 200×0.001=0.2200 \times 0.001 = 0.2. The breaker's real job is not to protect the dependency; it is to stop your own resources being consumed by a known-lost cause.

The half-open state is what makes it self-healing: after a cool-down, exactly one request is allowed through. If it succeeds the circuit closes; if it fails the timer restarts. Without half-open you need a human to reset it.

Bulkheads

Name from ship design: compartments that keep a hull breach from sinking the vessel. In software, give each dependency its own bounded pool so one slow dependency cannot consume all your capacity.

Text
  200 worker threads, unpartitioned:      with bulkheads:  ┌──────────────────────────────┐        ┌──────┬──────┬──────┬──────┐  │ all 200 blocked on dead      │        │ redis│ model│  db  │ spare│  │ Redis → service is down      │        │  20  │ 120  │  40  │  20  │  └──────────────────────────────┘        └──────┴──────┴──────┴──────┘                                            ↑ dead Redis costs 20 threads,                                              the other 180 keep serving

Graceful degradation

Decide in advance what the system does with less. Write it as a table so that the decision is made in daylight rather than at 02:41.

FailureNaive behaviourDegraded behaviourUser impact
Cache unavailable500 errorCompute directlySlower: 22 ms → 45 ms
Large model timing out504Fall back to the distilled modelAccuracy 94.1% → 91.3%
Personalisation service down500Serve globally popular resultsLess relevant, still useful
Feature store down500Use last-known features with a staleness flagSlightly stale scores
Analytics sink downBlocks the requestDrop to a local buffer, fire-and-forgetNone
GPU fleet saturatedEveryone times outShed load: 429 low-priority tiers firstSome clients throttled, rest fine

Load shedding is the least popular and most important row. When demand exceeds capacity, serving 70% of requests correctly beats serving 100% of them at a 40-second latency that every client will time out on anyway. Shed by priority — background batch jobs before interactive traffic — and return 429 with Retry-After so clients back off rather than hammering.

Health checks that tell the truth

Three probes, three questions, three different consequences:

ProbeQuestionFailure actionShould it check dependencies?
StartupHas initialisation finished?Keep waiting (long timeout)No
LivenessIs the process wedged?Restart the containerNever
ReadinessCan it serve a request right now?Remove from the load balancerOnly mandatory ones

The rule people break is putting dependency checks in liveness. If liveness pings the database and the database has a hiccup, every replica fails liveness, every replica restarts simultaneously, and now you have a cold fleet and a sick database. Liveness must only answer "is this process capable of executing code", and it must be served by a path that cannot be starved by worker threads.

Readiness may check mandatory dependencies — for a model server, "are the weights resident on the GPU" is mandatory; "is Redis reachable" is not. Gate readiness on a completed warm-up inference, otherwise every scale-up sends traffic to a replica that will time out for the next 45 seconds.

Python
@app.get("/healthz")           # liveness: no I/O, no locks, no dependenciesdef healthz():    return {"ok": True}@app.get("/readyz")            # readiness: mandatory dependencies onlydef readyz():    if not MODEL.loaded:        return JSONResponse({"ready": False, "reason": "loading"}, 503)    if not MODEL.warmed:        return JSONResponse({"ready": False, "reason": "warming"}, 503)    if SHUTTING_DOWN:           # drain: fail readiness ~15 s before exiting        return JSONResponse({"ready": False, "reason": "draining"}, 503)    return {"ready": True, "model": MODEL.version, "gpu_mem_free_mb": MODEL.free_mb}

The SHUTTING_DOWN flag implements graceful drain. On SIGTERM, fail readiness immediately but keep serving for a period longer than the load balancer's deregistration delay — typically 15–30 seconds — then exit. Without it, every deployment drops the requests that were in flight, and a rolling update of 20 replicas produces 20 small error spikes.

Balancing and state

Least outstanding requests should be your default for inference, because service times vary by an order of magnitude with input size. Round robin will hand a request to a replica already two seconds deep in a long generation while another idles; least-outstanding will not. On a fleet with 40× service-time spread, this single setting can move p99 by 20×.

Session affinity — pinning a client to a replica — is occasionally necessary (a conversation's KV cache lives on one GPU) and always costly. It defeats even balancing, it makes scale-in disruptive, and it turns one replica's death into a specific set of users' failure rather than a spread. If you need it, prefer consistent hashing with bounded loads so a hot key cannot overwhelm one node, and always have a path to rebuild state elsewhere.

Databases and connection pools

Size the pool from Little's Law, not from a round number. If each request makes one 3 ms query at 600 req/s, average concurrent queries is 600×0.003=1.8600 \times 0.003 = 1.8; a pool of 10 per replica is generous. With 20 replicas that is 200 connections against a Postgres instance configured for 100 — which fails, loudly, at the worst moment. Pool sizing must be computed fleet-wide: replicas × pool_size ≤ max_connections − headroom. When the arithmetic does not fit, put a connection proxy in front rather than raising the limit.

Read replicas scale reads and introduce lag. A replica 200 ms behind the primary breaks read-your-own-writes: a user updates a preference, the next request reads a stale value, and the UI appears to lose the change. The fixes, in order of preference: route a user's reads to the primary for a few seconds after they write; version the entity and retry the read if the version is older than the one just written; or accept staleness and show it explicitly.

Consistency you can live with

DataConsistency neededMechanism
Billing and usage countersStrongPrimary writes, transactional, idempotency keys
Model registry / active versionStrongSmall, rarely written, read through a short-TTL cache
Prediction resultsEventualCache keyed by input + model version
Feature valuesBounded stalenessTTL plus a staleness flag on the response
Analytics and logsEventual, lossy acceptableAsync, buffered, never on the request path

Deploying without breaking things

Text
  Blue-green                          Canary  ┌──────────┐                        ┌──────────┐  │  100% ───┼──► blue  (v7)          │  95%  ───┼──► v7  └──────────┘    green (v8) idle     │   5%  ───┼──► v8       switch ─────────────►          └──────────┘  ┌──────────┐                          then 25% → 50% → 100%  │  100% ───┼──► green (v8)            rollback = shift weight back  └──────────┘    blue  (v7) idle  rollback = flip back (seconds)
Blue-greenCanaryRolling
Blast radius of a bad release100% until rollbackOnly the canary shareGrows as the roll proceeds
Rollback speedSeconds (traffic flip)Seconds (weight change)Minutes (roll back through)
Extra capacity required2× during cutover~5–10%~10%
Detects subtle quality regressionsPoorly — all or nothingWell — compare cohortsPoorly
Cost for GPU fleetsHigh (double the GPUs)LowLow

Put numbers on the blast radius. A new model version has an 8% error rate and takes 10 minutes to detect. Under blue-green at 100% traffic the error exposure is 10×1.00×0.08=0.810 \times 1.00 \times 0.08 = 0.8 error-minutes. Under a 5% canary it is 10×0.05×0.08=0.0410 \times 0.05 \times 0.08 = 0.04 error-minutes — twenty times less. Against a 99.9% monthly SLO with a 43.2-minute error budget, neither is fatal on its own; but a team shipping twenty times a month burns 16 minutes of budget with blue-green and 0.8 minutes with canary.

For models, canary needs one addition that software canaries do not: compare quality, not just errors and latency. A new model that returns 200 OK for everything while its mean confidence has dropped from 0.91 to 0.68, or whose positive-class rate has doubled, is failing silently. Gate promotion on prediction-distribution comparison between the canary and the incumbent, not only on the error rate.

The shape of a resilient service

Text
                    ┌───────────────────────────┐   clients ────────►│ LB · least-outstanding    │                    │ health = /readyz, 5 s     │                    └────────────┬──────────────┘             ┌───────────────────┼───────────────────┐             ▼                   ▼                   ▼       ┌──────────┐        ┌──────────┐        ┌──────────┐       │ zone A   │        │ zone B   │        │ zone C   │  N = 6 replicas       │ ×2       │        │ ×2       │        │ ×2       │  50% util normal       └────┬─────┘        └────┬─────┘        └────┬─────┘  75% after AZ loss            └───────────────────┼───────────────────┘                                ▼        ┌───────────────────────────────────────────────────┐        │ per-replica bulkheads                             │        │  ├─ cache pool 20  · timeout 50 ms  · breaker on  │  optional        │  ├─ model pool 120 · timeout 6 s    · breaker on  │  degrade→small model        │  ├─ db pool 10     · timeout 400 ms · breaker on  │  mandatory        │  └─ telemetry: async buffer, never blocks         │  fire-and-forget        └───────────────────────────────────────────────────┘

Every element there answers a specific way the opening outage happened: bounded pools so one dependency cannot take the process, explicit timeouts so nothing waits forever, breakers so a dead dependency costs a millisecond not five seconds, a liveness probe with no dependencies so a sick cache cannot trigger a fleet-wide restart, and replica sizing that survives losing a whole zone at 75% utilisation.

Turning this into a checklist you can act on

Take your service and answer six questions with numbers, today. What is the timeout on every outbound call, and how does each compare with that call's measured p99? Which dependencies are mandatory and which are optional, and does the code actually degrade for the optional ones or does it throw? At what utilisation does the fleet run after losing its largest failure domain? Does liveness touch any dependency? Does readiness fail before the process exits, and for longer than the load balancer's deregistration delay? And how many layers between the client and the model perform retries?

Then run the experiment rather than trusting the answers. Kill a replica in production during business hours and watch whether p99 moves. Block the cache with a firewall rule for sixty seconds and confirm the service degrades to 45 ms instead of returning errors. Introduce 500 ms of latency on the database and see whether the bulkhead holds. Untested fault tolerance is not fault tolerance — it is a set of beliefs, and the difference is only ever revealed at 02:41 with an audience.