Skip to main content

Command Palette

Search for a command to run...

Do It Later, On Purpose: Async Processing and Queues, Explained Like You're New

Blueprints of Scale — Core Concepts #10

Updated
•53 min read•View as Markdown

In the URL shortener post, the interviewer asked a question I answered in one word: when someone clicks a short link, how do you record the click analytics without slowing down the redirect? "Asynchronously," I said, and moved on. It was the right answer and a completely unexplained one. This post is the explanation I owed you.

Here's the idea in full: most of the work your system does doesn't need to happen while the user is waiting. The charge needs to go through now. The confirmation page needs to render now. But the receipt email, the analytics event, the search index update, the thumbnail render, the fraud check, the warehouse notification: none of those need the user's attention, and doing them inline is what makes the spinner spin. Async processing is the discipline of separating "must happen now" from "must happen eventually," and building the machinery that guarantees "eventually." That machinery is the queue, and it's one of the highest-leverage ideas in backend engineering: it makes systems faster for the user, tougher under spikes, and simpler to reason about at scale. It also has failure modes that will page you if you skip the boring parts. This post covers both.

Here's what's covered: the seven-second checkout and the one question that fixes it; what "async" actually means (it's about whose time you're spending); the queue as a shock absorber, and the shapes queues come in (point-to-point, pub/sub fan-out, and the log); why delivery is at-least-once, why your consumers must be idempotent, and where the at-most-once switch actually lives; the dual-write problem, the transactional outbox, and the relay's own failure modes; workers: parallelism and its ceiling, ordering versus throughput, poison messages, dead-letter queues, priorities, batching, big payloads, and noisy tenants; delayed and scheduled jobs, distributed locks, and retries with backoff; backpressure, properly defined, and the three things you can do when producers outrun consumers; and the principal-level toolkit: choreography versus orchestration and the engines that do it, the lag arithmetic, event-driven pitfalls (ordering, replay, schema evolution), what async does to the user experience, tracing across queues, and when not to go async at all.

Sections 1 and 2 assume nothing. Sections 3 through 8 are the machinery every backend engineer meets in production. Section 9 is the judgment. The cheat sheet is at the end under Do It Later, On Purpose, distilled, and every diagram is also described in the text around it.


Section 1 — The seven-second checkout

An online store's "Place Order" button takes seven seconds to respond. Every time, not just sometimes. The team profiles it and finds the checkout handler doing its work like a diligent clerk with no sense of priority: charge the card (1.2 seconds), reserve the inventory (0.8), send the confirmation email (1.5), send the receipt email (1.5), update the analytics pipeline (1.0), add the loyalty points (1.2). Total: 7.2 seconds of spinner. Every step correct. Every step necessary. Every step happening while the user stares at the screen.

A seven-step synchronous checkout taking 7.2 seconds: charge card, reserve inventory, emails, analytics, loyalty — all before responding

One picture, and it's the whole problem: a chain where every link is someone else's latency. And it gets worse than slow. Each of those six steps is a dependency (the email service, the analytics pipeline, the loyalty system), which means the checkout now fails when any of them fails. The loyalty service has a bad deploy? Checkout is down. The email provider is slow today? Every customer waits.

Synchronous work doesn't just add latency. It multiplies your blast radius.

The resilience post (#2) called this out: every inline dependency is a way to be down.

Now the question this entire post is built on: which of those six steps does the user actually need before they see "order placed"? The charge, obviously; you can't confirm an order you haven't taken money for. The inventory reservation, yes; you can't sell what you don't have. The other four? The user will never know whether the confirmation email sent in 400 milliseconds or 4 minutes. The analytics pipeline definitely doesn't care. The loyalty points can land whenever. Four of the six steps are hostages to the user's attention for no reason.

The fix the team shipped: the request handler does the charge and the inventory reservation, the two things the answer depends on, drops four messages onto a queue ("send confirmation email," "send receipt," "record analytics," "add loyalty points"), and returns. Checkout went from 7.2 seconds to 2.0. The emails still send. The analytics still update. The loyalty points still land. They just happen after the user has their answer. (The remaining two seconds are the payment provider's and the inventory system's, and they're a different kind of problem: those are real synchronous dependencies, and making them faster is the estimation post's (#13) and the resilience post's (#2) business, not this one's.)

Nothing got faster. The user just stopped waiting for it.

There's a lesson hiding in that last line: latency has an audience. Work the user waits for has a budget measured in hundreds of milliseconds. Work nobody waits for has a budget measured in minutes, and a much bigger budget means much cheaper engineering. The rest of this post is the machinery that makes "later" reliable.


Section 2 — The one question: does the user need the result now?

Every piece of work your system does, run it through one question: does the user need the result now? Not "is it important" (the receipt email is important). Not "is it part of the flow" (the analytics event is part of the flow). The question is strictly about the response: does the answer I'm about to send depend on this work being done?

The one question: if the user needs the result now, do it inline; otherwise put the work on the queue and answer fast

The "yes" bucket, the truly synchronous work, is smaller than most codebases suggest: authentication (you can't serve the request without knowing who it's for), the decision the response depends on (is the seat available? did the charge clear? is the username taken?), and reads (fetching data to display is inherently synchronous; the page needs it). The "no" bucket is everything else, and it's enormous: notifications of every kind, analytics and metrics, search index updates, thumbnail and preview rendering, cache warming, webhooks to third parties, reports and exports, fraud and abuse scoring that doesn't gate the response, data warehouse syncs.

A week later, the same team from Section 1 did the exercise properly, to check the fix they'd shipped in a hurry. They listed every side effect in the checkout flow and asked the question about each one, in a room, with the product manager present. Charge: yes. Inventory: yes. Confirmation email: no. Receipt: no. Analytics: no. Loyalty: no. Fraud check, which the first pass had missed entirely: interesting. The fraud score didn't gate the order, it gated the shipment, which happened hours later. So: no. Seven items, two yeses. The exercise took twenty minutes, confirmed the five seconds they'd already saved, and found one more item to move. Most "slow endpoints" are slow the same way: not because any step is slow, but because steps the user doesn't need are spending the user's time.

Underneath is a useful mental model: every request has a latency budget, and synchronous work spends it. The classic thresholds from usability research (Jakob Nielsen's, from 1993, and still cited) are that 100 milliseconds feels instant, one second keeps the user's flow of thought intact, and ten seconds is the limit of attention. Most teams budget a few hundred milliseconds for an interactive response, and every inline dependency, every email send, every analytics call, draws from the same small account. Async work draws from a different account entirely, one measured in minutes, with effectively infinite funds.

One caveat before the machinery: "later" has to actually happen, reliably and observably, and that's Sections 3 through 8. "Fire and forget" without the machinery is just "forget." The question tells you what to defer; the queue tells you how to keep the promise.

Async is about moving work to the budget where it's cheap, not about making it faster.


Section 3 — The queue as a shock absorber

A queue is embarrassingly simple, which is why it works. Three roles. Producers put messages on the queue ("send this email," "process this image") and move on; they don't know or care who handles the message or when. The broker holds the messages until someone takes them, and a properly configured broker holds them durably, on disk and replicated, so a crashed broker doesn't lose them (that "properly configured" is doing work: RabbitMQ's transient queues, a Kafka topic with a replication factor of one, or a Redis list with default persistence are not durable, whatever the marketing says). Consumers, also called workers, pull messages off, do the work, and acknowledge completion (an "ack" is the consumer telling the broker "done, you can delete this"). The producer and the consumer never meet. They don't even have to exist at the same time: the producer can be long gone before the consumer wakes up. The queue decouples work in time and in knowledge, and that decoupling is what absorbs shocks.

The queue as a shock absorber: a 100x traffic spike buffers 100,000 messages per minute while ten workers drain 2,000 each per minute

Take the concert ticket on-sale: a hundred times normal traffic at 10 AM sharp, 100,000 messages a minute for ten minutes, a million messages. Without a queue, that spike hits your workers directly: they saturate, requests time out, the on-sale collapses, and angry fans trend on social media. With a queue, the spike hits the broker, which just holds messages. It's very good at holding messages; that's its whole job. Ten workers at 2,000 a minute each chew through the pile at their sustainable pace, and the last confirmation email goes out about fifty minutes later. Nobody at the concert cares that their email arrived at 10:50 instead of 10:05. The queue converted a capacity emergency into a latency bill, and for background work, latency is the cheapest currency there is.

Queues also act as bulkheads, the resilience post's term for walls that stop one failure from sinking the ship. If the email service goes down for an hour, the messages wait in the queue instead of failing the checkout. When the service recovers, the workers drain the backlog. The producer never knew there was an outage. Without the queue, a downstream outage is your outage. With it, it's a delay.

Now, "queue" covers three different shapes, and picking wrong causes real pain:

Point-to-point (SQS standard and FIFO queues, RabbitMQ's classic and quorum queues, Azure Service Bus queues): a message is delivered to one consumer, and once acknowledged, it's gone. This is the work-distribution shape, sometimes called competing consumers: a pool of workers pulling tasks, each task done once (well, at-least-once; Section 4). Use it when the work is tasks: send emails, render thumbnails, process uploads. Parallelism is simple: add workers.

Publish/subscribe (SNS fanning out to SQS queues, RabbitMQ exchanges with several bound queues, Service Bus topics with subscriptions, Google Pub/Sub): one message, many readers, each getting its own copy. "Order placed" goes to the warehouse's queue, the email service's queue, and the analytics queue at once, and each consumes at its own pace. This is the fan-out shape, and it's how you get "many readers" without a log. What it doesn't give you, by default, is history: once each subscriber has consumed its copy, the message is gone (some, like Google Pub/Sub, can retain and replay if you turn it on).

Log-based (Kafka, Redpanda, Pulsar, Kinesis, and RabbitMQ's newer streams): the broker is an append-only log. Messages are never deleted on read; they're kept for a retention period (Kafka's default is seven days, configurable per topic), and each consumer group tracks its own position (its offset) in the log independently. Five different systems can each read the same events at their own pace, and a new consumer can rewind and replay history from last Tuesday. This is the event shape: "order placed" is a fact that the warehouse, the analytics team, and the fraud system all want to hear about, separately, and that a system built next year will want to hear about retroactively. The price is that parallelism is capped by the number of partitions the topic was created with (each partition is read by one consumer in a group; extra consumers sit idle), changing the partition count is a migration, and on the classic protocol every time a consumer joins or leaves the group, the group pauses to rebalance (Kafka 4.0's incremental protocol removes most of that pause).

The way teams learn the difference is by choosing point-to-point for events. Order events go into a classic RabbitMQ queue; then the analytics team wants the same events, but they've been consumed and deleted; then someone wants to reprocess last month after a bug fix, and the messages are gone. The migration that follows moves the event streams to a log and keeps the queue for the actual tasks (send the email, charge the card), and everyone stops fighting the tool.

As a default: tasks go point-to-point, fan-out goes pub/sub, facts go on the log. It's a default, not a law; Kafka gets used for task queues at plenty of companies and now has share groups built for it. What matters is that when someone proposes using one shape for another's job, you ask where the replay button is, or where the second reader gets its copy.

Two properties that belong in every queue's design and rarely make it in:

  • Expiry. Messages have a lifespan. SQS keeps a message for four days by default and fourteen at most; RabbitMQ has per-message and per-queue TTLs; Kafka has retention. A backlog that takes five days to drain silently expires the work at the back of the line, and a dead-letter queue (Section 6) with a shorter retention than its source queue loses the very messages it was meant to preserve. Set both deliberately.
  • Cost and latency. The queue hop itself costs milliseconds and money. Hosted queues bill per request (SQS charges per million API calls, which is why batching sends and receives and using long polling matter), and a self-run broker bills in engineer-hours. Neither is a reason not to use a queue. Both are reasons to know the number.

Section 4 — At-least-once: the delivery guarantee you actually get

Exactly-once delivery is impossible. The idempotency post (#6) walked through the argument in its Section 4, and it fits in two sentences: the broker can't know the consumer finished unless the consumer says so, and the acknowledgment itself can be lost. So this section is about what you build instead of chasing the impossible, and about a distinction the vendors blur: delivery is impossible to make exactly-once; processing is not, as long as the record of "I did this" and the effect itself are committed together.

Queues make three promises about delivery:

  • At-most-once: the message is delivered zero or one times. Fast, simple, and it can lose your message: the broker hands it to a consumer, the consumer crashes before processing, and nobody ever knows. Use it for metrics and sampling, where a lost data point is noise.
  • At-least-once: the message is delivered one or more times. Nothing is lost, but duplicates happen. This is what every serious queue promises by default, because it's the only reliable promise available.
  • Exactly-once: delivered exactly one time. See above: not available at the boundary between the broker and your code. Any vendor claiming it means exactly-once within a scope: Kafka's transactions give exactly-once processing for pipelines that read from Kafka and write back to Kafka; SQS FIFO deduplicates sends that carry the same deduplication ID within a five-minute window. Both are real and useful. Neither reaches the code that sends the email.

Where the switch between at-most-once and at-least-once actually lives is worth knowing, because it's a configuration setting rather than a product choice. It's the ack. A consumer that acknowledges before doing the work (auto-ack in RabbitMQ, committing the Kafka offset as soon as the batch arrives) has chosen at-most-once: crash after the ack and the message is gone. A consumer that acknowledges after the work (manual ack, commit after processing) has chosen at-least-once: crash before the ack and the message comes back. Choose the second and build for duplicates.

Sequence diagram of at-least-once delivery: a lost ack causes redelivery, the consumer checks its dedup store and skips — effectively once

Six steps, and that is the entire contract. The consumer did the work, the ack got lost, the broker redelivered (correctly, by its contract), and the consumer's dedup check turned the duplicate into a no-op. The reliability lives at the edge, not in the pipe. This is the idempotency post's central result in work clothes: at-least-once delivery plus an idempotent consumer (one that can safely handle the same message twice) equals effectively-once processing.

The welcome-email worker with no dedup is the standard cautionary tale. It sends the email, then crashes before acknowledging; a deploy rolled mid-processing. The broker redelivers. It sends the email again, crashes again (same deploy, same bad luck). One new user gets the welcome email three times, and the third copy's "We're so glad you're here!" reads as mockery. The fix is a processed_messages(message_id PRIMARY KEY) table and a check before sending. That table has a name, the inbox (the receiving-side twin of Section 5's outbox), and it has one rule that the naive version breaks: the check and the business write have to happen in one database transaction, so that "insert the message ID, then send, then crash before commit" rolls back cleanly and "insert the ID and commit, then crash before sending" can't happen. Check-then-act with a gap between them is the race the idempotency post's Section 5 is about. Duplicates are rare until the day they're a flood, and the inbox is flood insurance.

The mechanics you'll configure: the visibility timeout is SQS's term (Google Pub/Sub calls it the ack deadline) for how long a message stays hidden from other consumers after one takes it, 30 seconds by default, up to 12 hours. Finish and ack within the window, and the message is deleted. Crash or run long, and the message becomes visible again, redelivered by design. Set the timeout longer than your slowest legitimate processing time, or extend it with a heartbeat call while the work is running, or slow messages get processed twice while the first attempt is still running. RabbitMQ works differently: an unacknowledged message is redelivered when the consumer's connection or channel closes, and a consumer that holds a delivery unacknowledged past the acknowledgement timeout (30 minutes by default) has its channel closed by the broker, which requeues everything on it. Same outcome, different trigger.

One more interaction, since it bites the ordered queues: if a message in a FIFO message group or a Kafka partition keeps failing and being redelivered, everything behind it in that group waits. Ordering plus redelivery means a stuck message stalls its whole key, which is Section 6's poison-message problem with a narrower blast radius.

So the default for everything that matters: at-least-once, acked after the work, with an idempotent consumer. At-most-once is for telemetry you can afford to lose. And when a vendor's "exactly-once" marketing reaches your inbox, translate it: exactly-once inside their boundary, at-least-once delivery to your code, and the dedup still yours.


Section 5 — The dual-write problem and the transactional outbox

Your handler does two things: writes the order to the database, and publishes "order placed" to the queue. Two writes, two systems, no shared transaction. Now crash between them. Two ways to lose:

  • The lost event: the database commit lands, the process dies before publishing. The order exists; the warehouse never hears about it. Nobody ships anything, and nobody knows.
  • The phantom event: the publish lands, the database transaction rolls back. The warehouse ships an order that doesn't exist. Someone gets a free television.
The transactional outbox: the order row and its event land in one database commit, then a relay publishes the event to the queue

That diagram is the fix: the transactional outbox. Instead of publishing to the queue directly, the handler inserts the event into an outbox table in the same database transaction as the order itself. The transaction is atomic (both rows commit or neither does), so the lost event and the phantom event are both impossible by construction. A separate relay process reads new outbox rows and publishes them to the queue, then marks them sent (or deletes them). If the relay crashes mid-publish, it retries on restart, and the event might publish twice, which is exactly what Section 4's idempotent consumer is for. The outbox doesn't eliminate duplicates. It eliminates the two failure modes duplicates are preferable to.

The naive dual-write runs fine for a year. Then a deploy rolls during peak and kills a dozen handlers between the two writes. Twelve orders exist that the warehouse never saw, and customer support finds them three days later, by hand, in a spreadsheet. The postmortem's fix is the outbox table and a fifty-line relay, and the engineer's comment in the code review is the one to remember: "We had a distributed transaction and didn't know it. Now we don't."

The relay is a small program with its own failure modes, and they're worth listing because they're where outbox implementations go wrong in their second year:

  • It has to be a singleton, or coordinated. Two relays reading the same outbox publish everything twice (tolerable) and, worse, out of order (not tolerable if order matters). Either elect one leader (Section 7 has the mechanics) or have several relays claim rows with SELECT … FOR UPDATE SKIP LOCKED so each row has exactly one publisher.
  • Its poll interval is your latency. A relay that polls every second adds up to a second before the event exists. That's fine for emails and not for fraud checks; know which you have.
  • The outbox grows. Sent rows have to be deleted or archived, or the table becomes the largest one in the database and the relay's scan slows down.
  • It can fall behind, silently. Relay lag (the age of the oldest unsent row) is its own metric with its own alert. An outbox with ten thousand unsent rows is a system that thinks it's publishing and isn't.

The industrial version, for when you outgrow the polling relay: change data capture (CDC), a tool like Debezium reading the database's transaction log (Postgres's write-ahead log, MySQL's binlog) and publishing row changes as events. Same guarantee (the log is the transaction's truth), no polling, more infrastructure. Debezium's outbox recipe even deletes the outbox row in the same transaction that inserts it, since CDC sees the insert regardless. Start with the outbox table; graduate to CDC when the relay becomes a bottleneck you can measure.

Two patterns sit next to the outbox and get confused with it. The first is the reason you need it at all: there's no practical way to run one transaction across your database and your broker (the two-phase commit protocols that try, XA and friends, are slow, block on failures, and most brokers don't support them), so the outbox gives you the same effect by putting both writes in the database and letting the relay carry one of them out. The second is the saga: when a business operation spans several services (charge, reserve, ship), each step is a local transaction and each has a compensating action that undoes it (refund, release, cancel), so a failure at step three runs the compensations for steps two and one. Sagas are Section 9's orchestration problem, and the sharding post (#3, Section 7) and the idempotency post (#6, Section 6) both cover why the compensations must be idempotent. And a footnote for the pattern people ask about: event sourcing stores the events as the source of truth and derives the current state by replaying them, often paired with CQRS (separate read and write models). It's powerful for audit-heavy domains and a poor fit for ordinary CRUD applications, where it multiplies complexity for no one's benefit.

If the event is best-effort (analytics, metrics), the dual-write is fine and the outbox is ceremony. If the event drives money or fulfillment, the outbox is the difference between "the warehouse was told" and "we hope the warehouse was told."


Section 6 — Workers: parallelism, ordering, and poison

Workers are the unglamorous half of the system, and they're where most async outages actually live. The happy path is simple: add more workers, chew through messages faster. Parallelism is nearly linear until the work itself bottlenecks (the database the workers write to, the third-party API they call) or, on a log, until you run out of partitions. The queue scales easily; the things the workers touch do not. Size the worker pool against the downstream, not the queue.

Then the tension that shapes every worker design: ordering versus throughput. One queue with ten workers means ten messages processed concurrently: fast, and in no particular order. For sending emails, that's perfect. But some work has to happen in sequence: you can't process "order cancelled" before "order placed" for the same order, and you can't apply "set address to B" before "set address to A" and end up correct. The standard answer is partitioning by key: all messages for the same order go to the same partition (Kafka) or message group (SQS FIFO), each partition is consumed by one worker at a time, and order is preserved within the key while different keys still process in parallel. (This is the sharding post's idea of splitting by a key, applied to a stream instead of a table.) You don't choose between ordering and throughput globally. You choose per key.

Worker retry loop: process cleanly and ack, crash and retry up to N times, then park poison messages in the dead-letter queue for a human

The right side of that diagram is the section's real subject: poison messages. Some message (malformed, unexpected, a 2 GB "image" that isn't an image) crashes the worker every single time. The worker picks it up, dies, the message becomes visible again, another worker picks it up, dies. In an ordered queue, everything behind it in that partition or group stops; that's head-of-line blocking, the async equivalent of one broken-down car stopping the highway. In an unordered queue the damage is subtler: each redelivery costs one worker one crash-and-restart cycle, so throughput drops by however many workers are busy dying at any moment, and if the crash is an out-of-memory kill (OOM: the operating system terminates a process that has used more memory than it's allowed), the restart takes long enough that a pool can spend most of its time restarting.

The thumbnail service version: uploads processed happily for months, until a user uploads a "photo" that's actually a 2 GB corrupted TIFF. The thumbnail worker loads it into memory, gets OOM-killed, restarts, picks up the same message, gets OOM-killed again. Every worker in the pool takes its turn. Thumbnails stop company-wide, for all users, because of one file and a crash loop. The fix has two parts, and you want both: a retry cap (try three to five times, then stop; a message that fails five times isn't unlucky, it's poison), and a dead-letter queue, the DLQ, where exhausted messages go to wait. In SQS that's a redrive policy with a maxReceiveCount; in RabbitMQ it's a dead-letter exchange with a delivery limit on quorum queues. The DLQ is the morgue with a purpose: the poison is quarantined instead of blocking the living, a human can inspect it, and after a fix the messages can be replayed.

Every queue that carries work you can't afford to lose gets a DLQ. A queue without one is a system where one bad message is a company-wide outage.

The two exceptions to "every queue": an at-most-once telemetry queue where dropping is the design, and a FIFO queue where the vendor warns (SQS does) that dead-lettering a message breaks the ordering guarantee for its group, so you decide deliberately between "stall the key" and "skip the message." And a rule for the DLQ itself: give it a longer retention than its source, because a standard SQS queue counts a message's age from its original send, and a DLQ with the same four-day retention can expire a message the moment it arrives (FIFO queues reset the clock on the move).

Several more worker mechanics that decide whether a system is boring or exciting:

Prefetch. A worker shouldn't grab a hundred messages at once, because a crashed worker returns all hundred for redelivery, and because a hundred messages in one worker's hands is a hundred messages no other worker can process. The knob has different names (RabbitMQ's basic.qos prefetch count, Kafka's max.poll.records; SQS caps a single receive at ten). Keep it near the worker's actual concurrency.

Batching. The opposite pressure: sending, receiving, and acking messages in batches is often five to ten times cheaper and faster than one at a time, and producer clients can wait a few milliseconds to fill a batch before sending (Kafka's linger.ms). The trade is latency for throughput, and the sharp edge is partial failure: if message 7 of a batch of 10 fails, your code has to know which nine to ack.

Lag, not depth, as the metric. How many messages are waiting is less useful than how old the oldest one is. Kafka's consumer lag (the offset gap divided by the consume rate, which turns it into seconds) and SQS's ApproximateAgeOfOldestMessage are the numbers to graph and alert on, and when you autoscale workers, scale on backlog per worker or on message age. Scaling on raw depth oscillates: depth spikes, you add workers, depth crashes, you remove them, depth spikes again.

Priorities. "Send the password-reset email before the newsletter" wants a priority mechanism, and the simplest one is separate queues per priority with workers that drain the urgent one first (SQS has no priorities; RabbitMQ's classic queues take an x-max-priority argument and its quorum queues have built-in levels). Strict priority has a failure mode: under sustained load the low-priority lane starves entirely. Either weight the draining or give the low lane its own small worker pool.

Big payloads. Brokers cap message size (SQS at 1 MiB as of 2025, Kafka around 1 MB by default), and even under the cap, a queue full of 900 KB messages is slow and expensive. The claim-check pattern stores the payload in object storage and sends only a reference, and it's what the 2 GB TIFF should have been from the start: the worker downloads what it needs, and a corrupted file fails one download instead of one process.

Noisy tenants. One customer's bulk import of ten million records shouldn't delay everyone else's password resets. Options, from cheapest to most robust: a per-tenant rate limit on the producer (post #8), a queue per tenant or per tenant tier, or shuffle-sharding tenants across a set of queues so no two large tenants share all the same ones. SQS's fair queues are the hosted version of the same idea.

Visibility timeout versus long processing (Section 4): if legitimate work can take ten minutes, the timeout must exceed ten minutes, or the worker extends it with a heartbeat while it works, or the message gets processed twice concurrently and your idempotent consumer earns its keep.

Partition by key whenever order matters per entity (orders, users, accounts) and accept unordered everywhere else, because ordering costs throughput. Set retry caps low, DLQs on everything that matters, prefetch small, batches deliberate. And monitor the DLQ like a pager: a growing DLQ is either poison or a bug, and both want a human.


Section 7 — Later, precisely: delays, schedules, and retries with backoff

Not all "later" is "as soon as a worker is free." Sometimes later means precisely later:

Delayed jobs. "Retry this webhook in ten minutes." "Send the trial-ending reminder in three days." The message sits invisibly until its time comes. The delay is part of the message, not a worker's sleep call, because a sleeping worker is a worker that can't do anything else, and a sleep that's mid-way when the deploy rolls is a lost job. What the brokers actually offer is more limited than people assume: SQS delay queues and message timers max out at fifteen minutes, and FIFO queues don't support per-message timers at all; RabbitMQ's delayed-message exchange plugin stores its schedule on a single node (a node failure loses the delayed messages), isn't built for hundreds of thousands of pending messages or multi-day delays, and, having always been a separately shipped community plugin, was deprecated and archived in 2026, with its replacement available only in the commercial edition. For anything beyond minutes, use something designed for it: a scheduler service (EventBridge Scheduler on AWS), a scheduled_jobs table with a run_at column that the relay polls, or a workflow engine's durable timers (Section 9).

Scheduled jobs. "Every night at 2 AM, generate the report." "Every five minutes, check for stuck orders." This is cron's territory, but cron on one box is a single point of failure, and cron on ten boxes runs the job ten times. Production systems use a distributed scheduler with exactly one active runner, and "exactly one" is a leader election problem: the runners compete for a lock and only the holder runs the schedule. The lock can be a database advisory lock, a row with an expiry, a lease in etcd or ZooKeeper, or a Redis key with a TTL, and every one of them has the same sharp edge: the holder can pause (a long garbage collection, a stalled VM) past its lease's expiry, a second runner takes over, and now two runners believe they're the leader. The idempotency post (#6) has the fencing-token argument for surviving that. The practical rule that follows: a scheduled job will occasionally fire twice, so the job itself must be idempotent, and the job should usually just enqueue the real work rather than doing it. The schedule is the trigger; the queue is the muscle.

Exponential backoff with jitter: attempts 1-3 fail with growing waits, attempt 4 succeeds, so a recovering service isn't stampeded

That's the retry discipline in one picture, and it's the twin of the resilience post's backoff chapter: exponential backoff with jitter, on every retry, no exceptions. The version that works best, from AWS's analysis, is full jitter: each wait is a random duration between zero and min(cap, base × 2^attempt), so the ceiling doubles each time but the actual wait is spread across the whole range. The jitter isn't a sprinkle on top of a fixed delay; it's the whole delay, because the point is that a thousand workers never retry in the same instant. A payment service that skips this goes down for four minutes, every client retries immediately in a tight loop, and when the service comes back it's greeted by a stampede many times larger than normal traffic. It falls over again. Then again. The outage lasts forty minutes, of which thirty-six are self-inflicted. With backoff and jitter, the retries trickle in over minutes and the service recovers on the first try. The retry that kills a recovering service is the immediate one.

Cap the attempts (three to five is the usual range) and then the message goes to the DLQ (Section 6), not into infinite retry. Infinite retry is a slow-motion outage: a broken downstream means every message retries forever, workers churn, and the queue never drains. A retry policy without a cap is a hope, not a policy. Give retries a budget as well as a cap: no more than some fraction of total traffic may be retries at any moment, so a sick downstream isn't buried under retries of its own failures. And remember that retries stack across layers: a consumer that makes three attempts, calling a service whose client retries three times, calling a database whose driver retries three times, turns one message into up to 27 attempts. The resilience post has the arithmetic. Make retries visible: log the attempt number, alert on retry-rate spikes, because a sudden jump in retries is often the first signal that a downstream is sick, minutes before its own alarms fire.

The pattern is the same whether it's a queue consumer, an HTTP client, or a saga step: wait longer each time, randomize, stop eventually.


Section 8 — Backpressure: when producers outrun consumers

Every queue is a bet: that on average, consumption keeps up with production. Averages lie. The marketing campaign multiplies signups by ten for a weekend. A downstream slows every worker by 30%. A deploy halves the fleet for an hour. And then the arithmetic is merciless: 1,100 messages arriving per second, 1,000 processed per second, the queue grows by 100 a second, 360,000 an hour. The queue is doing its job, holding messages, right up until the holding becomes the emergency.

Backpressure when producers outrun consumers: add consumers to drain the pile, drop low-value work deliberately, or slow the producers

First, the word, because it gets used for everything. Backpressure is a signal that flows against the direction of data: the consumer telling the producer "slow down," and the producer actually slowing. TCP does it with its receive window; reactive-streams libraries do it with request(n); Kafka's producer does it by blocking when its local buffer fills. A queue with an unbounded buffer has no backpressure at all; it just absorbs until it can't. So the three options below are what you do when the buffer is filling, and only one of them is backpressure in the strict sense. You choose among them in a design review, not during the incident:

  1. Buffer. Add consumers, drain faster. Works when the surge is temporary and the downstream can take it; autoscaling workers on lag (Section 6) is the automated version. But buffering has a limit: if the arrival rate permanently exceeds capacity, you're renting a bigger pile.
  2. Drop. Deliberately discard low-value work. Analytics events can be sampled; the tenth "user viewed page" event this minute can go. Dropping is a business decision disguised as an engineering one. It must be chosen per queue, in advance, with the product's blessing, and every drop must be counted, because "we dropped 40% of analytics on Tuesday" is a fact someone needs.
  3. Push back. Slow the producers: bounded buffers that block, or, at the API edge, HTTP 429 with a Retry-After header (the rate limiting post's, #8, territory). This is the truthful option: the system admits it's full instead of accepting work it can't do. For user-facing producers, pushing back degrades gracefully. For internal ones, it moves the queue's problem to the caller, which is correct, because the caller can decide what matters.

The way this goes wrong slowly: a notification service that's 2% underwater, producers adding 2% more messages a day than the workers drain, a gap invisible on any single day's dashboard. On its own that gap would grow the backlog by about half an hour a day, three or four hours by the end of the week, which someone might catch. What actually happens is that it compounds. As lag grows, the downstream email provider starts timing out under the steadier load, timeouts become retries, retries widen the gap, and by Friday the "your order shipped" email arrives at midnight and the "flash sale ends tonight" email arrives Saturday. The postmortem finds the underlying gap has existed for months. The fix is an alert on the growth rate of lag, not just on lag: a flat-but-high queue is a spike; a steadily growing one is a capacity emergency with a date.

Lag is the vital sign. Not worker CPU, not message rate: the age of the oldest message, or depth divided by consume rate, which is the same number in seconds (360,000 messages at 1,000 a second is six minutes of lag, the answer to "how stale is the newest work?"). Alert on lag crossing the work's freshness budget and on lag growing over time. And size consumer capacity for a lag target rather than for the average arrival rate: decide how stale the work may get (six minutes is fine for analytics, fatal for fraud checks), then provision enough consumers that the peak-hour backlog drains inside that budget, with the queue absorbing the difference.

Document the policy per queue before the spike: buffer, drop, or push back, and for drop, what may be dropped and who approved it. Set the lag alerts on day one. The queue's promise is that it turns capacity emergencies into latency bills. Just remember: an unbounded latency bill is an emergency with slower paperwork.


Section 9 — Going deep: the principal-level toolkit

Everything so far was machinery. This section is judgment: the calls that separate a system that works from one that survives its own success. How to coordinate multi-step workflows, how to read the lag arithmetic, which vendor claims to believe, the traps in event-driven design, what async does to the user and to the person debugging it, and when not to go async at all.

Choreography versus orchestration. A multi-step workflow (order placed, payment charged, inventory reserved, warehouse notified, email sent) can be coordinated two ways. In choreography, there's no boss: each service listens for events and reacts. The order service emits order.placed; the payment service hears it and charges; it emits payment.charged; the warehouse hears that and ships. Nobody owns the flow. In orchestration, a central orchestrator runs the show: it calls step 1, waits, calls step 2, and when step 3 fails it runs the compensating actions for steps 2 and 1. That's the saga pattern from Section 5, and workflow engines (Temporal and its ancestor Cadence, AWS Step Functions, the Netflix-born Conductor OSS) are orchestration as a product: they give you durable timers ("wait three days, then continue"), automatic retries per step, and a complete history that answers "what state is order 4821 in, and what happened at step three?" What they cost is a stateful service you have to run or pay for, and a programming model your team has to learn.

Choreography vs orchestration: services reacting independently to events on a bus, versus one orchestrator driving charge, reserve, and ship

Choreography is beautifully decoupled (add a new reaction, the fraud service starts listening, without touching anything) and beautifully illegible: the workflow exists nowhere. It's spread across five codebases' event handlers. Debugging is archaeology ("who emitted what, when, and who heard it?"), and onboarding a new engineer onto a choreographed flow takes weeks. Orchestration inverts the trade: the flow is written down in one place, failures and compensations are explicit, but the orchestrator is a component to build, scale, and keep highly available, and every step's latency now includes a round trip to the boss. The rule of thumb: choreography for simple flows (three steps or fewer, no rollback needed); orchestration when the flow has branches, compensations, timers, or anyone will ever ask what state a particular order is in. Whichever you pick, the idempotency post's rule follows you: every step and every compensation must be idempotent, because steps retry.

The lag arithmetic, and the law behind it. Section 8 gave you the vital sign; here's why it works. If messages arrive at λ per second and each spends W seconds in the system (waiting plus processing), then on average there are L = λ × W messages in the system. That's Little's law, and it's the estimation post's (#13) favorite equation. Rearranged, W = L ÷ λ: lag equals depth divided by throughput, which is the formula from Section 8. And the growth rule is simpler still, plain conservation: the pile grows by (arrivals − departures) per second, so a queue that's growing has an arrival rate above its consume rate, and no amount of tuning changes that until departures beat arrivals. 360,000 messages at 1,000 per second is 360 seconds of lag; the newest message waits six minutes. That's the number for the dashboard, because "depth: 360,000" means nothing to a human and "lag: 6 minutes" means everything.

Computing queue lag: 360,000 messages at 1,000 per second means 6 minutes of lag — scale consumers until departures beat arrivals

Exactly-once versus effectively-once, one last time. You'll keep meeting vendors who say "exactly-once." Kafka's transactions give you exactly-once processing for pipelines whose input and output are both Kafka topics, with an idempotent producer and a read_committed consumer. SQS FIFO deduplicates sends within a five-minute window. Both are real, scoped, useful claims, and neither changes your architecture, because the boundary between the broker and the code that sends the email or charges the card is still at-least-once, and your consumer still dedups (Section 4). Design for effectively-once at the edges regardless of what the pipe promises. The pipe's guarantee is a bonus, not a foundation.

Event-driven pitfalls: the three that bite. First, ordering: a log gives you order within a partition, and nothing across partitions. If your design needs "every consumer sees events in global order," redesign; global order at scale costs you all your throughput, and Section 6's per-key ordering is the version that works. Second, replay: the log keeps history, so a new consumer can rewind to day one, which is a superpower for backfills and a loaded gun in two ways. For privacy, "delete my data" against a log is either retention expiry (the data ages out), compaction with a tombstone (for keyed, compacted topics), or crypto-shredding (encrypt each user's events with a per-user key and destroy the key), and you need to know which one your topics support before you promise anyone a deletion. For reprocessing, replaying six months of events through a fixed consumer is a migration, and it re-runs every side effect unless you plan for it: a replay mode flag that suppresses emails and external calls, a shadow consumer that computes results without acting, idempotency keys that span the replay, and a DLQ redrive that doesn't resurrect deleted work. Third, schema evolution: the producer ships event v2 with a renamed field, and the v1 consumer, deployed last quarter and still running, chokes. Events are a public API. The tooling for that is a schema registry (Confluent's, or the AWS Glue one) holding Avro, Protobuf, or JSON Schema definitions, with a compatibility rule enforced on every new version: backward compatible (new consumers can read old events), forward compatible (old consumers can read new events), or full. The operational rule that follows is to pick a compatibility mode and deploy in its order (with the usual backward mode, consumers can read old events, so consumers deploy first; with forward, producers deploy first), and the design rule is additive changes only: new optional fields, never a renamed or repurposed one. The Knight Capital lesson from the idempotency post applies double here: a flag whose meaning changed cost them more than $440 million in forty-five minutes.

An append-only log kept 7 days feeding three consumer groups: live traffic, analytics replaying from day one, and reprocessing after a bug

What async does to the user. "Later" is a promise the interface has to keep visibly. The HTTP shape is 202 Accepted plus a status resource (GET /jobs/4821) the client can poll, or a webhook or push notification when the work completes. In the interface it's optimistic UI (show the order as placed, mark the email as "sending") with a real pending state rather than a fake completed one. And it collides with the consistency the user expects: "I just placed an order, where's my order?" now races the queue. The replication post (#5) calls this read-your-writes, and the answers are the same here: route the user's next read to the source that has their write, or show the pending state honestly, or write the user-visible record synchronously and defer only the invisible work. What you don't do is let the user refresh into an empty page and wonder.

Tracing across the queue. A request's work is now scattered across time and machines: the charge at 10:00, the email at 10:04, the warehouse event at 10:06, on three hosts. Without distributed tracing, "why didn't the email send?" is a multi-system scavenger hunt. The mechanism is simple and constantly forgotten: the producer writes the trace context (the W3C traceparent value) into the message's headers or attributes, and the consumer reads it back and starts its span as a child of the producer's. For batch consumers, OpenTelemetry's span links let one consumer span point at the many producer spans it's handling. Put the message ID in every log line the consumer writes. The observability post (#9) covers the rest of the machinery; this is the one habit that keeps traces from going cold at the queue.

Security and testing, briefly, because they're skipped. A queue is an attack surface: control who may publish and who may consume (IAM policies, Kafka ACLs), encrypt in transit and at rest, keep personal data out of messages where you can (the claim-check pattern helps here too), and treat message contents as untrusted input in the consumer, exactly as the security post (#12) says to treat any input. For testing, run the real broker locally (LocalStack, testcontainers), write contract tests against the schema registry, and inject the failures on purpose: duplicate every message in a test environment, deliver some out of order, plant a poison message, and check that the inbox, the ordering keys, and the DLQ behave. Async assertions need a wait-and-poll helper; a test that checks the side effect immediately after enqueueing is testing the race.

When not to go async. Async has costs, and a principal names them. The debugging cost above. The reasoning cost: "did it happen yet?" becomes a real user question, and support has to be able to answer it. The consistency cost, from the user-experience paragraph. So don't go async when the user needs the answer now (Section 2's question, still the law), when the flow is two steps and always fast (a queue for a 50 ms job is ceremony), or when the work's value expires in seconds (real-time bidding, live collaboration).

Async is a loan against future debugging. Take it when the interest is worth it.

The team that choreographs everything, twelve services reacting to events, no orchestrator, no tracing, has an elegant year. Then a payment succeeds, the warehouse never ships, and the investigation takes three engineers two days: the event was emitted, consumed, and dropped by a handler with a swallowed exception, and nothing anywhere recorded the flow's state. They introduce an orchestrator for the money path, not for all paths, just the one where "what state is order 4821 in?" is a question with dollar signs, and keep choreography for the rest. The principal move is not to pick a side but to draw the line between the flows that need a boss and the ones that don't.

And the gate, stated the way the idempotency post stated its own: a new queue ships if and only if it has a DLQ (or a documented reason not to), a lag alert, a documented drop-or-push-back policy, an idempotent consumer, and trace context in its messages. "We'll add monitoring later" is not a story. Five minutes in the design review, or a page in the middle of the night: those are the options, and they're priced accordingly.


Do It Later, On Purpose, distilled

For the skimmers and the revisitors: everything above, on one page.

Key numbers and rules:

The one question Does the user need the result now? Yes goes inline, no goes on the queue
The seven-second checkout, fixed 7.2 s → 2.0 s by queueing four of six steps; nothing got faster
Three queue shapes Point-to-point for tasks; pub/sub for fan-out; the log for facts with history. Defaults, not laws
Where at-least-once is chosen Ack after the work (not before); commit the offset after processing
The inbox processed_messages checked and written in the same transaction as the effect
Visibility timeout (SQS) 30 s default, 12 h max; longer than the slowest legitimate job, or heartbeat it
Broker delay limits SQS timers ≤ 15 minutes, none on FIFO; multi-day delays need a scheduler, a run_at table, or a workflow timer
Message size SQS 1 MiB, Kafka ~1 MB default; beyond that, claim-check
Retention SQS 4 days default / 14 max; Kafka 7 days default; give the DLQ longer retention than its source
Lag Age of the oldest message, or depth ÷ consume rate; 360,000 at 1,000/s = 6 minutes
Little's law L = λ × W; the growth rule is arrivals − departures
Retry cap before DLQ 3–5 attempts; then the dead-letter queue
Backoff Full jitter: wait random(0, min(cap, base × 2^attempt)); retries budgeted; layered retries multiply
Scheduled jobs One leader-elected runner; assume it fires twice; the job enqueues, the queue works
Schema rule Registry-enforced compatibility; additive changes only; deploy in the compatibility mode's order (consumers first under backward compatibility)
The code-review gate DLQ + lag alert + drop-or-push-back policy + idempotent consumer + trace context, or it doesn't ship

Every trade-off, in one table:

Decision Chose Over Why
Slow endpoint Ask "does the user need it now?" Optimize each step Most slowness is needed work spending the user's time, not slow work
Latency budgets Two budgets: ms for sync, minutes for async One budget for everything Async work draws from the cheap account
Spike handling Queue as shock absorber Scale workers for the peak The queue turns a capacity emergency into a latency bill
Downstream outage Queue as bulkhead Fail the request Messages wait; the producer never knows there was an outage
Queue shape, tasks Point-to-point (SQS, RabbitMQ queues) Log-based One worker pool, each message done once and forgotten
Queue shape, fan-out Pub/sub (SNS → SQS, exchanges, topics) Duplicating producers Many readers, each with its own copy, no history
Queue shape, events Log-based (Kafka, streams) Point-to-point Many independent readers; replay is a feature; parallelism capped by partitions
Delivery guarantee At-least-once + idempotent consumer (inbox) "Exactly-once" pipe Exactly-once delivery is impossible; effectively-once lives at the edge
DB write + event Transactional outbox Dual-write in code Same transaction or neither; duplicates beat lost and phantom events
The relay Singleton or SKIP LOCKED; lag alerted; rows cleaned up "It's just a loop" Two relays reorder; a stalled relay is silent
CDC Graduate to log-based capture Polling relay forever Same guarantee, less polling, when the relay is a measured bottleneck
Ordering needs Partition by key Global ordering Per-key order preserves throughput; global order kills it
Poison messages Retry cap + dead-letter queue Infinite retry One bad message must not stall the highway; the DLQ quarantines it
Worker scaling On lag or backlog per worker On raw depth Depth-based scaling oscillates
Big payloads Claim-check (store the blob, send a reference) Fat messages Size caps, cost, and a corrupted file fails a download instead of a process
Noisy tenants Per-tenant limits, queues, or shuffle-sharding One shared queue One customer's bulk import shouldn't delay everyone's password resets
"Later" that's precise Scheduler service / run_at table / workflow timers; leader-elected cron Worker sleep, bare cron, broker timers past their limits Sleeps die on deploy; single-box cron is a single point of failure; SQS timers stop at 15 minutes
Retries Full-jitter backoff, capped, budgeted Immediate retry Immediate retries stampede recovering services
Producers outrun consumers Buffer, drop, or push back, chosen in advance Decide during the incident Buffer for spikes, drop low-value work (counted), push back openly
Multi-step workflows Choreography ≤ 3 steps; orchestration (or an engine) beyond One pattern for everything Simple flows stay decoupled; money paths get a boss and a history
Schema changes Registry, compatibility rules, additive, deployed in the mode's order Renaming or repurposing fields Events are a public API
Replay Replay mode, shadow consumers, spanning idempotency keys Re-run and hope Replay re-runs side effects unless you stop it
User experience 202 + status resource, pending states, read-your-writes routing Empty page after refresh "Later" is a promise the interface keeps visibly
Going async at all When the user doesn't need it now Async everything Debugging, reasoning, and consistency costs are real; async is a loan

Three ideas to take with you:

  1. The only question is "does the user need the result now?" Everything else (the queue, the workers, the DLQ, the backoff) is machinery in service of that one decision. Most slow endpoints aren't slow because any step is slow; they're slow because work nobody's waiting for is spending the waiter's time.

  2. At-least-once delivery plus an idempotent consumer equals effectively-once, and the dedup is always at the edge. The pipe's guarantees are bonuses, not foundations. If you can't point to the inbox check in your consumer, and to the transaction it shares with the effect, you don't have reliability, you have optimism.

  3. A queue without a DLQ, a lag alert, a drop-or-push-back policy, an idempotent consumer, and trace context is a hope, not a design. The boring parts are the design. Set them in the review, when it's cheap.


Further reading


Where you'll meet this

This post leans on the idempotency post (#6) twice: its Section 4 walked through why exactly-once delivery is impossible, which is why every consumer here dedups, and its inbox and outbox sections are this post's Sections 4 and 5, the same pattern from both sides. It's the resilience post's (#2) bulkhead made concrete: the queue is the wall between "the email service is down" and "checkout is down," and Section 7's backoff is that post's retry chapter with the queue as the stage. The sharding post's (#3) sagas reappeared in Section 9 as orchestration, and its partition-by-key idea as Section 6's ordering. The rate limiting post (#8) is where "push back" lives once it reaches the API edge, the replication post (#5) owns the read-your-writes problem async creates, and the observability post (#9) is where the trace context in Section 9 ends up. And the URL shortener design was async from the start: the interview's answer to "how do click analytics avoid slowing the redirect?" was this post's one question. The user needs the redirect now; the analytics can wait. Next up is load balancing (#11), the box in every diagram that decides which worker gets the message in the first place.


Keep exploring: the Core Concepts series

Every post in the series stands alone. Read them in any order.


Blueprints of Scale — Core Concepts #10. If this helped, the best thanks is a share with someone who's learning.

Core Concepts

Part 8 of 14

Six ideas that show up in almost every system you'll ever design: caching, resilience, sharding, unique IDs, replication, and idempotency. Each post stands alone — read them in any order.

Up next

Do the Math First: Estimation for System Design, Explained Like You're New

Blueprints of Scale — Core Concepts #13