AI System Design and Architecture

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.

The three pieces, wired the way they have to beClient — POST returns a job id in 6 msLoad balancer — least outstanding requestsQueue — one shared pool, depth is the signalWorkers — OCR, embed, summarise, 38 s eachResult store — polled, or delivered by webhook
Queue depth, not CPU, is the scaling signal: it already accounts for how long each item takes to clear.

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\mu = 1/0.040 = 25 requests per second. Requests arrive randomly at rate λ\lambda. For a single-server queue with random arrivals, the average time a job spends waiting before it starts is:

Wq=ρμ−λ,ρ=λμW_q = \frac{\rho}{\mu - \lambda}, \qquad \rho = \frac{\lambda}{\mu}

Put real numbers through it:

Arrival rate λUtilisation ρAverage queue wait WqTotal time in system
10 req/s40%0.40 / 15 = 26.7 ms66.7 ms
17.5 req/s70%0.70 / 7.5 = 93.3 ms133 ms
20 req/s80%0.80 / 5 = 160 ms200 ms
22.5 req/s90%0.90 / 2.5 = 360 ms400 ms
24 req/s96%0.96 / 1 = 960 ms1,000 ms
24.75 req/s99%0.99 / 0.25 = 3,960 ms4,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=160W_q = 160 ms.

Option B — one shared queue, four workers pulling from it. Now the offered load is a=λ/μ=80/25=3.2a = \lambda/\mu = 80/25 = 3.2 and c = 4. The probability an arriving job has to wait at all comes from the Erlang C formula:

C(c,a)=acc!⋅cc−a∑k=0c−1akk!+acc!⋅cc−aC(c,a) = \frac{\frac{a^c}{c!}\cdot\frac{c}{c-a}}{\sum_{k=0}^{c-1}\frac{a^k}{k!} + \frac{a^c}{c!}\cdot\frac{c}{c-a}}

Working it through with a = 3.2, c = 4: the terms ak/k!a^k/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.36913.2^4/4! = 104.8576/24 = 4.3691, multiplied by c/(c−a)=4/0.8=5c/(c-a) = 4/0.8 = 5, giving 21.8453. So

C(4,3.2)=21.845314.7813+21.8453=0.5964C(4, 3.2) = \frac{21.8453}{14.7813 + 21.8453} = 0.5964

and the average wait is Wq=C/(cμ−λ)=0.5964/(100−80)=0.0298W_q = C/(c\mu - \lambda) = 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.4L_q = \lambda W_q = 80 \times 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.

Text
  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/j11

The 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:

GuaranteeWhat it meansCostUse when
At-most-onceMessage may be lost, never duplicatedCheapestMetrics, non-critical telemetry
At-least-onceNever lost, may be delivered twiceRequires acks and redeliveryAlmost all inference work
Effectively-onceNo loss, no visible duplicateTransactions or dedup stateBilling, 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

RabbitMQApache KafkaRedis Streams
ModelBroker routes messages to queues; consumers ackAppend-only partitioned log; consumers track offsetsIn-memory log with consumer groups
Message removed after read?Yes, on ackNo — retained for a configured windowNo, but trimmed by length or age
Replay historyNoYes — rewind to any offsetPartially, within retention
OrderingPer queue, lost with competing consumersStrict per partitionPer stream
Routing flexibilityVery high (direct, topic, fanout, headers)Low — topic and partition keyLow
Realistic throughputTens of thousands msg/sMillions msg/sHundreds of thousands msg/s
Operational weightModerateHeavy (brokers, KRaft controllers, partitions; Kafka 4 removed ZooKeeper)Light if you already run Redis
Best forTask queues with complex routing and per-message retryEvent streams, feature pipelines, multiple independent consumersSmall-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.

Python
import pika, hashlib, jsonconn = pika.BlockingConnection(pika.ConnectionParameters("rabbit"))ch = conn.channel()ch.queue_declare(queue="jobs.ocr", durable=True,                 arguments={"x-dead-letter-exchange": "dlx.ocr"})ch.basic_qos(prefetch_count=1)   # one in flight per workerdef handle(chan, method, props, body):    job = json.loads(body)    key = hashlib.sha256((job["doc_url"] + job["model_version"]).encode()).hexdigest()    if results.exists(key):              # idempotent: already done        chan.basic_ack(method.delivery_tag)        return    try:        results.put(key, run_ocr(job["doc_url"]))        chan.basic_ack(method.delivery_tag)    except TransientError:        chan.basic_nack(method.delivery_tag, requeue=True)    except PermanentError:        chan.basic_nack(method.delivery_tag, requeue=False)  # → dead letterch.basic_consume("jobs.ocr", handle)ch.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.

Text
                      ┌─► 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.

StrategyHow it picksBehaviour with 40× variable service timesVerdict for inference
Round robinNext in rotationSends a request to a backend already 2.4 s deep in work while another is idleAvoid
Least connectionsFewest open connectionsGood proxy for "least busy" when one connection = one in-flight requestGood
Least outstanding requestsFewest in-flight requestsDirectly approximates the shared-queue result — the best simple choiceBest default
WeightedProportional to declared capacityEssential for mixed fleets (an A100 node should get 3× a T4 node)Combine with least-outstanding
IP hash / consistent hashDeterministic by client or keyMaximises cache hits; unbalances badly if keys are skewedOnly for cache or session affinity
Least response timeLowest EWMA of recent latencyAdapts to slow nodes, but can oscillate if the window is shortGood 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

TypeExamplesStrengthsWeaknesses
Hardware applianceF5 BIG-IP, Citrix ADCVery high throughput, TLS offload in siliconCapital cost, slow to change, poor fit for elastic fleets
SoftwareNGINX, HAProxy, EnvoyCheap, scriptable, rich algorithms, runs as a sidecarYou operate it, patch it, and scale it
Cloud managedAWS ALB/NLB, GCP LBAutoscales itself, health checks, integrated certificatesLimited 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

MetricWhy it is temptingWhy it misleads for AI
CPU utilisationUniversal, freeA GPU-bound worker may sit at 12% CPU while the GPU is saturated. CPU says "idle", users say "slow".
MemoryEasy to readModel weights make it a flat constant — 17 GB whether serving 0 or 500 req/s.
Requests per secondDirectly proportional to loadOnly valid if per-request cost is uniform, which for variable-length prompts it is not.
GPU utilisationClosest to the bottleneckGood, but a batching server can show 95% while latency is still fine.
Queue depth per replicaMeasures unmet demand directlyThe 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=525 \times 0.2 = 5 jobs. So the scaling rule is: keep queue_depth / replicas ≤ 5, and desired replicas is ceil(queue_depth / 5).

Text
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:

StageTime
Metric scrape + evaluation window60 s
Cloud provisions a GPU node90 s
Pull 8 GB image at 250 MB/s32 s
Load weights to GPU, warm up CUDA kernels45 s
Total from breach to serving227 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,700100 \times 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=22722{,}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 replicaAction
0 – 2 (sustained 10 min)Remove 2 replicas
2 – 5Hold
5 – 15Add 25% of current replicas
15 – 40Add 100% of current replicas
> 40Jump to maxReplicas, page the on-call engineer

The three parts wired together

Text
                    ┌──────────────────┐   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.