Writing

Messaging with RabbitMQ

notes on RabbitMQ messaging, including queues, exchanges, acknowledgements, and streams.

Messaging with RabbitMQ

A task can be encapsulated as a message and sent to a queue. The message is taken by a worker, the work is completed, and RabbitMQ is notified when processing finishes. This is the main idea behind messaging: the sender does not need to wait for the task to be completed.

RabbitMQ provides more than a queue. The destination of a message, the number of messages delivered to a consumer at once, the response to a failure, and whether messages are kept for replay can all be configured.

Connections and channels

A connection is a TCP connection between an application and RabbitMQ. Multiple channels can be opened on that connection. Channels are like lightweight connections that share the same TCP connection, but operations on one channel stay separate from operations on another. Each protocol method carries a channel ID so RabbitMQ knows where it belongs. When the connection closes, all its channels close too.

For applications that process messages in multiple threads or processes, it is common to use a separate channel for each one instead of sharing a channel between them.

How does a message reach a queue?

The publisher sends a message to an exchange. The exchange routes it to zero or more queues according to its type and its bindings. A binding is a rule connecting an exchange to a queue.

For example, to route messages from exchange orders to queue email-worker, the queue is bound to the exchange. The exchange then uses that binding when a message arrives.

The most common exchange types are:

  • Default exchange: Its name is the empty string (""). Every declared queue is automatically bound to it using the queue name as the routing key. Publishing to the default exchange with routing key email-worker sends the message to that queue.
  • Fanout: Sends each message to every queue bound to the exchange. The routing key is ignored. This is useful when several services need the same notification.
  • Direct: Sends a message to queues whose binding key exactly matches the routing key. More than one queue can use the same binding key, so this can also send one message to several queues.
  • Topic: Matches a routing key made of dot-separated words against binding patterns. * matches exactly one word, and # matches zero or more words. A binding of order.* matches order.created, while order.# also matches order.payment.failed.

An alternate exchange can handle messages that the original exchange cannot route. A dead-letter exchange is used when a message leaves a queue without being processed normally. For example, a message can be dead-lettered after a consumer rejects it without requeueing, after its TTL expires, when it is dropped because of a queue length limit, or when a quorum queue’s delivery limit is reached. If the entire queue expires, its messages are not dead-lettered.

Work queues, broadcasts, and topics

With a work queue, several consumers can subscribe to the same queue. Each delivery goes to one consumer. Tasks can be distributed among workers this way: one email is sent by one worker, the next email by another, and so on.

With publish/subscribe, each service has its own queue bound to a fanout exchange. If an order created event is published, it can be received by both the email service and the analytics service. They read from separate queues, so one service consuming its copy does not remove the other’s copy.

With a topic exchange, services can receive only the categories they care about. A billing queue might bind to order.payment.*, while an audit queue binds to order.#.

The important question is: does one worker need to handle the task, or do several services each need a copy? The queue and binding setup follows from that answer.

Prefetch: how much work should a consumer receive?

RabbitMQ can deliver messages faster than a consumer can process them. Prefetch count limits how many deliveries may remain unacknowledged at once.

If the prefetch count is 5, a consumer can have up to five unacknowledged deliveries. Once it acknowledges one, RabbitMQ can send another. This helps prevent a slow consumer from accumulating too much work in memory.

RabbitMQ normally applies a prefetch limit separately to each consumer on a channel. If three consumers share a channel and basic.qos(5, false) is set, each consumer can have up to five unacknowledged messages. With basic.qos(5, true), the consumers on that channel share a limit of five. The shared limit requires more coordination, so it can be slower.

Prefetch also affects how work is distributed. Suppose one worker needs ten seconds per task and another needs one second. If the first worker already has unacknowledged messages, a low prefetch limit lets the faster worker receive more of the remaining work. Setting prefetch to 1 is a simple way to see this behavior, but it can reduce throughput, especially when network latency is high. The right value depends on the workload and should be measured rather than assumed.

If all workers stay busy, the queue can still grow. In that case, queue length should be monitored, and more workers or a change in how work is produced may be needed.

Ordering and priority

Queues normally provide FIFO behavior. Messages published on one channel are enqueued in publishing order in every queue they are routed to. That does not mean processing will always finish in that order.

If several consumers read from the same queue, one consumer may finish later than another. Requeueing and redelivery can also change the observed order. Priority queues deliberately deliver higher-priority messages ahead of lower-priority ones. If strict processing order matters, consumers, prefetch, and retries must be accounted for in the design.

RabbitMQ supports message priorities for classic queues and, in current versions, quorum queues. Their priority behavior and supported levels differ, so the queue type should be checked before a specific range is relied upon.

Consumer acknowledgements

Networks fail and consumers can crash while processing a message. RabbitMQ needs to know when it is safe to stop holding that message.

There are two acknowledgement modes:

  1. Automatic acknowledgement: RabbitMQ considers the message handled as soon as it sends it to the consumer’s socket. If the consumer crashes before doing the work, that message can be lost. This is the fire-and-forget option.
  2. Manual acknowledgement: The consumer tells RabbitMQ what happened after processing. With basic.ack, it reports success. With basic.nack or basic.reject, it reports failure and chooses whether the message should be requeued. If it is not requeued, RabbitMQ dead-letters it when a dead-letter exchange is configured; otherwise it discards it.

basic.nack can reject multiple outstanding deliveries at once. basic.reject handles only one. Acknowledgements can also use the multiple flag to cover all outstanding delivery tags up to a specified tag. Delivery tags belong to a channel, so an acknowledgement must be sent on the same channel that received the delivery.

If a consumer’s channel or connection closes before it acknowledges a delivery, RabbitMQ requeues the unacknowledged message. Another consumer may then receive it. This is why consumers should be ready for redelivery and why processing should be idempotent: doing the work twice should not create two payments or two records.

For example, a user may be saved to the database just before the connection fails, preventing RabbitMQ from receiving the acknowledgement. The message may be sent again even though the user was already saved. An application-level way to recognize completed work is needed.

Publisher confirms

Consumer acknowledgements cover delivery from RabbitMQ to the consumer. Publisher confirms cover the other direction: from the publisher to RabbitMQ. These two mechanisms are independent.

After a publisher enables confirms on a channel, RabbitMQ sends an asynchronous ack or nack for published messages. For a persistent message routed to a durable queue, a confirm is sent after the broker has persisted it. A confirm tells the publisher that RabbitMQ has taken responsibility for the publish. It does not say that a consumer processed the message.

There is one detail that is easy to miss: RabbitMQ can also confirm an unroutable message. If it must be known whether a message reached any queue, the mandatory flag should be used and basic.return handled, or an alternate exchange should be configured. A confirm by itself is not proof that the message reached a queue.

Enabling confirms but ignoring the responses does not give the application useful delivery information. Outstanding messages should be tracked, and both confirms and failures should be handled. When the result is uncertain after a connection failure, a publisher may have to retry; that retry can create a duplicate.

At-most-once, at-least-once, and duplicates

RabbitMQ does not make an end-to-end exactly-once promise for ordinary queue processing. The behavior depends on what the publisher and consumer do.

  • At-most-once: Publish without tracking confirms and use automatic consumer acknowledgements. A failed publish or a consumer crash can lose a message, but RabbitMQ will not redeliver a message that was automatically acknowledged.
  • At-least-once: Track publisher confirms, retry publishes whose outcome is uncertain, and acknowledge a consumed message only after the work succeeds. This protects against loss in common failure cases, but retries and redeliveries can produce duplicates.

For an important operation, an idempotency key can be stored with the result. If the same message arrives again, earlier processing can be detected and the message acknowledged without repeating the side effect. This can give the application an effectively-once result for that operation, provided the idempotency check and the operation are designed together.

When a queue is not enough: RabbitMQ Streams

A traditional queue is useful when consumers take work and acknowledge it so it can be removed. A stream keeps messages for a retention period and lets consumers read them again. Consumers can start at an offset and replay earlier messages. This makes streams useful for event sourcing, log collection, and analytics.

Streams are persistent and replicated. A single stream preserves order within that stream. When scaling across partitions is needed, RabbitMQ provides super streams; order is then meaningful within each partition.

There are two ways to interact with a stream, and the difference matters:

  • With the RabbitMQ Stream protocol, clients publish directly to a stream and use stream-specific operations for reading and offsets. There is no exchange routing in that publish path.
  • With AMQP 0-9-1, a stream can be declared as a queue type and bound to an exchange. Super streams also use exchanges and bindings in their topology.

Streams therefore do not remove exchanges from RabbitMQ. They provide a different storage and consumption model, with a protocol optimized for using it.

A note on STOMP and MQTT

RabbitMQ also supports protocols such as STOMP and MQTT. These protocols often present messaging in terms of topics, but queues still matter: they buffer messages for consumers and are where many RabbitMQ features apply. This is useful for dynamic consumers such as WebSocket or mobile clients that subscribe and unsubscribe as needed.

Putting it together

For a task such as sending an email, a message can be published to an exchange and routed to a work queue. A measured prefetch limit can be set, and the message can be acknowledged only after the email work succeeds. Publisher confirms should be tracked, and the consumer should be idempotent because either side may need to retry.

If several services need the same event, a separate queue can be assigned to each service and bound to a fanout or topic exchange. If the event must be retained and replayed later, a stream can be considered.

The main idea stays simple: put work into a message, decide who should receive it, and decide what should happen when any step fails.

References