KubeMQ
ConnectorsRabbitMQ (AMQP 0-9-1)How-to guides

Reliability

Delivery guarantees on the KubeMQ RabbitMQ connector — publisher confirms, mandatory/return, dead-letter exchanges, per-message TTL, and at-least-once delivery.

This guide covers publisher confirms, mandatory / basic.return, dead-letter exchanges (DLX), per-message TTL, delayed delivery (x-delay), and at-least-once delivery — plus the gotchas that bite RabbitMQ migrants. Several behaviours deviate from RabbitMQ; the most dangerous is the publish-then-close loss below.

Fire-and-forget publishes followed by an immediate close are silently lost unless you use publisher confirms (or keep the connection open until the consumer has drained). This is the single most dangerous behaviour for a multi-message producer — it fails with no error, no nack, and no log on the client. See Publish-then-close.

Publisher confirms

confirm.select enables a per-channel, monotonically-increasing publish sequence from 1. A sequence is acked only after all routed queues accept the message:

  • single routed queue → after the queue send;
  • multiple routed queues → after the batch send for all queues.

multiple=true coalesces consecutive acked sequences; confirmations may arrive out of order.

Publisher confirms have no rollback. On any-queue failure the connector sends basic.nack(seq) — but queues that already accepted the message stay delivered. A naive retry on a nack may therefore duplicate the message in the queues that already got it. Design retries to be idempotent.

Transactions (tx.*)

tx.select, tx.commit and tx.rollback work with the semantics measured on RabbitMQ 4.3.4, so frameworks that turn them on (Spring AMQP's channelTransacted=true) run unchanged:

  • tx.select makes the channel transactional and it stays transactional: every commit or rollback opens the next transaction. A second tx.select is a no-op.
  • Transaction mode and confirm mode are mutually exclusive — switching either way is 406 (cannot switch from confirm to tx mode / cannot switch from tx to confirm mode).
  • tx.commit or tx.rollback on a channel that never selected → 406 channel is not transactional.
  • A publish inside a transaction is validated when it arrives (a missing exchange is a 404 right then, a refused permission a 403) but is routed and stored only at commit; rollback discards it. An unroutable mandatory publish gets its basic.return at commit.
  • An ack, nack or reject inside a transaction is validated when it arrives and takes effect at commit. Rollback forgets it: the delivery stays unacknowledged — it is not requeued — and may be settled again.
  • Uncommitted work dies with the channel.

As on RabbitMQ, a transaction is not atomic across queues and adds no durability a publish does not already have. One deviation: an open transaction holds at most 65,536 operations and 128 MiB of message bodies; the operation that crosses either bound closes the channel 406 transaction too large — commit more often.

Publish-then-close silently loses unconfirmed messages

This is the single most dangerous behaviour for a multi-message producer to get wrong, because it fails silently — no error, no nack, no log on the client.

Fire-and-forget publishes then an immediate close are silently lost. A basic.publish without confirms does not block: it only hands the message to the connector's per-channel executor queue, which sends to the KubeMQ queue asynchronously. If you close the channel or connection before that executor has drained, every still-buffered publish is abandoned — never sent, with no error returned to the client. The buffer holds up to 64 pending publishes, so a tight publish-then-close loop can drop dozens of messages at once.

In an internal test run against this connector, with no confirms: 30 fire-and-forget publishes immediately followed by a close lost 16; 100 lost 84. These exact counts are timing-dependent (they reflect a race between the executor's send rate and how soon you close) and will vary by environment — the only guarantee is "more than zero." On a confirm channel that waits for acks before closing, the same loops lose 0, which is the invariant to rely on.

The fix — pick one

  1. Use a confirm channel and wait for all acks before closing (recommended). Call confirm.select, publish, then block until every publish is acked. A publish is acked only after the connector has actually sent it to the queue, so waiting for confirms forces the executor queue to drain before you close.
  2. Or keep the connection open until the consumer has drained. If you genuinely cannot use confirms, don't close immediately after publishing — keep the channel/connection alive until you have independent evidence (a consumer ack, a queue-depth check) that the messages were ingested.

Remember: confirms have no rollback, so make any retry idempotent.

// amqp091-go — confirm mode; wait for acks before closing.
_ = ch.Confirm(false)
confirms := ch.NotifyPublish(make(chan amqp.Confirmation, 1))
_ = ch.Publish("", "orders", false, false, amqp.Publishing{Body: body})
if c := <-confirms; !c.Ack {
    log.Println("nacked — retry idempotently")
}
// only now is it safe to close
# pika — confirm mode; BlockingChannel raises on a nack.
channel.confirm_delivery()
try:
    channel.basic_publish(exchange="", routing_key="orders", body=body)
except pika.exceptions.UnroutableError:
    ...  # retry idempotently
# publish_delivery blocks until confirmed, so it is now safe to close
// amqp-client — confirm mode; block until all publishes are confirmed.
channel.confirmSelect();
channel.basicPublish("", "orders", null, body);
channel.waitForConfirmsOrDie(5_000); // drains the executor queue before close
// amqplib — ConfirmChannel; await each publish callback before closing.
const ch = await connection.createConfirmChannel();
await new Promise<void>((resolve, reject) => {
  ch.publish("", "orders", body, {}, (err) => (err ? reject(err) : resolve()));
});
// safe to close now
// RabbitMQ.Client — confirm mode; wait for confirms before closing.
channel.ConfirmSelect();
channel.BasicPublish("", "orders", body: body);
channel.WaitForConfirmsOrDie(TimeSpan.FromSeconds(5));
# bunny — confirm mode; wait for confirms before closing.
channel.confirm_select
exchange.publish(body, routing_key: "orders")
channel.wait_for_confirms # blocks until the executor queue drains
// lapin — publisher confirms; await the returned confirmation.
let confirm = channel
    .basic_publish("", "orders", BasicPublishOptions::default(), body,
        BasicProperties::default())
    .await?
    .await?; // second await resolves the confirm
// confirm is now Ack/Nack — safe to close

Mandatory / return

basic.publish(mandatory=true) on an unroutable message returns basic.return(312 NO_ROUTE) with the full message content, sent before the ack in confirm mode. Without mandatory, an unroutable message is silently dropped.

Dead-letter exchange (DLX)

Configure with the queue arguments x-dead-letter-exchange / x-dead-letter-routing-key.

Dead-lettering fires on three triggers, as on RabbitMQ:

Triggerx-death reasonWhen
Rejectionrejectedbasic.reject / basic.nack(requeue=false); also a message that exceeds x-delivery-limit (RabbitMQ's delivery_limit reason string does not exist here)
ExpiryexpiredA per-message expiration or queue-level x-message-ttl that has passed — see Per-message TTL for when it is noticed
Length limitmaxlenA publish refused at an x-overflow: reject-publish-dlx queue's x-max-length

When a message is dead-lettered:

  • the x-death array is RabbitMQ-exact (queue, reason, time, exchange, routing-keys, count), most-recent-first;
  • the x-first-death-* / x-last-death-* convenience headers are set;
  • the original expiration property is moved to x-death[0].original-expiration;
  • a missing DLX exchange, or one that routes nowhere (declared, but with no binding matching the dead-letter routing key) → the message is dropped with a WARN and the original is acked. Declare the dead-letter exchange and its binding before messages start expiring.

The dead-letter hop ceiling destroys the message. DeadLetterMaxHops (default 16) caps the x-death count per (queue, reason). On a work-queue ↔ wait-queue retry ladder that count goes up once per cycle, so with the default a message gets 16 retries and is then dropped — not dead-lettered onward to a failure queue. A RabbitMQ ladder written to retry 20 times before parking a message loses it on the 17th cycle here. Count your cycles and set CONNECTORS_AMQP_DEAD_LETTER_MAX_HOPS above them; the drop is logged, the message is gone.

x-delivery-limit (RabbitMQ's quorum-queue poison cap) is enforced: a limit of N permits N deliveries and the (N+1)th diverts the message to the DLX. The declare is refused 406 unless the queue's x-dead-letter-exchange resolves to exactly one queue — otherwise the message would be silently destroyed at the limit instead of dead-lettered. 0 is refused, and a value above Queue.MaxReceiveCount (default 1024) is refused naming the ceiling.

Per-message TTL

Set the expiration property to milliseconds as a numeric string (^\d+$, else 406), or a queue-level x-message-ttl on queue.declare; where both are set the smaller applies, as on RabbitMQ. What happens at the deadline depends on whether the queue has somewhere to dead-letter to:

QueueAt the deadlinePrecision
No x-dead-letter-exchangeThe expired message is discarded, as RabbitMQ doesRounded up to whole seconds and clamped to Queue.MaxExpirationSeconds (default 12h); an x-message-ttl above that ceiling is refused 406 at declare
With x-dead-letter-exchangeThe expired message is dead-lettered with reason expired and original-expiration, matching RabbitMQHonored exactly as sent, to the millisecond, with no ceiling

Expiry is eager on one shape and lazy on everything else. A background sweep dead-letters expired messages with nothing consuming — but only for a queue with a queue-level x-message-ttl greater than 0, an x-dead-letter-exchange, and no consumer on any node. That is exactly the canonical wait-queue retry topology, and it fires on its own. Two shapes stay lazy and are noticed only when a consumer or basic.get reaches the message: a wait queue driven by a per-message expiration, and a queue with a TTL but no dead-letter exchange. The sweep visits at most 8 candidate queues per 30-second tick, so the worst case from deadline to dead-lettering is ceil(candidates / 8) × 30 s — pick your shortest retry step with that number in mind. basic.get dead-letters at most 8 expired messages per call before answering get-empty; nothing is lost, the next call continues.

Two more divergences: x-message-ttl: 0 means "no TTL" here (RabbitMQ: expire-on-arrival), and on a dead-letter-bearing queue the deadline binds AMQP readers only — a gRPC, REST or dashboard reader of the same channel can still receive an expired message (stale, never lost).

x-expires (idle-queue expiry) is enforced too: a queue with no consumer, no redeclare and no basic.get for its window is deleted, and every message on it destroyed. Publishing does not count as use. Do not set it on a queue whose reader is periodic rather than continuous — see Migrating from RabbitMQ.

Delayed delivery (x-delay)

Set the x-delay header to milliseconds; the connector computes ceil(ms/1000) seconds, clamped to a per-queue maximum (12h), and strips the x-delay header on delivery (matching the RabbitMQ delayed-message-exchange plugin).

At-least-once delivery

Unacked deliveries are requeued on disconnect — zero loss, even on an ungraceful disconnect. Graceful shutdown sends connection.close(320), nacks pending, and requeues unacked. Durable queues persist across restart with Redelivered == true on recovery. Exactly-once is NOT provided. See Queues and consumers.

Cluster behavior

Exclusivity and direct reply-to are cluster-wide. Exclusive consumers, exclusive queue names and single-active queues are held across the cluster by gossip, and a Direct Reply-To reply published on another node is forwarded to the requester's node. Two nodes that grant the same thing inside one gossip delay settle it deterministically (the earlier claim stands; the other side cancels its consumers or evicts its queue), and a node silent for 20 seconds is declared dead and its claims freed. Two residues: single-active takeover order is by node presence, not individual registration, and consumer priorities stay per node. No load-balancer session affinity is needed.

The idle clock behind x-expires is the one thing that is still node-local: a consumer on any node protects the queue everywhere, but a queue kept alive only by basic.get traffic on one node can be expired by another. Keep a consumer attached rather than polling.

Error quick reference

TriggerCode
mandatory=true + unroutable312
expiration not ^\d+$406
tx.select on a confirm-mode channel, or confirm.select on a transactional one406
tx.commit / tx.rollback without tx.select406
Transaction over 65,536 operations or 128 MiB of bodies406
x-max-length without x-overflow, or with drop-head406 at declare
x-delivery-limit whose DLX does not resolve to exactly one queue406 at declare
Graceful shutdown / connection limit320
Fire-and-forget publish then immediate closenone — silently dropped

Was this page helpful?

On this page