RabbitMQ Work Queues and Exchange Patterns Implementation Guide

Work Queues Pattern

The Work Queues pattern involves multiple consumers subscribing to a single queue to collectively process messages. When message processing is time-intensive, production speed may exceed consumption speed, leading to message accumulation.

Using the Work Queues model allows multiple consumers to share the workload, significantly improving processing speed.

Message Production

To simulate high message volume, implement a loop that sends numerous messages:

@Test
public void demonstrateWorkQueues() throws InterruptedException {
    String queueName = "task.queue";
    String baseMessage = "processing_task_";
    
    for (int counter = 0; counter < 50; counter++) {
        rabbitTemplate.convertAndSend(queueName, baseMessage + counter);
        Thread.sleep(20); // Send every 20ms
    }
}

Message Consumption

Create multiple consumers sharing the same queue:

@RabbitListener(queues = "task.queue")
public void processTaskFirst(String message) throws InterruptedException {
    System.out.println("Worker A received: [" + message + "] at " + LocalTime.now());
    Thread.sleep(20); // Simulate quick processing
}

@RabbitListener(queues = "task.queue")
public void processTaskSecond(String message) throws InterruptedException {
    System.err.println("Worker B received: [" + message + "] at " + LocalTime.now());
    Thread.sleep(200); // Simulate slow processing
}

The workers have different processing speeds - Worker A processes 50 messages per second while Worker B handles only 5 messages per second. Without proper configuration, messages distribute evenly regardless of processing capacity, causing inefficiency.

Fair Dispatch Configuration

Configure the consumer's application.yml to implement fair dispatch:

spring:
  rabbitmq:
    listener:
      simple:
        prefetch: 1

This setting ensures each consumer processes one message at a time, allowing faster consumers to handle more work while slower ones manage their load appropriately.

Exchange Types

In advanced RabbitMQ patterns, exchanges route messages from publishers to queues. Exchanges don't store messages but forward them based on routing rules.

Fanout Exchange

Fanout exchanges broadcast messages to all bound queues without considering routing keys.

Setup

Create a fanout exchange named broadcast.exchange and bind two queues to it.

Publishing

@Test
public void sendBroadcastMessage() {
    String exchange = "broadcast.exchange";
    String content = "broadcast notification";
    rabbitTemplate.convertAndSend(exchange, "", content);
}

Consuming

@RabbitListener(queues = "broadcast.queue.one")
public void receiveFromQueueOne(String message) {
    System.out.println("Subscriber 1 received: [" + message + "]");
}

@RabbitListener(queues = "broadcast.queue.two")
public void receiveFromQueueTwo(String message) {
    System.out.println("Subscriber 2 received: [" + message + "]");
}

Direct Exchange

Direct exchanges route messages based on exact routing key matches betwean the message and queue bindings.

Configuration Example

  • Create routing.exchange as direct type
  • Bind direct.queue.a with routing keys green and purple
  • Bind direct.queue.b with routing keys green and orange

Publishing

@Test
public void sendDirectMessage() {
    String exchange = "routing.exchange";
    String message = "priority alert";
    rabbitTemplate.convertAndSend(exchange, "green", message);
}

Topic Exchange

Topic exchanges support wildcard routing keys using * for single words and # for multiple words.

Wildcard Examples

  • product.* matches product.create
  • order.# matches order.customer.update

Publishing

@Test
public void sendTopicMessage() {
    String exchange = "topic.exchange";
    String message = "important update";
    rabbitTemplate.convertAndSend(exchange, "us.orders.new", message);
}

Queue and Exchange Declaration

Instead of manual creation through the management interface, declare queues and exchanges programmatically:

@Configuration
public class ExchangeConfiguration {
    
    @Bean
    public TopicExchange topicExchange() {
        return new TopicExchange("dynamic.topic");
    }
    
    @Bean
    public Queue primaryQueue() {
        return new Queue("primary.queue");
    }
    
    @Bean
    public Binding queueBinding(Queue primaryQueue, TopicExchange topicExchange) {
        return BindingBuilder.bind(primaryQueue)
            .to(topicExchange)
            .with("*.notifications");
    }
}

Message Serialization

Default JDK serialization has performance and security isues. Configure JSON serialization instead:

@Bean
public MessageConverter jsonMessageConverter() {
    Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter();
    converter.setCreateMessageIds(true);
    return converter;
}

Publisher Reliability

Retry Configuration

Handle network failures with automatic retries:

spring:
  rabbitmq:
    template:
      retry:
        enabled: true
        initial-interval: 1000ms
        multiplier: 1
        max-attempts: 3

Confirmation Mechanisms

Enable publisher confirms and returns:

spring:
  rabbitmq:
    publisher-confirm-type: correlated
    publisher-returns: true

Configure return callbacks:

@PostConstruct
public void setupReturnCallback() {
    rabbitTemplate.setReturnsCallback(returned -> {
        log.error("Return callback triggered");
        log.debug("Exchange: {}", returned.getExchange());
        log.debug("Routing Key: {}", returned.getRoutingKey());
    });
}

Consumer Reliability

Acknowledgment Modes

Configure acknowledgment behavior:

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: auto

Retry Mechanisms

Implement local retries to prevent message requeuing:

spring:
  rabbitmq:
    listener:
      simple:
        retry:
          enabled: true
          initial-interval: 1000ms
          max-attempts: 3

Handling Failed Messages

Route failed messages to error queues:

@Bean
public MessageRecoverer errorHandler(RabbitTemplate template) {
    return new RepublishMessageRecoverer(template, "error.exchange", "failed");
}

Message Idempotency

Prevent duplicate processing using business state checks:

public void processPayment(Long orderId) {
    lambdaUpdate()
        .set(Order::getStatus, PAYED)
        .set(Order::getPayTime, LocalDateTime.now())
        .eq(Order::getId, orderId)
        .eq(Order::getStatus, UNPAID)
        .update();
}

Alternatively, use unique message IDs:

@Bean
public MessageConverter idGenerator() {
    Jackson2JsonMessageConverter converter = new Jackson2JsonMessageConverter();
    converter.setCreateMessageIds(true);
    return converter;
}

Tags: RabbitMQ messaging Work Queues Exchanges Fanout Exchange

Posted on Thu, 23 Jul 2026 16:40:36 +0000 by connex