The design should support reliable message delivery, multiple producers and consumers, retries, and horizontal scalability.
Requirements
A correct message queue should meet the following requirements.1. Allow producers to publish messages.
2. Allow consumers to receive messages.
3. Support multiple producers and consumers.
4. Support reliable message delivery.
5. Support message acknowledgements.
6. Retry failed messages.
7. Support Dead Letter Queue (DLQ).
8. Maintain message ordering where required.
Design
The system consists of a Message Broker that receives messages from producers and delivers them to consumers. Messages are persisted before delivery to prevent data loss.Consumers acknowledge messages after successful processing. Failed messages are retried or moved to a Dead Letter Queue.
Producer 1 Producer 2
β β
βββββββββ¬ββββββββ
βΌ
Message Broker
β
ββββββββββββ΄βββββββββββ
β β
βΌ βΌ
Consumer 1 Consumer 2
Java Implementation
The following implementation demonstrates a simple message broker using a producer-consumer model, where producers publish messages and consumers process them asynchronously.Message
The Message class represents a message containing a unique identifier and the payload to be processed.public class Message {
private final String id;
private final String payload;
public Message(String id, String payload) {
this.id = id;
this.payload = payload;
}
public String getId() {
return id;
}
public String getPayload() {
return payload;
}
}
Message Queue
The MessageQueue provides a thread-safe queue for publishing and consuming messages using a BlockingQueue.import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
public class MessageQueue {
private final BlockingQueue queue = new LinkedBlockingQueue<>();
public void publish(Message message) throws InterruptedException {
queue.put(message);
}
public Message consume() throws InterruptedException {
return queue.take();
}
}
Producer
The Producer creates messages and publishes them to the message queue.public class Producer {
private final MessageQueue queue;
public Producer(MessageQueue queue) {
this.queue = queue;
}
public void publish(String message) throws InterruptedException {
queue.publish(new Message(String.valueOf(System.nanoTime()), message));
}
}
Consumer
The Consumer continuously retrieves messages from the queue and processes them in a separate thread.public class Consumer implements Runnable {
private final MessageQueue queue;
public Consumer(MessageQueue queue) {
this.queue = queue;
}
@Override
public void run() {
while (true) {
try {
Message message = queue.consume();
System.out.println(Thread.currentThread().getName() + " processed " + message.getPayload());
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
break;
}
}
}
}
Example
This example starts two consumer threads and publishes multiple messages, which are processed asynchronously as they become available.public class Main {
public static void main(String[] args) {
MessageQueue queue = new MessageQueue();
Producer producer1 = new Producer(queue);
Producer producer2 = new Producer(queue);
new Thread(new Consumer(queue), "Consumer-1").start();
new Thread(new Consumer(queue), "Consumer-2").start();
new Thread(() -> {
try {
producer1.publish("Order Created");
producer1.publish("Order Packed");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "Producer-1").start();
new Thread(() -> {
try {
producer2.publish("Payment Completed");
producer2.publish("Order Shipped");
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}, "Producer-2").start();
}
}
Output:
The output shows how the published messages are distributed among the available consumer threads for processing.
Consumer-1 processed Order Created
Consumer-2 processed Payment Completed
Consumer-1 processed Order Packed
Consumer-2 processed Order Shipped
Complexity
Publishing inserts a message into the queue. Consumption removes the next available message from the queue.Time Complexity
Publish β O(1)
Consume β O(1)
Space Complexity
O(n)
n is the number of messages currently stored in the queue.