Course Content
AI System Design and Architecture
3 sections · 7 lessons
Message Queues, Load Balancing, and Autoscaling
A document-processing API takes a PDF, runs OCR, embeds every page, and returns a summary. It takes 38 seconds. The team exposes it as a synchronous POST /analyse that blocks until done. It works in testing with three users.
On launch day 400 people upload documents in the first ten minutes. The load balancer's idle timeout is 60 seconds, so requests that take longer than that are killed mid-flight with a 504 — and the worker keeps burning GPU on a result nobody will receive. Users hit refresh, which starts a second job for the same document. The GPU fleet is now doing double work on a queue it cannot see. By 10:40 the p99 is 190 seconds and the error rate is 46%.
Nothing here is a model problem. It is a flow control problem: work arriving faster than it can be served, with no buffer to absorb the difference and no signal to add capacity. The three tools that fix it are a queue, a load balancer, and an autoscaler — and each of them fails in a specific way if you do not understand the arithmetic underneath.
Why the maths forces you to buffer
Start with one worker that takes 40 ms per inference. Its service rate is μ=1/0.040=25 requests per second. Requests arrive randomly at rate λ. For a single-server queue with random arrivals, the average time a job spends waiting before it starts is:
Put real numbers through it:
| Arrival rate λ | Utilisation ρ | Average queue wait Wq | Total time in system |
|---|---|---|---|
| 10 req/s | 40% | 0.40 / 15 = 26.7 ms | 66.7 ms |
| 17.5 req/s | 70% | 0.70 / 7.5 = 93.3 ms | 133 ms |
| 20 req/s | 80% | 0.80 / 5 = 160 ms | 200 ms |
| 22.5 req/s | 90% | 0.90 / 2.5 = 360 ms | 400 ms |
| 24 req/s | 96% | 0.96 / 1 = 960 ms | 1,000 ms |
| 24.75 req/s | 99% | 0.99 / 0.25 = 3,960 ms | 4,000 ms |
Read the last two rows again. Going from 96% to 99% utilisation — a 3% increase in traffic — multiplies wait time by four. This is why "our servers are only at 85% CPU, we have headroom" is one of the most dangerous sentences in operations.
Queue delay does not grow linearly with load. It grows as 1/(1−ρ), so the last 10% of capacity costs more delay than the first 90% combined.
Pooling: one queue beats many queues
Here is a result that changes how you wire systems. Suppose you have four workers, each μ = 25 req/s, and total demand of 80 req/s.
Option A — four separate queues, 20 req/s each (a load balancer that assigns a request to a fixed backend and lets it wait there). Each is an independent single-server queue at ρ = 0.8, so from the table above, Wq=160 ms.
Option B — one shared queue, four workers pulling from it. Now the offered load is a=λ/μ=80/25=3.2 and c = 4. The probability an arriving job has to wait at all comes from the Erlang C formula:
Working it through with a = 3.2, c = 4: the terms ak/k! for k = 0..3 are 1, 3.2, 5.12 and 5.4613, summing to 14.7813. The final term is 3.24/4!=104.8576/24=4.3691, multiplied by c/(c−a)=4/0.8=5, giving 21.8453. So
and the average wait is Wq=C/(cμ−λ)=0.5964/(100−80)=0.0298 s = 29.8 ms.
Same hardware, same traffic, same utilisation. 160 ms versus 29.8 ms — a 5.4× improvement purely from where the queue lives. The reason is simple: in Option A a job can be waiting behind a busy worker while another worker sits idle. In Option B that cannot happen.
Little's Law tells you the backlog: Lq=λWq=80×0.0298=2.4 jobs waiting on average. That number matters later, because it is what an autoscaler should watch.
Message queues
A queue is a durable buffer between producers and consumers. The producer writes a message and returns immediately; a consumer picks it up when it has capacity.
POST /jobs │ ▼ ┌────────────┐ enqueue ┌───────────────────────────┐ │ API (fast) │ ───────────► │ queue: jobs.ocr │ └────────────┘ │ [j7][j8][j9][j10][j11] │ │ └────────────┬──────────────┘ │ 202 Accepted │ pull (prefetch=1) │ {"job_id":"j11"} ┌────────┴────────┬─────────┐ ▼ ▼ ▼ ▼ client worker-1 worker-2 worker-3 │ │ │ └────────┬────────┴─────────┘ ▼ result store ◄── GET /jobs/j11The API's job becomes: validate, persist, enqueue, return 202 Accepted with a job id. That takes 6 ms instead of 38 seconds, so the 60-second load balancer timeout is no longer in play, and a refresh no longer duplicates work because the client already holds an id.
Delivery semantics, and why you need idempotency
Every queue offers one of three guarantees, and only one of them is cheap:
| Guarantee | What it means | Cost | Use when |
|---|---|---|---|
| At-most-once | Message may be lost, never duplicated | Cheapest | Metrics, non-critical telemetry |
| At-least-once | Never lost, may be delivered twice | Requires acks and redelivery | Almost all inference work |
| Effectively-once | No loss, no visible duplicate | Transactions or dedup state | Billing, anything that mutates money |
In practice you choose at-least-once and make the consumer idempotent. The pattern: derive a deterministic key (for example a SHA-256 of the input document plus the model version), and have the worker write results with that key as a unique constraint. A duplicate delivery then costs one wasted inference at worst, and never a duplicate charge or a duplicate email.
This is where people get it wrong: they set the visibility timeout — the window a message stays hidden after being picked up — shorter than the actual processing time. A 38-second job with a 30-second visibility timeout is redelivered while still running, so two workers process it, and under load this compounds into a self-inflicted traffic multiplier. Set the timeout to at least p99 processing time, and heartbeat-extend it for long jobs.
Choosing a broker
| RabbitMQ | Apache Kafka | Redis Streams | |
|---|---|---|---|
| Model | Broker routes messages to queues; consumers ack | Append-only partitioned log; consumers track offsets | In-memory log with consumer groups |
| Message removed after read? | Yes, on ack | No — retained for a configured window | No, but trimmed by length or age |
| Replay history | No | Yes — rewind to any offset | Partially, within retention |
| Ordering | Per queue, lost with competing consumers | Strict per partition | Per stream |
| Routing flexibility | Very high (direct, topic, fanout, headers) | Low — topic and partition key | Low |
| Realistic throughput | Tens of thousands msg/s | Millions msg/s | Hundreds of thousands msg/s |
| Operational weight | Moderate | Heavy (brokers, KRaft controllers, partitions; Kafka 4 removed ZooKeeper) | Light if you already run Redis |
| Best for | Task queues with complex routing and per-message retry | Event streams, feature pipelines, multiple independent consumers | Small-to-mid task queues, low latency, minimal ops |
For an inference task queue, RabbitMQ or a managed equivalent is usually the right default. Reach for Kafka when several unrelated consumers need the same stream — for example, when a "document uploaded" event must trigger OCR, index refresh, and an analytics rollup, each at its own pace and each able to replay after a bug fix.
Three patterns you will actually use
Work queue (competing consumers). One queue, N workers, each message handled exactly once by one worker. This is the shape that gives you the 29.8 ms pooling result above. Set prefetch = 1 for long, variable jobs so a worker cannot hoard messages it will not get to for two minutes.
1import pika, hashlib, json23conn = pika.BlockingConnection(pika.ConnectionParameters("rabbit"))4ch = conn.channel()5ch.queue_declare(queue="jobs.ocr", durable=True,6 arguments={"x-dead-letter-exchange": "dlx.ocr"})7ch.basic_qos(prefetch_count=1) # one in flight per worker89def handle(chan, method, props, body):10 job = json.loads(body)11 key = hashlib.sha256((job["doc_url"] + job["model_version"]).encode()).hexdigest()12 if results.exists(key): # idempotent: already done13 chan.basic_ack(method.delivery_tag)14 return15 try:16 results.put(key, run_ocr(job["doc_url"]))17 chan.basic_ack(method.delivery_tag)18 except TransientError:19 chan.basic_nack(method.delivery_tag, requeue=True)20 except PermanentError:21 chan.basic_nack(method.delivery_tag, requeue=False) # → dead letter2223ch.basic_consume("jobs.ocr", handle)24ch.start_consuming()Note the dead-letter exchange. Without it, a message that always fails — a corrupt PDF, say — is requeued forever and can occupy a worker permanently. That single misconfiguration has taken down more pipelines than any model bug.
Publish/subscribe. One event, many independent consumers, each with its own copy. Use it when the producer must not know who cares.
┌─► embeddings-indexer (own offset) doc.uploaded ──────►├─► thumbnail-generator (own offset) └─► analytics-rollup (own offset)Request/reply. Asynchronous transport, synchronous feel: the caller sends a message with a reply_to queue and a correlation id, then waits. Useful for putting a queue's back-pressure and batching in front of a GPU while keeping a simple call-and-wait API. The trap is that you have reintroduced a synchronous timeout, so you must cap the wait and fall back to a job id when it expires.
Load balancing
A load balancer decides which of N healthy backends receives each request. For stateless web servers with uniform 5 ms responses, the algorithm barely matters. For inference it matters enormously, because service times are heavy-tailed: a 50-token prompt might take 60 ms and a 4,000-token prompt 2,400 ms on the same model. A 40× spread breaks the assumption every simple algorithm rests on.
| Strategy | How it picks | Behaviour with 40× variable service times | Verdict for inference |
|---|---|---|---|
| Round robin | Next in rotation | Sends a request to a backend already 2.4 s deep in work while another is idle | Avoid |
| Least connections | Fewest open connections | Good proxy for "least busy" when one connection = one in-flight request | Good |
| Least outstanding requests | Fewest in-flight requests | Directly approximates the shared-queue result — the best simple choice | Best default |
| Weighted | Proportional to declared capacity | Essential for mixed fleets (an A100 node should get 3× a T4 node) | Combine with least-outstanding |
| IP hash / consistent hash | Deterministic by client or key | Maximises cache hits; unbalances badly if keys are skewed | Only for cache or session affinity |
| Least response time | Lowest EWMA of recent latency | Adapts to slow nodes, but can oscillate if the window is short | Good with a long smoothing window |
Quantify the round-robin failure. Two backends, requests arriving every 100 ms. A 2,400 ms request lands on backend A. Under round robin, A also receives requests 3, 5, 7, 9… so twelve subsequent requests queue behind it, each waiting an average of about 1.2 s — while backend B idles between its 60 ms jobs. Under least-outstanding-requests, A is skipped until it drains, and those twelve requests see roughly 60 ms. Same fleet, same traffic, a 20× difference in p99 caused only by the balancing rule.
Where the balancer lives
| Type | Examples | Strengths | Weaknesses |
|---|---|---|---|
| Hardware appliance | F5 BIG-IP, Citrix ADC | Very high throughput, TLS offload in silicon | Capital cost, slow to change, poor fit for elastic fleets |
| Software | NGINX, HAProxy, Envoy | Cheap, scriptable, rich algorithms, runs as a sidecar | You operate it, patch it, and scale it |
| Cloud managed | AWS ALB/NLB, GCP LB | Autoscales itself, health checks, integrated certificates | Limited algorithm choice, per-GB and per-LCU charges, opaque under stress |
One AI-specific gotcha: cloud application load balancers typically default to round robin and a 60-second idle timeout. Both defaults are wrong for inference. Turn on least-outstanding-requests and raise the idle timeout above your p99 — or better, put the long work behind a queue so the timeout stops mattering.
Autoscaling
Autoscaling adds and removes replicas in response to a signal. The whole discipline is choosing the right signal and the right time constants.
The signal
| Metric | Why it is tempting | Why it misleads for AI |
|---|---|---|
| CPU utilisation | Universal, free | A GPU-bound worker may sit at 12% CPU while the GPU is saturated. CPU says "idle", users say "slow". |
| Memory | Easy to read | Model weights make it a flat constant — 17 GB whether serving 0 or 500 req/s. |
| Requests per second | Directly proportional to load | Only valid if per-request cost is uniform, which for variable-length prompts it is not. |
| GPU utilisation | Closest to the bottleneck | Good, but a batching server can show 95% while latency is still fine. |
| Queue depth per replica | Measures unmet demand directly | The right default: it is the only metric that tells you a user is waiting. |
Derive the target from your SLA rather than guessing. If a worker serves 25 jobs/s and you want queue wait under 200 ms, then by Little's Law the acceptable backlog per worker is 25×0.2=5 jobs. So the scaling rule is: keep queue_depth / replicas ≤ 5, and desired replicas is ceil(queue_depth / 5).
apiVersion: autoscaling/v2kind: HorizontalPodAutoscalermetadata: name: ocr-workersspec: scaleTargetRef: {apiVersion: apps/v1, kind: Deployment, name: ocr-worker} minReplicas: 2 maxReplicas: 40 metrics: - type: External external: metric: {name: rabbitmq_queue_messages_ready} target: type: AverageValue averageValue: "5" # 5 queued jobs per replica behavior: scaleUp: stabilizationWindowSeconds: 30 # react fast policies: [{type: Percent, value: 100, periodSeconds: 60}] scaleDown: stabilizationWindowSeconds: 600 # retreat slowly policies: [{type: Pods, value: 2, periodSeconds: 120}]The asymmetry between the two stabilisation windows is the single most important line in that file. Scaling up late costs you an SLA breach; scaling down early costs you an SLA breach and a cold start. Be eager to grow and reluctant to shrink.
Cold start: the number that decides everything
Work out what a scale-up actually costs on a GPU fleet:
| Stage | Time |
|---|---|
| Metric scrape + evaluation window | 60 s |
| Cloud provisions a GPU node | 90 s |
| Pull 8 GB image at 250 MB/s | 32 s |
| Load weights to GPU, warm up CUDA kernels | 45 s |
| Total from breach to serving | 227 s ≈ 3.8 minutes |
Now suppose traffic jumps from 80 to 200 req/s with 4 workers (capacity 100 req/s) in place. For 227 seconds you accumulate a deficit of 100 req/s, so the backlog reaches 100×227=22,700 jobs. Scaling to exactly 8 workers (200 req/s) would hold the backlog steady forever — it never drains. To clear it you must overshoot: 12 workers gives 300 req/s, a drain rate of 100/s, and 22,700/100=227 s to recover. Total time outside SLA: about 7.6 minutes.
Autoscaling does not remove the need for headroom; it removes the need for permanent headroom. You still have to survive the provisioning window on the capacity you already have.
Three ways to shrink that 227 s, in order of value: bake weights into the image or a pre-warmed volume rather than downloading at boot (−40 s or more); keep a small pool of pre-provisioned warm nodes (−122 s, at the cost of paying for idle GPUs); and scale on a leading indicator such as upstream request arrivals rather than a lagging one such as queue depth, which buys back most of the 60 s evaluation window.
Thrashing
Thrashing is oscillation: scale out, metric drops below target, scale in, metric spikes, scale out again. Each cycle pays a full cold start and throws away warm cache. It happens when the scale-in threshold is close to the scale-out threshold and the cooldown is short.
The fix is hysteresis with real numbers. Scale out at 5 queued jobs per replica; scale in only when the figure has stayed below 2 for ten consecutive minutes; never remove more than two replicas per two minutes; never go below a floor that covers your normal trough. Step scaling helps too — react proportionally to how far out of bounds you are:
| Queue depth per replica | Action |
|---|---|
| 0 – 2 (sustained 10 min) | Remove 2 replicas |
| 2 – 5 | Hold |
| 5 – 15 | Add 25% of current replicas |
| 15 – 40 | Add 100% of current replicas |
| > 40 | Jump to maxReplicas, page the on-call engineer |
The three parts wired together
┌──────────────────┐ clients ────────►│ Cloud LB (ALB) │ least-outstanding-requests └────────┬─────────┘ ▼ ┌────────────────────┐ │ API pods (×6) │ validate → enqueue → 202 │ HPA on RPS │ └─────────┬──────────┘ ▼ ┌────────────────────────────────────────────┐ │ jobs.ocr depth=12 ── DLQ ─► alert│ └────────────────────────┬───────────────────┘ │ prefetch=1 ┌────────────────────────┴────────────────────────┐ ▼ ▼ ▼ ┌───────────┐ ┌───────────┐ ┌───────────┐ │ worker-1 │ │ worker-2 │ ... │ worker-N │ │ GPU │ │ GPU │ │ GPU │ └─────┬─────┘ └─────┬─────┘ └─────┬─────┘ └────────────────────────┴────────────────────────┘ ▼ ┌──────────────────┐ │ result store │◄── GET /jobs/{id} └──────────────────┘ HPA target: queue_depth / replicas ≤ 5 (= 200 ms wait at μ=25/s)Each piece does one job. The load balancer spreads short synchronous work evenly. The queue absorbs the difference between arrival rate and service rate, so a burst becomes latency instead of errors. The autoscaler converts sustained backlog into capacity. Remove any one of them and a specific failure returns: no queue and bursts become 504s; no least-outstanding balancing and one slow request poisons a backend; no autoscaler and the backlog grows without bound.
What this means when you build one
Before you write any of it, measure two numbers on a single worker: service time at p50 and at p99. Everything above is derived from them. μ comes from p50, your visibility timeout comes from p99, your autoscaler target comes from p50 and your latency SLA, and your maximum safe utilisation comes from how far apart p50 and p99 are — the wider the spread, the lower the ρ you can run at.
Then set three guard rails that people routinely forget. Give every queue a dead-letter destination and alert when anything lands in it, or a single poison message will silently eat a worker. Cap maxReplicas at a number whose bill you have actually calculated — a runaway autoscaler responding to a retry storm can quadruple your spend in twenty minutes, and 40 g5.xlarge instances is USD 40.24 an hour, or USD 966 a day if nobody notices. And make every consumer idempotent from the first commit, because at-least-once delivery means duplicates are not an edge case, they are Tuesday.
Finally, load-test the scale-up path, not just the steady state. A system that handles 200 req/s comfortably and a system that can get from 80 to 200 req/s without a seven-minute breach are different systems, and only one of them survives a launch.