Lesson 6 / 25

Work Queues, Pub/Sub, Routing and RPC

Recognise the standard messaging patterns and implement request-reply.

Patterns built from the same parts

The same building blocks produce several standard patterns. Work queue (competing consumers): many workers consume one queue, each message goes to one worker, and adding workers increases throughput. Publish-subscribe: a fanout or topic exchange copies each message to one queue per subscriber. Routing: a direct or topic exchange sends messages to different queues by key, for example by severity or region. Request-reply (RPC): the client publishes a request with a reply_to queue name and a correlation_id, the server processes it and publishes the response to reply_to with the same correlation ID, and the client matches responses to requests. RabbitMQ's direct reply-to feature (the pseudo-queue amq.rabbitmq.reply-to) avoids creating a reply queue per client. Use RPC over a broker sparingly; a direct HTTP or gRPC call is usually simpler when the caller must wait for an answer anyway.

Request-reply with correlation IDs

The client waits for the response whose correlation_id matches its request.

import uuid, pika

corr_id = str(uuid.uuid4())
response = None

def on_reply(ch, method, props, body):
    global response
    if props.correlation_id == corr_id:
        response = body

channel.basic_consume(queue="amq.rabbitmq.reply-to", on_message_callback=on_reply, auto_ack=True)
channel.basic_publish(
    exchange="",
    routing_key="rpc.price-quote",
    properties=pika.BasicProperties(reply_to="amq.rabbitmq.reply-to", correlation_id=corr_id),
    body=b'{"sku": "notebook", "qty": 3}',
)
while response is None:
    connection.process_data_events(time_limit=1)   # add an overall timeout in real code

Always time out RPC

If the server is down, an RPC client waiting on a reply queue can hang forever. Set an overall timeout and decide what to show the user when it expires.

Quick check: In a work-queue pattern with three consumers on one queue, how many consumers receive each message?

  • Exactly one
  • All three
  • None until acknowledged
  • Two, for redundancy
Answer

Exactly one — Competing consumers share the queue; each message is delivered to one of them.