Queues and Consumers
Declaring queues, consuming, acknowledging, prefetch (QoS), and basic.get on the KubeMQ RabbitMQ connector — every AMQP queue maps to a KubeMQ Queue channel.
Every AMQP queue maps to a KubeMQ Queue channel amqp.{vhost}.{queue}. This guide covers
declaring queues, consuming, acknowledging (ack / nack / reject), prefetch (QoS), and
basic.get.
Consumption is at-least-once: unacked deliveries are requeued on disconnect, so there is zero
loss even on an ungraceful disconnect. Exactly-once is NOT provided — plan for redelivery
(Redelivered == true). See Reliability.
Declaring queues
queue.declare supports durable / exclusive / auto-delete / arguments.
| Aspect | Behavior |
|---|---|
| Idempotency | Identical args → ok; mismatch → 406 precondition-failed. |
| Passive declare | passive=true: exists → ok; missing → 404 not-found. |
| Server-named queues | queue.declare("") → the server mints amq.gen-{uuid22}. |
| Exclusive queues | Connection-scoped. Cross-connection access (including passive declare) → 405 resource-locked. Auto-deleted when the owning connection closes. The name is claimed cluster-wide — a declare of a name a peer node holds exclusively is refused 405 too. |
| Auto-delete queues | Deleted when the last (cluster-aware) consumer cancels. |
| Enforced arguments | x-message-ttl, x-expires, x-max-length + x-overflow, x-delivery-limit, x-single-active-consumer and the dead-letter arguments all change behavior; x-max-length without x-overflow is refused 406. See Capabilities. |
Exclusivity is cluster-wide. Exclusive queue names, exclusive consumers and single-active queues are held across the cluster by gossip. Two nodes that grant the same thing inside one gossip delay settle it deterministically — the earlier claim stands and the other side cancels its consumers or evicts its queue — and a node silent for 20 seconds is declared dead and its claims freed. See Reliability and Migration from RabbitMQ.
Consuming
basic.consume registers a consumer (an auto-generated ctag-{n} if the tag is empty). The
server then delivers basic.deliver + content header + body.
| Trigger | Result |
|---|---|
| Duplicate consumer tag on a channel | 530 not-allowed (connection error, RabbitMQ dialect) |
| Exclusive consumer over existing consumers — on any node | 403 access-refused |
| Any consume while a peer node holds an exclusive consumer | 403 access-refused (in exclusive use) |
Non-integer x-priority argument | 406 precondition-failed |
A minimal consume loop:
// amqp091-go — consume and manually ack each delivery.
deliveries, _ := ch.Consume("orders", "", false /* autoAck */, false, false, false, nil)
for d := range deliveries {
process(d.Body)
_ = d.Ack(false) // multiple=false
}# pika — consume with manual ack.
def on_message(ch, method, props, body):
process(body)
ch.basic_ack(delivery_tag=method.delivery_tag)
channel.basic_consume(queue="orders", on_message_callback=on_message, auto_ack=False)
channel.start_consuming()// amqp-client — consume with manual ack.
DeliverCallback cb = (tag, delivery) -> {
process(delivery.getBody());
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
};
channel.basicConsume("orders", false /* autoAck */, cb, t -> {});// amqplib — consume with manual ack.
await channel.consume("orders", (msg) => {
if (!msg) return;
process(msg.content);
channel.ack(msg);
}, { noAck: false });// RabbitMQ.Client — consume with manual ack.
var consumer = new EventingBasicConsumer(channel);
consumer.Received += (_, ea) =>
{
Process(ea.Body.ToArray());
channel.BasicAck(ea.DeliveryTag, multiple: false);
};
channel.BasicConsume("orders", autoAck: false, consumer);# bunny — consume with manual ack.
queue.subscribe(manual_ack: true, block: true) do |delivery_info, _props, body|
process(body)
channel.ack(delivery_info.delivery_tag)
end// lapin — consume with manual ack.
let mut consumer = channel.basic_consume(
"orders", "", BasicConsumeOptions::default(), FieldTable::default()).await?;
while let Some(delivery) = consumer.next().await {
let delivery = delivery?;
process(&delivery.data);
delivery.ack(BasicAckOptions::default()).await?;
}Ack / nack / reject
| Method | Effect |
|---|---|
basic.ack(tag, multiple) | AckRange — message(s) consumed. |
basic.nack / basic.reject(tag, requeue=true) | NAckRange — requeued at the tail. |
basic.reject / basic.nack(tag, requeue=false) | Dropped, or dead-lettered if a DLX is configured. |
| Unknown delivery tag | 406 precondition-failed. |
Requeue lands at the tail. Requeued messages re-enter at the queue tail, not the head (a deviation from RabbitMQ classic head-requeue). Fairness ordering therefore differs.
At-least-once consumption
Unacked deliveries are requeued on disconnect (a downstream nack-all safety net), so there is
zero loss even on an ungraceful disconnect. Exactly-once is NOT provided — plan for
redelivery (Redelivered == true).
Prefetch (QoS)
basic.qos(prefetch-size, prefetch-count, global) limits in-flight unacked deliveries.
| Scope | global | Meaning |
|---|---|---|
| Per-consumer | false (RabbitMQ default) | The budget applies to each consumer. |
| Per-channel | true | The budget is shared across all consumers on the channel. |
prefetch-size != 0is refused540 not-implemented— RabbitMQ refuses a byte-based prefetch rather than ignoring it.- Default = unlimited.
- Effective allowance =
min(per-consumer budget, remaining channel-global budget).
basic.get (pull)
basic.get returns get-ok (with a delivery tag) or get-empty.
basic.get on an empty queue answers in about 50 ms. The connector runs a short broker
poll, so the get-empty is authoritative rather than a cached count; RabbitMQ answers in about
1 ms. A tight poll loop pays that on every miss — prefer basic.consume for throughput.
GetBatchSize (default 32) bounds per-Get pulls. On a queue with expired messages,
basic.get dead-letters at most 8 of them per call before answering get-empty.
Recover / flow
| Method | Behavior |
|---|---|
basic.recover(requeue=true) | Nack-all unacked on the channel (tail requeue). |
basic.recover-async(requeue=true) | The same requeue, with no reply (the deprecated asynchronous form). |
basic.recover(requeue=false) / basic.recover-async(requeue=false) | 540 not-implemented, as on RabbitMQ. |
channel.flow(active=true) | Replies flow-ok, takes no action. |
channel.flow(active=false) | 540 not-implemented — RabbitMQ dropped client-driven flow control long ago. |
Purge / delete
queue.purgereturns the exact message count purged.queue.deletehonorsif-unused/if-empty, sends a server-initiatedbasic.cancel(consumerTag)to live consumers (on every node), and returns the residual count. Deleting a queue that does not exist repliesdelete-okwith a count of0and performs no side effects.if-emptyis refused406rather than assumed satisfied when the count cannot be determined (the broker monitoring endpoint is unavailable) or while the connector still holds unacknowledged deliveries for the queue — that second case needs no outage at all, one manual-ack consumer with messages in flight is enough. Match on the reply text, and retry after the state settles. A queue that has never been published to is treated as empty.
Related
Exchanges and routing
How a publish resolves through default/direct/fanout/topic/headers exchanges into these queues.
Reliability
Publisher confirms, dead-letter exchanges, per-message TTL, and at-least-once delivery.
Error codes
The 404, 405, 406, and 530 AMQP codes the connector returns on declare, consume, and ack failures.
Was this page helpful?
Pub/Sub (Fanout)
Broadcast every message to all subscribers over AMQP 0-9-1 — a fanout exchange copies to exclusive queues, each backed by its own KubeMQ Queue channel.
Reliability
Delivery guarantees on the KubeMQ RabbitMQ connector — publisher confirms, mandatory/return, dead-letter exchanges, per-message TTL, and at-least-once delivery.