How message brokers actually deliver
What actually happens between publish and ack, and why does every queue end up delivering duplicates?
Brokers hand out leases, not messages: a delivery stays owned by the broker until the consumer acks, and anything unacked is delivered again. Practical systems are at-least-once, so exactly-once is something you build from dedupe plus idempotent side effects.
THE MENTAL MODEL
Think of a delivery as a lease with a timer. The broker keeps the message, lends it to one consumer, and takes it back if the lease (visibility timeout, ack wait, ack deadline, session) ends without an ack. Ack after the side effect and a crash causes a duplicate (at-least-once); ack before it and a crash loses the work (at-most-once). Nothing in between exists at the network boundary, so "exactly-once" always means a narrower promise: no duplicates within a window, within a region, or within the broker's own read-process-write loop. Ordering is equally scoped: per partition, message group, or ordering key, never globally across parallel consumers.
HOW IT FITS TOGETHER
- Producer → Brokerpublish order.created (eventId=e42)
- Broker → Producerconfirm: durably storedWithout a confirm, a producer retry can itself create a duplicate.
- Broker → Consumerdeliver #1, lease starts (e.g. 30 s)Consumer charges the card, then crashes before acking.
- Broker → Consumerlease expired: redeliver #2redelivered=true / receive count 2. Handler must see e42 is done.
- Consumer → Brokerack (delete, XACK, commit offset)
- Broker → Consumerdeliver poison message, attempt N
- Consumer → Brokernack or no ack, every time
- Broker → Dead-letter queuedelivery limit hit: dead-letterAlert on DLQ depth; redrive only after fixing the cause.
- Committed with business dataOutbox row written in the same DB transaction as the state change, so no event is lost or invented.
- Published and confirmedRelay publishes and waits for a broker confirm; it may publish twice after a crash.
- Stored and readyDurable on the broker, waiting in a queue, partition, or stream.
- In flight (leased)Delivered to one consumer; invisible to others until ack or lease expiry.
- Processed idempotentlyHandler records the event ID with its side effect, so a repeat is a no-op.
- Acked or committedBroker deletes it or the group offset moves past it.
- Redelivered on failureLease expiry, nack, or consumer disconnect returns it with a higher delivery count.
- Dead-letteredAfter the delivery limit it moves to a DLQ (or is parked) for inspection and redrive.
KEY TERMS
- At-most / at-least / exactly-once
- Ack-then-process can lose work; process-then-ack can repeat it. Brokers that say exactly-once mean a scoped guarantee, not a property of your HTTP calls or emails.
- Ack and lease
- SQS visibility timeout (default 30 s, max 12 h), JetStream AckWait, Pub/Sub ack deadline, RabbitMQ channel lifetime, Redis Streams pending entry list. Long jobs must extend the lease or expect a second worker.
- Consumer groups and partitions
- Kafka assigns each partition to one consumer in a group, so partition count caps parallelism. Queues (SQS, RabbitMQ) instead compete per message, which spreads load but drops ordering.
- Ordering scope
- Order holds per Kafka partition (key), SQS FIFO message group, or Pub/Sub ordering key. Retries, DLQs, and parallel consumers all break wider ordering assumptions.
- Dedupe windows
- SQS FIFO drops resends with the same deduplication ID for 5 minutes. Kafka's idempotent producer stops retry duplicates in the log; its transactions plus read_committed give exactly-once only for Kafka-to-Kafka processing.
- Dead-letter queue
- Where a message goes after N failed deliveries (SQS maxReceiveCount, quorum queue delivery-limit, JetStream MaxDeliver). Without one, a poison message blocks or loops forever.
- Transactional outbox
- Write the event to an outbox table in the same transaction as the state change, then relay it. It removes lost or phantom events, not duplicates, so consumers stay idempotent.
IN YOUR STACK
TypeScript · AWS SDK v3 (SQS) Server
import { SQSClient, ReceiveMessageCommand, DeleteMessageCommand } from "@aws-sdk/client-sqs";
const sqs = new SQSClient({});
const QueueUrl = process.env.QUEUE_URL!;
while (true) {
const { Messages = [] } = await sqs.send(new ReceiveMessageCommand({
QueueUrl,
MaxNumberOfMessages: 10,
WaitTimeSeconds: 20, // long polling
VisibilityTimeout: 60, // lease: must exceed worst-case handling time
MessageSystemAttributeNames: ["ApproximateReceiveCount"],
}));
for (const msg of Messages) {
const event = JSON.parse(msg.Body!);
await handleOnce(event.eventId, event); // idempotent by business key
// The ack: until this succeeds, the message comes back after the timeout.
await sqs.send(new DeleteMessageCommand({ QueueUrl, ReceiptHandle: msg.ReceiptHandle! }));
}
}- Standard queues are at-least-once even inside the visibility timeout; FIFO queues dedupe sends for 5 minutes and order per MessageGroupId, but a slow consumer still gets a redelivery.
- A redrive policy with maxReceiveCount moves repeat failures to the DLQ; extend long jobs with ChangeMessageVisibility rather than a huge timeout.
- AttributeNames is deprecated on ReceiveMessage; use MessageSystemAttributeNames.
TypeScript · RabbitMQ (amqplib) Server
import amqp from "amqplib";
const conn = await amqp.connect(process.env.AMQP_URL!);
const ch = await conn.createChannel();
await ch.assertExchange("orders.dlx", "fanout", { durable: true });
await ch.assertQueue("orders", {
durable: true,
arguments: { "x-queue-type": "quorum", "x-dead-letter-exchange": "orders.dlx" },
});
await ch.prefetch(50); // max unacked deliveries in flight on this channel
await ch.consume("orders", async (msg) => {
if (!msg) return; // consumer was cancelled by the broker
try {
await handleOnce(JSON.parse(msg.content.toString()), msg.fields.redelivered);
ch.ack(msg);
} catch {
ch.nack(msg, false, false); // requeue=false: dead-letter instead of hot-looping
}
}, { noAck: false });- Unacked deliveries are requeued when the channel or connection closes, flagged redelivered; there is no per-message timer like SQS.
- nack with requeue=true on a persistent failure creates a tight redelivery loop. Quorum queues default to a delivery-limit of 20 (RabbitMQ 4.0+), then drop or dead-letter.
- Prefetch 0 means unlimited; without a limit one slow consumer hoards the backlog.
Java · Apache Kafka client Server
Properties props = new Properties();
props.setProperty("bootstrap.servers", "localhost:9092");
props.setProperty("group.id", "billing");
props.setProperty("enable.auto.commit", "false");
props.setProperty("isolation.level", "read_committed"); // skip aborted txn writes
props.setProperty("key.deserializer", StringDeserializer.class.getName());
props.setProperty("value.deserializer", StringDeserializer.class.getName());
try (var consumer = new KafkaConsumer<String, String>(props)) {
consumer.subscribe(List.of("orders"));
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
// Same key -> same partition -> ordered for that key only.
handleOnce(r.topic(), r.partition(), r.offset(), r.value());
}
consumer.commitSync(); // commit after processing = at-least-once
}
}- Offsets are the ack: committing before processing is at-most-once, after is at-least-once. A rebalance between process and commit replays the batch.
- Exactly-once (idempotent producer, transactions, sendOffsetsToTransaction, read_committed) covers Kafka-to-Kafka pipelines; writes to a DB or API still need their own idempotency.
- There is no built-in DLQ for plain consumers; publish failures to a dead-letter topic yourself (Kafka Connect and Kafka Streams have their own handlers).
Python · Redis Streams (redis-py) Server
import redis
r = redis.Redis(decode_responses=True)
STREAM, GROUP, ME = "orders", "billing", "worker-1"
try:
r.xgroup_create(STREAM, GROUP, id="0", mkstream=True)
except redis.exceptions.ResponseError:
pass # BUSYGROUP: the group already exists
while True:
# Take over entries another consumer read but never acked (idle > 60 s).
_, claimed, _deleted = r.xautoclaim(STREAM, GROUP, ME, min_idle_time=60_000,
start_id="0-0", count=50)
fresh = r.xreadgroup(GROUP, ME, {STREAM: ">"}, count=50, block=5_000)
entries = claimed + [e for _stream, batch in fresh for e in batch]
for entry_id, fields in entries:
handle_once(fields["eventId"], fields)
r.xack(STREAM, GROUP, entry_id)- Read-but-unacked entries sit in the group's pending entries list forever; nothing redelivers them unless a consumer runs XAUTOCLAIM or XCLAIM.
- Each claim increments the delivery counter (XPENDING shows it). There is no DLQ: past a threshold, copy the entry to a dead-letter stream and XACK it.
- Trimming the stream (MAXLEN) can delete entries that are still pending; XAUTOCLAIM returns those IDs as deleted.
TypeScript · NATS JetStream (nats.js v3) Server
import { connect, nanos } from "@nats-io/transport-node";
import { jetstream, jetstreamManager, AckPolicy } from "@nats-io/jetstream";
const nc = await connect({ servers: "nats://localhost:4222" });
const jsm = await jetstreamManager(nc);
await jsm.consumers.add("ORDERS", {
durable_name: "billing",
ack_policy: AckPolicy.Explicit,
ack_wait: nanos(30_000), // the lease, in nanoseconds
max_deliver: 5, // default -1: redeliver forever
});
const consumer = await jetstream(nc).consumers.get("ORDERS", "billing");
for await (const m of await consumer.consume({ max_messages: 50 })) {
try {
await handleOnce(m.json(), m.info.deliveryCount);
m.ack();
} catch (err) {
if (m.info.deliveryCount >= 5) m.term(); // stop redelivery
else m.nak(5_000); // retry after 5 s
}
}- Call m.working() during long handlers to reset the ack timer instead of raising ack_wait for everyone.
- Messages that reach max_deliver stay in the stream and the server emits an advisory; there is no automatic DLQ subject.
- Publisher-side dedupe uses the Nats-Msg-Id header within the stream's duplicate window; it does not dedupe consumer redeliveries.
Python · Google Cloud Pub/Sub Server
from google.cloud import pubsub_v1
from google.cloud.pubsub_v1.subscriber import exceptions as sub_exceptions
subscriber = pubsub_v1.SubscriberClient()
# Pull subscription created with enable_exactly_once_delivery=True.
path = subscriber.subscription_path("my-project", "orders-billing")
def callback(message: pubsub_v1.subscriber.message.Message) -> None:
handle_once(message.attributes["event_id"], message.data)
try:
# Only a successful ack future guarantees no redelivery.
message.ack_with_response().result()
except sub_exceptions.AcknowledgeError as e:
log.warning("ack failed (%s); expect redelivery", e.error_code)
flow = pubsub_v1.types.FlowControl(max_messages=100)
future = subscriber.subscribe(path, callback=callback, flow_control=flow)
with subscriber:
future.result()- Exactly-once delivery is regional and pull-only; it stops redelivery after a successful ack but not publish-side duplicates, which arrive with new message IDs.
- A missed ack deadline or a nack still redelivers; the client library extends deadlines while the callback runs, up to its configured maximum.
- Ordering keys order messages per key; a dead-letter topic is configured on the subscription with max delivery attempts.
WHERE IT BITES
- Acking before the side effect commits (auto-ack, auto-commit, noAck) turns every crash into silent data loss.
- Using the broker's message ID as the idempotency key: producer retries and outbox relays create new IDs for the same event. Dedupe on a business event ID stored with the side effect.
- A lease shorter than the slowest handler, or a GC pause, makes two workers process the same message concurrently; extend the lease or use a DB unique constraint.
- Requeueing a poison message immediately with no delivery limit burns CPU and blocks ordered groups; cap attempts and dead-letter.
- Assuming global order: parallel consumers, retries, and DLQ redrives reorder events. Carry a version or sequence and reject stale updates.
- Treating a DLQ as a bin: without alerts on depth and a tested redrive path, failures just disappear more slowly.
CLOSE THE AI. EXPLAIN THIS.
Your consumer charges a card, then crashes before acking. Walk through what SQS, Kafka, and RabbitMQ each do next, and name the one piece of state that stops the second charge.WHEN IT BREAKS IN PRODUCTION
SOURCES
- Amazon SQS visibility timeout (opens in new tab)Checked
- Exactly-once processing in Amazon SQS (FIFO) (opens in new tab)Checked
- RabbitMQ: Consumer acknowledgements and publisher confirms (opens in new tab)Checked
- Apache Kafka design: message delivery semantics (opens in new tab)Checked
- Pub/Sub exactly-once delivery (opens in new tab)Checked
- Pattern: Transactional outbox (opens in new tab)Checked
Explainer reviewed