Course Content
AI System Design and Architecture
3 sections · 7 lessons
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.
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.
| Arrangement | Formula | With 99.9% components | Downtime per 30-day month |
|---|---|---|---|
| One component | A | 99.900% | 43.2 min |
| Three in series | A3 | 99.700% | 129.6 min |
| Five in series | A5 | 99.501% | 215.5 min |
| Two in parallel | 1−(1−A)2 | 99.9999% | 0.043 min |
| Three in parallel | 1−(1−A)3 | 99.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% 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, so N−1≥4 and N=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.
1import redis2from redis.backoff import NoBackoff3from redis.retry import Retry45# The line that would have prevented the outage.6cache = redis.Redis(7 host="cache.internal",8 socket_connect_timeout=0.2, # p99 connect is 8 ms; 200 ms is generous9 socket_timeout=0.05, # p99 GET is 0.9 ms; 50 ms means "it is broken"10 retry=Retry(NoBackoff(), 0), # a cache is optional; do not retry it11 health_check_interval=30,12)1314def cached_predict(key, compute):15 try:16 hit = cache.get(key)17 if hit is not None:18 return decode(hit)19 except redis.RedisError:20 CACHE_ERRORS.inc() # degrade, never fail the request21 value = compute()22 try:23 cache.setex(key, 3600, encode(value))24 except redis.RedisError:25 pass26 return valueA 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=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] 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.
1import random, time23def call_with_retry(fn, attempts=3, base=0.25, cap=4.0, budget_s=8.0):4 deadline = time.monotonic() + budget_s5 for n in range(attempts):6 try:7 return fn(timeout=min(2.0, deadline - time.monotonic()))8 except Permanent: # 4xx: retrying cannot help9 raise10 except Transient:11 if n == attempts - 1 or time.monotonic() >= deadline:12 raise13 time.sleep(random.uniform(0, min(cap, base * 2 ** n))) # full jitterNote 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.
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 OPENQuantify 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,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.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.
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 servingGraceful 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.
| Failure | Naive behaviour | Degraded behaviour | User impact |
|---|---|---|---|
| Cache unavailable | 500 error | Compute directly | Slower: 22 ms → 45 ms |
| Large model timing out | 504 | Fall back to the distilled model | Accuracy 94.1% → 91.3% |
| Personalisation service down | 500 | Serve globally popular results | Less relevant, still useful |
| Feature store down | 500 | Use last-known features with a staleness flag | Slightly stale scores |
| Analytics sink down | Blocks the request | Drop to a local buffer, fire-and-forget | None |
| GPU fleet saturated | Everyone times out | Shed load: 429 low-priority tiers first | Some 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:
| Probe | Question | Failure action | Should it check dependencies? |
|---|---|---|---|
| Startup | Has initialisation finished? | Keep waiting (long timeout) | No |
| Liveness | Is the process wedged? | Restart the container | Never |
| Readiness | Can it serve a request right now? | Remove from the load balancer | Only 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.
1@app.get("/healthz") # liveness: no I/O, no locks, no dependencies2def healthz():3 return {"ok": True}45@app.get("/readyz") # readiness: mandatory dependencies only6def readyz():7 if not MODEL.loaded:8 return JSONResponse({"ready": False, "reason": "loading"}, 503)9 if not MODEL.warmed:10 return JSONResponse({"ready": False, "reason": "warming"}, 503)11 if SHUTTING_DOWN: # drain: fail readiness ~15 s before exiting12 return JSONResponse({"ready": False, "reason": "draining"}, 503)13 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.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
| Data | Consistency needed | Mechanism |
|---|---|---|
| Billing and usage counters | Strong | Primary writes, transactional, idempotency keys |
| Model registry / active version | Strong | Small, rarely written, read through a short-TTL cache |
| Prediction results | Eventual | Cache keyed by input + model version |
| Feature values | Bounded staleness | TTL plus a staleness flag on the response |
| Analytics and logs | Eventual, lossy acceptable | Async, buffered, never on the request path |
Deploying without breaking things
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-green | Canary | Rolling | |
|---|---|---|---|
| Blast radius of a bad release | 100% until rollback | Only the canary share | Grows as the roll proceeds |
| Rollback speed | Seconds (traffic flip) | Seconds (weight change) | Minutes (roll back through) |
| Extra capacity required | 2× during cutover | ~5–10% | ~10% |
| Detects subtle quality regressions | Poorly — all or nothing | Well — compare cohorts | Poorly |
| Cost for GPU fleets | High (double the GPUs) | Low | Low |
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.8 error-minutes. Under a 5% canary it is 10×0.05×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
┌───────────────────────────┐ 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.