Message Queues and At-Least-Once Delivery
Why a broker can promise at-least-once or at-most-once but never exactly-once, queue versus log, and how idempotent consumers and the outbox make delivery effectively once.
A message queue decouples the producer of work from its consumer in time and in load: the producer writes a message and moves on, the consumer processes it later at its own pace, and the broker holds it in between. The catch is delivery semantics. A broker can promise at-most-once (may lose a message) or at-least-once (may deliver it twice), but never exactly-once on its own, because the acknowledgement that says "done" can be lost after the work was done. So every consumer must tolerate duplicates, and "exactly once" is a property you build end to end with idempotent handlers and a transactional outbox, not a checkbox on the broker.
Context
Queues are older than the web: IBM MQ shipped in 1993 to let mainframe programs talk without both being up at once, Java standardised the API as JMS in 2001, and the AMQP protocol (2003) produced RabbitMQ in 2007. Amazon SQS was the very first AWS service (2004). Then LinkedIn open-sourced Kafka in 2011 and changed the shape of the thing: not a queue that deletes messages once consumed, but a durable, replayable log that many independent consumers read at their own offsets. Redis Streams (2018), NATS JetStream and Google Pub/Sub sit somewhere between the two models. The "exactly-once" argument peaked in 2015-2017: engineers pointed out that no network protocol can guarantee it, and Kafka 0.11 shipped "exactly-once semantics" that, read carefully, apply to Kafka-to-Kafka pipelines only.
You have met queues as background jobs (Sidekiq, Celery, BullMQ), as the thing that sends the welcome email after sign-up, as webhook retries, and as an incident titled "customer received the same invoice three times". You have met their vocabulary as an SQS visibility timeout, a Kafka consumer lag graph, or a dead-letter queue nobody looked at for a month.
The simplest possible producer and consumer:
// producer: fire and forget, returns in a millisecond
await queue.send({type: 'order.paid', orderId: 'o_42', messageId: 'm_9f3a'})
// consumer: runs later, on another machine, possibly more than once
for await (const msg of queue.receive()) {
await sendInvoiceEmail(msg.body.orderId) // the side effect
await msg.ack() // "delete it, I'm done"
}- Broker
- The server that stores messages between producer and consumer: RabbitMQ, Kafka, SQS, Redis.
- Ack
- The consumer telling the broker a message is fully processed. Until then the broker considers it in flight.
- Redelivery / visibility timeout
- If no ack arrives within a window (SQS: visibility timeout, RabbitMQ: consumer disconnect), the broker gives the message to a consumer again.
- Dead-letter queue (DLQ)
- Where a message goes after failing N times, so one poison message does not block everything behind it.
- Offset
- In a log (Kafka), a consumer’s position: the index of the last message it has committed. Acking is moving the offset.
- Consumer group
- A set of consumer instances that share the work of one subscription; each partition is read by exactly one member at a time.
- Lag / backlog
- How far consumers are behind producers: messages waiting in a queue, or offsets behind the log head.
Why it matters
Once a system has more than one service, a queue is how it survives one of them being slow or down, and how it absorbs a traffic spike without dropping requests. The failure modes are correspondingly systemic: a handler that is not idempotent turns every retry into a duplicate charge, a poison message with no dead-letter path stops a whole queue, and a service that writes to its database and then publishes an event can lose the event forever if it crashes between the two.
Why exactly-once delivery is impossible
The consumer has two moments to acknowledge: after doing the work or before it. Each choice loses in one crash scenario, and no protocol trick removes the gap, because the ack is itself a message that can be lost.
| Guarantee | Broker promises | Consumer must | Use for |
|---|---|---|---|
| At-most-once | Deliver zero or one time; never retry | Nothing special; accept loss | Metrics samples, presence pings, anything where a fresh value replaces a lost one |
| At-least-once | Deliver until acked; may repeat | Be idempotent, or deduplicate by message id | Almost everything: orders, emails, payments, sync jobs |
| "Effectively once" | At-least-once underneath | Dedup key stored atomically with the business write; producer uses an outbox | Anything where a duplicate is visible to a user or costs money |
Queue versus log
RabbitMQ, SQS, Sidekiq and BullMQ are queues: a message is delivered to one of several competing consumers and deleted on ack. Kafka, Kinesis, Redis Streams and Pulsar are logs: messages are appended to partitions and kept for a retention period; each consumer group tracks its own offset and can replay. The choice shapes ordering, fan-out and operations.
| Queue (RabbitMQ, SQS, Sidekiq) | Log (Kafka, Kinesis, Redis Streams) | |
|---|---|---|
| After ack | Message deleted | Message stays until retention expires; offset moves |
| Fan-out | One consumer per message; multiple subscribers need exchanges / SNS in front | Every consumer group reads everything independently |
| Ordering | Per queue at best; competing consumers and retries reorder. FIFO queues limit throughput per group | Strict within a partition; choose the key to get per-entity order |
| Replay | Not possible once acked | Reset the offset and reprocess history |
| Parallelism | Add consumers freely; the broker balances per message | Bounded by partition count; one member per partition |
| Per-message retry / delay | Built in: requeue, delayed exchanges, visibility timeout, DLQ | Not built in: skip and write to a retry topic, or block the partition |
| Good fit | Job processing, task fan-out to workers, request buffering | Event streams many services consume, audit trails, change data capture, high volume |
Building effectively-once processing
Two places leak: the consumer may do the side effect twice, and the producer may do the business write without ever publishing. Both are fixed by making the queue-related write and the business write one atomic operation.
Consumer side: idempotent handler with a dedup record
async function handle(msg: Message) {
await db.transaction(async tx => {
// 1. claim the message id; a duplicate hits the unique index and stops
const claimed = await tx.execute(
`INSERT INTO processed_messages (id) VALUES ($1)
ON CONFLICT (id) DO NOTHING`, [msg.id],
)
if (claimed.rowCount === 0) return // already done → just ack
// 2. the business write, in the SAME transaction as the claim
await tx.execute(
`UPDATE orders SET status = 'invoiced' WHERE id = $1`, [msg.body.orderId],
)
})
// 3. side effects that live outside the DB need their own idempotency key
await email.send({idempotencyKey: msg.id, template: 'invoice', ...})
await msg.ack()
}
// crash anywhere → redelivery → step 1 short-circuits or the tx rolled back;
// either way the order is invoiced once and the email provider dedups on msg.idProducer side: the transactional outbox
"Save the order, then publish" is two writes to two systems, and a crash between them loses the event; publishing first and then saving risks an event for an order that does not exist. The outbox pattern writes the event into an outbox table in the same database transaction as the order, and a separate relay publishes rows from that table to the broker, marking them sent. The relay is at-least-once too, which is fine, because consumers deduplicate.
Operating a queue
| Problem | Mechanism | Detail |
|---|---|---|
| Transient failure | Retry with exponential backoff and jitter | Requeue with a delay that grows; never immediate, or a down dependency gets hammered by every message at once. |
| Poison message | Max attempts → dead-letter queue | After N failures move it aside with the error attached; alert on DLQ depth; provide a replay tool. |
| Slow handler | Visibility timeout ≥ processing time, or heartbeat | If the timeout is shorter than the work, a second consumer starts the same message while the first is still running. Extend the lease periodically for long jobs. |
| Ordering | Partition key / FIFO group per entity | Global order is unavailable at scale; per-entity order (all events of order 42) is cheap with a key. |
| Backlog growth | Monitor lag; scale consumers; shed load | Lag that grows for an hour is an outage in slow motion. Alert on age of oldest message, not just count. |
| Large payloads | Claim-check pattern | Put the blob in object storage and send its key; brokers cap message size (SQS 256 KB, Kafka default 1 MB). |
- 1Receive with a lease long enough for the p99 processing time, and a heartbeat that extends it for outliers.
- 2Check the dedup store first (or rely on the unique constraint inside the transaction) so a redelivery is cheap.
- 3Do the work. Any external call gets an idempotency key derived from the message id.
- 4Ack only after the work is durably committed. On failure, do not ack: let the backoff schedule redeliver, and after N attempts let the broker move it to the DLQ.
- 5Export lag, oldest-message age, retry rate and DLQ depth as metrics; those four graphs describe the health of the whole pipeline.
Pitfalls
- Acking before the work is done
Often accidental: a framework auto-acks on receive, or the handler returns a promise the runner does not await. A crash or a thrown error after that point loses the message silently. Turn off auto-ack, ack in a
finallyonly on success, and test by killing the consumer mid-handler. - A handler that is "usually" idempotent
Setting
status = paidis safe to repeat; incrementing a counter, appending a row or sending an email is not. Every handler will be called twice eventually. Give each a dedup key stored in the same transaction as its write, and pass that key to any external API that supports one. - Visibility timeout shorter than the job
The broker assumes the first consumer died and hands the message to a second one; now two workers run the same job concurrently and both may succeed. Size the timeout from measured p99 duration, extend it with a heartbeat for long work, and keep the dedup check so the loser exits early.
- No dead-letter queue
One message whose payload crashes the handler is retried forever, and in a FIFO or single-partition setup everything behind it waits. Cap attempts, route failures to a DLQ with the exception attached, alert on it, and have a way to fix and replay.
- Dual write: save to the database, then publish
Two systems, no shared transaction. A crash after the commit means downstream services never hear about the order; a broker outage means the same. Write the event to an outbox table in the same transaction and let a relay publish it, or use change data capture on the table itself.
- Assuming global order
Competing consumers, retries and multiple partitions all reorder messages. Design handlers to tolerate "shipped" arriving before "paid" (store and reconcile, or use version numbers), and use a partition key when one entity's events must stay in sequence.
Interview questions
Q1What is the difference between at-least-once and exactly-once delivery, and which one do real brokers give you?
At-least-once means the broker keeps redelivering until it receives an ack, so a message may be processed twice if the ack is lost; exactly-once would mean every message is processed precisely one time end to end. Brokers can only give at-least-once or at-most-once, because the acknowledgement can be lost after the work is done and the broker cannot tell that apart from the work never happening. Exactly-once is achieved by the application: idempotent handlers keyed by message id, so duplicates are harmless.
Q2The consumer sends the email and crashes before acking. What happens, and how do you design for it?
The broker's lease expires and it redelivers the message to another consumer, which sends the email again. I would give the email call an idempotency key derived from the message id so the provider drops the second send, and record the message id in a processed table inside the same database transaction as any state change, so the second handler short-circuits. Then the duplicate delivery is invisible to the user.
Q3Walk me through implementing an idempotent consumer.
Open a transaction, insert the message id into a table with a unique constraint using ON CONFLICT DO NOTHING; if zero rows were inserted, the message was already processed, so commit and ack. Otherwise perform the business writes in the same transaction and commit. External side effects get the message id as their idempotency key. Ack only after commit. A crash at any point either rolls back the claim, so redelivery retries cleanly, or leaves a committed claim that makes the redelivery a no-op.
Q4Kafka or RabbitMQ for this system?
RabbitMQ or SQS when the workload is jobs: each message goes to one worker, needs per-message retry and delay, and nobody needs to re-read history. Kafka when the workload is events that several services consume independently, when order per entity matters, when volume is high, or when replaying history is a requirement. Kafka's cost is operational and in the lack of per-message retry semantics; RabbitMQ's cost is no replay and weaker fan-out.
Q5How do you keep events for one order in order while processing millions of orders in parallel?
Partition by the order id: in Kafka the key hashes to a partition, and one consumer reads a partition sequentially, so all events for one order arrive in order while different orders spread across partitions and consumers. In SQS FIFO the same idea is a message group id. Global order across all orders is not available at scale and almost never needed.
Q6What is a poison message and what do you do about it?
A message the handler can never process successfully, because of a bad payload or a bug, so it fails, is redelivered, fails again, and in a strictly ordered setup blocks everything behind it. The fix is a maximum attempt count after which the broker moves it to a dead-letter queue with the error attached, an alert on that queue's depth, and a tool to inspect, fix and replay.
Q7You save an order and then publish an event. How can the event get lost, and how do you prevent it?
The process can crash, or the broker can be unreachable, between the database commit and the publish, so the order exists but no downstream service ever learns about it; reversing the order risks publishing events for orders that were rolled back. The transactional outbox fixes it: insert the event into an outbox table in the same transaction as the order, and have a relay, polling or change data capture, publish rows and mark them sent. The relay may publish twice, which consumers already tolerate.
- A queue decouples in time and load. The broker can promise at-least-once (duplicates) or at-most-once (loss); exactly-once delivery is impossible because the ack can be lost after the work.
- Ack after the work, and make the work idempotent: claim the message id in the same transaction as the business write, and pass it as the idempotency key to external calls.
- Queues delete on ack and balance per message; logs retain, partition and let every consumer group replay. Per-entity ordering comes from the partition key, never from global order.
- The transactional outbox closes the producer-side gap: event row and business row in one commit, a relay publishes.
- Operate with backoff plus jitter, a max-attempts dead-letter queue, a lease longer than p99 processing, and alerts on lag age and DLQ depth.
- Kafka's exactly-once covers Kafka-to-Kafka pipelines; any side effect outside it is at-least-once again.