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

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.

AspectBehavior
IdempotencyIdentical args → ok; mismatch → 406 precondition-failed.
Passive declarepassive=true: exists → ok; missing → 404 not-found.
Server-named queuesqueue.declare("") → the server mints amq.gen-{uuid22}.
Exclusive queuesConnection-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 queuesDeleted when the last (cluster-aware) consumer cancels.
Enforced argumentsx-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.

TriggerResult
Duplicate consumer tag on a channel530 not-allowed (connection error, RabbitMQ dialect)
Exclusive consumer over existing consumers — on any node403 access-refused
Any consume while a peer node holds an exclusive consumer403 access-refused (in exclusive use)
Non-integer x-priority argument406 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

MethodEffect
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 tag406 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.

ScopeglobalMeaning
Per-consumerfalse (RabbitMQ default)The budget applies to each consumer.
Per-channeltrueThe budget is shared across all consumers on the channel.
  • prefetch-size != 0 is refused 540 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

MethodBehavior
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.purge returns the exact message count purged.
  • queue.delete honors if-unused / if-empty, sends a server-initiated basic.cancel(consumerTag) to live consumers (on every node), and returns the residual count. Deleting a queue that does not exist replies delete-ok with a count of 0 and performs no side effects.
  • if-empty is refused 406 rather 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.

Was this page helpful?

On this page