Message queues: what they actually guarantee

JR

Jai Rao

August 24, 202610 min read

Why exactly-once delivery is something you build rather than buy, why checking before inserting is not idempotency, and the dual-write bug that passes review.


A queue does not make work disappear. It moves the work somewhere the user is not waiting, and in exchange hands you a set of problems that synchronous code never had. That trade is usually worth making. It is not free, and most of the difficulty shows up in places the architecture diagram does not mark.

Start with the situation that motivates it. A signup handler creates an account and sends a welcome email inline. The user's request now depends on an external mail provider: when that provider is slow, signup is slow, and when it is down, signup fails. A registration has been made contingent on something that has nothing to do with registering.

What moving it to a queue actually changes

Publish an event and return. Signup no longer waits on the mail provider, no longer fails when it is down, and the user gets their response in milliseconds.

Now the honest part. You have told the user "done" for work that has not happened. If the email never sends, nothing in that request will ever notice. The failure has moved out of the request path, which is exactly what you wanted, and into a background system that has to be operated — which is the cost nobody budgets for. A queue converts a visible synchronous failure into an invisible asynchronous one, and invisible failures need monitoring to become visible again.

The reasons this trade is usually right:

  • Absorbing bursts. Ten thousand simultaneous signups become a queue that drains at whatever rate your consumers manage, instead of ten thousand concurrent SMTP connections.
  • Isolating unreliable dependencies. A slow third party stops being able to slow your own endpoints.
  • Fanning out. One event, several independent consumers — email, analytics, provisioning — none of which the publisher needs to know about.
  • Genuinely long work. Video transcoding or report generation cannot happen inside a request no matter how fast your code is.

And when it is the wrong choice: if the user needs the result to continue, a queue just adds a polling problem. Asynchronous means the answer arrives later, so there must be somewhere for it to arrive.

Exactly-once delivery does not exist

Three delivery semantics get discussed, and the third is widely misunderstood.

At-most-once: acknowledge before processing. A crash loses the message. Acceptable for a metrics sample, unacceptable for a payment.

At-least-once: acknowledge after processing. A crash before the acknowledgement means redelivery. This is what almost every broker gives you.

Exactly-once is not something a broker can hand you end to end, and the reason is worth understanding rather than accepting. Suppose your consumer processes a message successfully and then dies before its acknowledgement reaches the broker. From the broker's side, a completed message and a failed one look identical — the acknowledgement is absent in both cases. It cannot distinguish them, so it must choose: redeliver, and risk a duplicate, or drop, and risk a loss. There is no third option. This is not an engineering shortfall to be fixed by a better broker; it is a property of message-passing over an unreliable channel.

So exactly-once is something you build, not something you buy: at-least-once delivery plus idempotent consumers. Duplicates arrive, and processing one twice has the same effect as processing it once.

Idempotency, and why checking first is not enough

Idempotency needs a key that identifies the unit of work — a message id, or better, something derived from the business operation such as an order id plus an action. Record keys you have completed, and skip ones you have seen.

The obvious implementation is wrong:

Text
# BROKEN: the check and the write are separate operations.def handle(key, payload):    if store.exists(key):        return "already done"    charge_card(payload)              # a second worker can be here too    store.insert(key)

Two workers receive the same redelivered message. Both check, both find nothing, both charge the card. Simulating exactly that interleaving:

Text
  check-then-insert, two workers interleaved: side effects performed = 2  (both read 'absent' before either wrote -- the charge happens twice)

The check proves nothing, because it describes the past rather than reserving the future. What enforces uniqueness is a uniqueness constraint. Attempt the insert and let the database arbitrate:

Text
def handle(key, payload):    try:        db.execute("insert into processed (idem_key) values (?)", (key,))    except IntegrityError:        return "duplicate ignored"     # another worker owns this one    charge_card(payload)    db.commit()
Text
delivering the same message twice:  1st: processed  2nd: duplicate ignored  rows for msg-1: 1

One row, one charge, no coordination needed between workers. Reserving before the side effect leaves a window where a crash means the key is claimed but the work is unfinished, so record a state — claimed, then completed — and let a stale claim expire back for retry. That is more code, and it is the version that is actually correct.

Where you can, prefer operations that are naturally idempotent. "Set status to shipped" can run any number of times; "increment the shipped count" cannot. Choosing the first formulation removes the problem instead of managing it.

Ordering: buy less of it than you think

Global ordering across a queue is expensive, because it means one consumer at a time — total order and parallelism are directly opposed. Systems built for strict global ordering do not scale out, by construction.

The good news is that applications almost never need it. What they need is that events about the same entity arrive in order: this account's balance changes, this document's edits. Two unrelated accounts have no meaningful ordering between them.

Partitioning by entity key delivers exactly that. Hash the key, route to a partition, one consumer per partition. You get per-entity ordering with as much parallelism as you have partitions — which is why brokers ask you for a partition key, and why choosing the entity id rather than something random is what makes ordering work.

One consequence to plan for: a partition is only as fast as its slowest message. One poisonous message at the head of a partition blocks every subsequent message for that entity, which is a stall rather than a failure and looks healthy from the outside.

Retries, backoff, and the queue nobody reads

Retrying immediately is worse than not retrying. A dependency failing under load, hit again instantly by every consumer, stays down. Back off exponentially, and add jitter so retries do not synchronise into a pulse:

Text
import randomdef backoff(attempt, base=0.5, cap=30.0):    raw = min(cap, base * (2 ** attempt))    return random.uniform(raw / 2, raw)      # jitter: spread the retries
Text
  attempt 0: window up to   0.5s   chosen  0.31s  attempt 3: window up to   4.0s   chosen  3.21s  attempt 6: window up to  30.0s   chosen 15.20s  attempt 7: window up to  30.0s   chosen 27.56s

The cap matters: unbounded doubling reaches delays measured in days. The jitter matters more than it looks — without it, a thousand consumers that failed together retry together, producing exactly the synchronised thundering herd the backoff was supposed to prevent.

Retry the right things. A timeout or a 503 is worth retrying; a malformed payload or a validation failure will fail identically every time, and retrying it wastes capacity and delays real work. Distinguish them explicitly, and send permanent failures straight to the dead-letter queue rather than through the retry ladder.

A dead-letter queue is where messages go after exhausting retries. It is operationally essential and routinely becomes a liability, because a dead-letter queue nobody looks at is a silent data-loss bucket with good intentions. Alert on messages arriving in it, not on its depth — depth tells you about accumulated history, arrival tells you something is wrong now. And keep enough context on each message to replay it after a fix, because that is the entire point of not discarding it.

The dual-write problem

Here is a bug that survives code review because both lines look correct:

Text
def place_order(order):    db.insert(order)               # 1. commit to the database    queue.publish(order_event)     # 2. tell everyone else

These are two separate systems and there is no transaction spanning them. Crash between the lines and the order exists with no event — no confirmation email, no warehouse notification, no analytics. Swap the order and a crash publishes an event for an order that does not exist, so consumers process something unfindable. There is no arrangement of two statements that is safe, which is why this pattern is worth recognising on sight.

The fix is to make the event part of the same transaction by writing it to a table in the same database:

Text
def place_order(order):    with db.transaction():                    # both, or neither        db.insert("orders", order)        db.insert("outbox", {            "id": new_uuid(),            "topic": "order.placed",            "payload": serialise(order),        })# A separate relay polls `outbox`, publishes, and marks rows sent.

Now the order and the intent to publish commit atomically. A relay reads unsent rows, publishes them, and marks them sent — and if it crashes after publishing but before marking, it republishes. Which is fine, because it delivers at-least-once into consumers you already made idempotent. The transactional outbox is the standard answer here, and it is worth reaching for as soon as a database write and a message must agree.

A queue and a log are different tools

Two models get called "message queue" and behave differently enough to matter.

In a work queue, a message is a task claimed by one consumer and removed when done. Many workers compete for messages; each unit of work happens once. This is the model for sending an email or transcoding a file.

In an append-only log, messages are written in sequence and retained, and each consumer group tracks its own position. Several independent consumers read the same messages without competing, and a new consumer can start from the beginning. This is the model for event streams that several systems care about — and it enables things a work queue cannot: replaying history into a new service, or reprocessing after fixing a bug, because the messages are still there.

Choose by whether the message is a task or a fact. Tasks belong in a queue and disappear when done. Facts belong in a log, because a fact does not stop being true once one system has read it.

Backpressure, and the queue that only grows

A queue is a buffer, not storage. A buffer absorbs a burst and then drains. If yours only ever grows, it is not buffering — it is deferring a failure and adding latency while it does so.

The number to watch is not depth, which is noisy and gives no sense of scale. It is consumer lag, ideally expressed in time: how old is the oldest unprocessed message? A depth of fifty thousand means nothing on its own; "the oldest message is forty minutes old" is immediately actionable and comparable across services. Alert on that.

When consumers cannot keep up you have four options, and it is worth knowing which you are choosing: add consumers (only helps if partitioning allows it), make each message cheaper to handle, shed load by dropping messages you can afford to lose, or slow the producers down. Doing none of them is also a choice — the one where the queue grows until it hits a retention limit and starts discarding your oldest messages silently.

Alongside lag, watch the ratio of processed to failed messages, the rate of arrivals into the dead-letter queue, and the age of the oldest unacknowledged message. Between them those tell you whether the system is keeping up, whether it is correct, and whether anything is stuck — which is everything you need to know about a queue you cannot see.