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.exchangeas direct type - Bind
direct.queue.awith routing keysgreenandpurple - Bind
direct.queue.bwith routing keysgreenandorange
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.*matchesproduct.createorder.#matchesorder.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;
}