Task queues vs logs, producer/consumer design, retries and backoff, leases, DLQs, lag/backpressure, and LLM worker operations.
Lesson 4 · Async queues
Workers, retries, backoff, DLQ, poison messages
Async is how LLM features stay interactive
Provider calls, embeddings, eval suites, and webhook fan-out do not belong on the user request thread. Queues buy isolation, retries, and backpressure — and introduce at-least-once delivery, lag, and poison messages. Interviewers grade whether you operate workers, not whether you can import Celery.
Sync path: validate, authorize, persist intent, enqueue, return 202. Async path: worker leases message, does side effects, acks. If you do heavy work inline, deploys, timeouts, and traffic spikes all become user-visible outages.
Task queue vs event log
Task/work queue (SQS, cloud tasks, Sidekiq, BullMQ): each message is a unit of work; competing consumers; success deletes/acks. Event log (Kafka): durable ordered stream, multiple consumer groups, replay. Use task queues for “do this job”; logs when many independent subscribers need the same history.
text
1PRODUCER (API)2 1. AuthZ + validate3 2. INSERT run status=queued (TX)4 3. Outbox/enqueue run_id5 4. 202 + run_id67CONSUMER (worker)8 1. Receive/lease message9 2. Idempotency check (run_id processed?)10 3. Transition queued→running (conditional)11 4. Call provider / do work12 5. Persist result + status13 6. Ack/delete message14On failure: retry with backoff or route to DLQ after N attempts
Retries and backoff
Retry only transient failures: 429/503, network blips, deadlocks. Do not infinite-retry permanent 400s or bad payloads — fix or DLQ. Use exponential backoff with jitter so thundering herds do not stampede the provider. Cap attempts; page on DLQ depth.
SQS-style visibility: message hidden while worker works; if worker dies, message reappears. Set timeout above p99 work duration or implement heartbeat extension. Too short → duplicate workers; too long → slow recovery after crashes.
Dead-letter queues and poison messages
A poison message crashes consumers repeatedly (bad payload, null deref, huge blob). After N failures, move to DLQ so the pipeline breathes. Ops: alarm on DLQ depth, runbook to inspect, fix, and replay. Ignoring DLQ is how async systems silently rot.
Concurrency, lag, and backpressure
Measure queue lag (oldest message age) and processing rate. Scale workers on lag, not only CPU. Backpressure: if provider 429s, slow consumers (token bucket) rather than spawning infinite workers. Prefer bounded concurrency per tenant so one whale cannot starve everyone.
Ordering and fairness
Global strict order kills throughput. Order only where required (per run_id, per account). FIFO queues have cost and throughput limits — justify them. Fairness: shuffle tenants or use per-tenant queues/rate limits so noisy neighbors do not pin the world.
LLM worker specifics
Provider timeouts are long; set HTTP client timeouts deliberately. Cancel: cooperative cancellation when user aborts run — stop paying for tokens. Budget: max concurrent provider calls; circuit break on elevated error rates. Store token usage even on failure when the API returns it.
Separate queues by urgency when needed: interactive “user waiting” vs bulk backfill. One shared queue means a 100k-embedding job starves chat. Priority queues or distinct worker pools encode product SLOs in infrastructure — name that tradeoff in interviews.
Poison vs transient — classification table
Build an error taxonomy in code: network/timeouts → retry; 429 → retry-after; 401 to provider → page (bad secret); 400 validation → terminal fail the run; unhandled exception → retry once then DLQ with stack fingerprint. Without classification, every failure looks the same and ops cannot act.
text
1RETRY CLASS ACTION2-------------------- ----------------------------------3Timeout / 503 backoff retry4429 honor Retry-After + jitter5401/403 to provider terminal + page (config/secret)6400/422 user input terminal fail run (no DLQ spam)7Worker bug (5× crash) DLQ + alert + fingerprint8Cancel requested stop; status=cancelled; ack
Deploying workers safely
Rolling deploys: stop consuming, finish in-flight, then exit (graceful shutdown). Drain period must fit max job time or you need checkpoint/resume. Version workers and messages when payloads evolve — old workers must not hard-crash on new fields (forward compatible consumers).
Include a sweeper for runs stuck in running past a deadline (worker died after lease extension bugs). Sweepers requeue or fail explicitly — silent stuck rows are the async equivalent of a hung HTTP worker with no timeout.
Observability hooks (preview of L7)
Emit: enqueue rate, process rate, lag age, retry count histogram, DLQ depth, per-tenant concurrency. Log run_id on every worker step. Without lag SLOs, you only notice async failure when a human complains the button “spun forever.”
Interview answers — queues & workers
01Q: Why async? Isolate long/unreliable work from interactive p99; enable retries and scaling workers independently of API replicas.
02Q: At-least-once? Default for most queues — design idempotent consumers; exactly-once is usually effect-once via dedupe.
03Q: Retry policy? Transient only, exponential backoff + jitter, max attempts, then DLQ; never retry poison without a code fix.
04Q: Task queue vs Kafka? Jobs/work units → task queue; multi-subscriber replayable facts → log. Do not force Kafka for email sends without need.
05Q: Visibility timeout? Must exceed work duration or use heartbeats; too low causes duplicate processing, too high delays redelivery after death.
06Q: Poison message? Crashes/loop fails; quarantine to DLQ; alert; fix consumer or payload; replay carefully.
07Q: Lag SLO? Age of oldest message / time-to-start work; scale workers; watch provider limits as the real bottleneck.
08Q: Per-tenant fairness? Concurrency caps and fair scheduling so one tenant’s bulk job cannot starve others.
09Q: Outbox? Write event in the same DB transaction as state change; publisher emits to queue — avoids lost enqueues.
10Q: When not to queue? Tiny pure CPU work that must be synchronous for UX and cannot usefully retry independently — do not add hops for fashion.
A worker fails calling the LLM provider with HTTP 503. What is the right default?
AAck the message and mark the run succeeded so the queue stays empty.BNack/retry with exponential backoff and jitter; after max attempts send to DLQ and mark run failed/retryable per policy.CSpin up infinite workers until one succeeds immediately.
AAny message larger than 1KB that should have been stored in S3 instead.BA message that repeatedly fails processing (bad payload/bug) and must be isolated (DLQ) so other work proceeds.CA message encrypted with a key the security team likes.
Visibility timeout is 30s but jobs take 2–5 minutes. What goes wrong?
ANothing — visibility timeout only affects metrics dashboards.BAnother worker may pick up the same message while the first is still running → duplicate side effects unless idempotent.CThe queue permanently deletes messages after 30s always.
When do you choose a durable log (Kafka-style) over a simple work queue?
AAlways — logs are strictly better for every email send and thumbnail job.BWhen multiple independent consumers need the same ordered history and replay (analytics, search index, audit projections).COnly when you do not care about durability at all.
API returns 202 after enqueue, but the process crashes before the message hits SQS (no outbox). User sees a run that never starts. Root issue?
AMissing transactional outbox / atomic intent+enqueue — dual write failure between DB state and queue.BHTTP 202 is illegal for LLM products so clients ignore the run id.CWorkers should poll the entire runs table every 10ms instead of using a queue forever.