Producer-Consumer
One producer, multiple consumers. Each message processed once. Perfect for task distribution.
Synchronous systems: Service A calls Service B and waits for response.
Problems:
Solution: Message Queues - Asynchronous, decoupled communication!
A message queue is a buffer that stores messages between senders (producers) and receivers (consumers). It acts as an intermediary component that enables asynchronous, decoupled communication between services or components in a distributed system.
At its essence, a message queue is a temporary storage mechanism that decouples the production of messages from their consumption. Instead of a producer directly calling a consumer and waiting for a response, the producer sends a message to the queue and continues with other work. The consumer retrieves messages from the queue when it’s ready to process them, completely independently of when the messages were produced.
This decoupling provides several critical benefits: temporal decoupling (producer and consumer don’t need to be active at the same time), spatial decoupling (they don’t need to know each other’s location), and synchronization decoupling (they don’t need to wait for each other).
One producer, one or more consumers. Each message processed once.
Characteristics:
Use cases:
import queueimport threadingimport timefrom typing import Optional
class MessageQueue: """Simple message queue implementation"""
def __init__(self, maxsize: int = 1000): self.queue = queue.Queue(maxsize=maxsize) self.consumers = [] self.running = False
def produce(self, message: dict): """Producer: Add message to queue""" try: self.queue.put_nowait(message) print(f"Produced: {message}") except queue.Full: print("Queue full! Message rejected.")
def consume(self, handler): """Consumer: Process messages from queue""" while self.running: try: message = self.queue.get(timeout=1) # Process message handler(message) # Acknowledge self.queue.task_done() except queue.Empty: continue
def start_consumer(self, handler, consumer_id: int): """Start a consumer thread""" def consumer_loop(): print(f"Consumer {consumer_id} started") self.consume(handler)
thread = threading.Thread(target=consumer_loop, daemon=True) thread.start() self.consumers.append(thread)
def start(self): """Start queue processing""" self.running = True
def stop(self): """Stop queue processing""" self.running = False self.queue.join() # Wait for all tasks to complete
# Usagequeue = MessageQueue(maxsize=100)queue.start()
# Message handlerdef process_order(message): print(f"Processing order: {message['order_id']}") # Process order... time.sleep(1) # Simulate work print(f"Order {message['order_id']} processed")
# Start multiple consumersqueue.start_consumer(process_order, consumer_id=1)queue.start_consumer(process_order, consumer_id=2)
# Producer sends messagesfor i in range(10): queue.produce({'order_id': i, 'amount': 100 + i}) time.sleep(0.1)
time.sleep(5) # Let consumers processqueue.stop()import java.util.concurrent.BlockingQueue;import java.util.concurrent.LinkedBlockingQueue;import java.util.concurrent.ExecutorService;import java.util.concurrent.Executors;import java.util.function.Consumer;
public class MessageQueue { private final BlockingQueue<Message> queue; private final ExecutorService executor; private volatile boolean running = false;
public MessageQueue(int maxSize) { this.queue = new LinkedBlockingQueue<>(maxSize); this.executor = Executors.newCachedThreadPool(); }
public void produce(Message message) { // Producer: Add message to queue if (queue.offer(message)) { System.out.println("Produced: " + message); } else { System.out.println("Queue full! Message rejected."); } }
public void consume(Consumer<Message> handler) { // Consumer: Process messages from queue while (running) { try { Message message = queue.take(); // Blocks until message available handler.accept(message); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } }
public void startConsumer(Consumer<Message> handler, int consumerId) { // Start a consumer thread executor.submit(() -> { System.out.println("Consumer " + consumerId + " started"); consume(handler); }); }
public void start() { running = true; }
public void stop() { running = false; executor.shutdown(); }}
// UsageMessageQueue queue = new MessageQueue(1000);queue.start();
// Message handlerConsumer<Message> processOrder = message -> { System.out.println("Processing order: " + message.getOrderId()); // Process order... try { Thread.sleep(1000); // Simulate work } catch (InterruptedException e) { Thread.currentThread().interrupt(); } System.out.println("Order " + message.getOrderId() + " processed");};
// Start multiple consumersqueue.startConsumer(processOrder, 1);queue.startConsumer(processOrder, 2);
// Producer sends messagesfor (int i = 0; i < 10; i++) { queue.produce(new Message(i, 100 + i)); Thread.sleep(100);}
Thread.sleep(5000); // Let consumers processqueue.stop();import { EventEmitter } from 'events';
interface Message { orderId: number; amount: number;}
class MessageQueue { private queue: Message[] = []; private maxSize: number; private consumers: Array<() => void> = []; private running: boolean = false; private eventEmitter: EventEmitter;
constructor(maxSize: number = 1000) { this.maxSize = maxSize; this.eventEmitter = new EventEmitter(); }
produce(message: Message): void { // Producer: Add message to queue if (this.queue.length >= this.maxSize) { console.log("Queue full! Message rejected."); return; }
this.queue.push(message); console.log(`Produced: ${JSON.stringify(message)}`); this.eventEmitter.emit('message'); }
consume(handler: (message: Message) => void): void { // Consumer: Process messages from queue const processMessages = () => { while (this.running && this.queue.length > 0) { const message = this.queue.shift(); if (message) { handler(message); } } };
this.eventEmitter.on('message', processMessages); processMessages(); // Process any existing messages }
startConsumer(handler: (message: Message) => void, consumerId: number): void { // Start a consumer console.log(`Consumer ${consumerId} started`); this.consume(handler); }
start(): void { this.running = true; }
stop(): void { this.running = false; }}
// Usageconst queue = new MessageQueue(100);queue.start();
// Message handlerfunction processOrder(message: Message): void { console.log(`Processing order: ${message.orderId}`); // Process order... setTimeout(() => { console.log(`Order ${message.orderId} processed`); }, 1000);}
// Start multiple consumersqueue.startConsumer(processOrder, 1);queue.startConsumer(processOrder, 2);
// Producer sends messagesfor (let i = 0; i < 10; i++) { queue.produce({ orderId: i, amount: 100 + i }); setTimeout(() => {}, 100);}
setTimeout(() => { queue.stop();}, 5000);#include <queue>#include <thread>#include <mutex>#include <condition_variable>#include <iostream>#include <functional>
struct Message { int orderId; int amount;};
class MessageQueue {private: std::queue<Message> queue; size_t maxSize; std::mutex mutex; std::condition_variable cv; bool running;
public: MessageQueue(size_t maxSize = 1000) : maxSize(maxSize), running(false) {}
void produce(const Message& message) { std::unique_lock<std::mutex> lock(mutex);
// Producer: Add message to queue if (queue.size() >= maxSize) { std::cout << "Queue full! Message rejected." << std::endl; return; }
queue.push(message); std::cout << "Produced: orderId=" << message.orderId << ", amount=" << message.amount << std::endl; cv.notify_one(); }
void consume(std::function<void(const Message&)> handler) { // Consumer: Process messages from queue while (running) { std::unique_lock<std::mutex> lock(mutex); cv.wait(lock, [this] { return !queue.empty() || !running; });
if (!running && queue.empty()) break;
if (!queue.empty()) { Message message = queue.front(); queue.pop(); lock.unlock();
handler(message); } } }
void startConsumer(std::function<void(const Message&)> handler, int consumerId) { // Start a consumer thread std::thread([this, handler, consumerId]() { std::cout << "Consumer " << consumerId << " started" << std::endl; this->consume(handler); }).detach(); }
void start() { running = true; }
void stop() { running = false; cv.notify_all(); }};
// Usagevoid processOrder(const Message& message) { std::cout << "Processing order: " << message.orderId << std::endl; // Process order... std::this_thread::sleep_for(std::chrono::seconds(1)); std::cout << "Order " << message.orderId << " processed" << std::endl;}
int main() { MessageQueue queue(100); queue.start();
// Start multiple consumers queue.startConsumer(processOrder, 1); queue.startConsumer(processOrder, 2);
// Producer sends messages for (int i = 0; i < 10; i++) { queue.produce({i, 100 + i}); std::this_thread::sleep_for(std::chrono::milliseconds(100)); }
std::this_thread::sleep_for(std::chrono::seconds(5)); queue.stop(); return 0;}using System;using System.Collections.Concurrent;using System.Threading;using System.Threading.Tasks;
public class Message { public int OrderId { get; set; } public int Amount { get; set; }}
public class MessageQueue { private readonly BlockingCollection<Message> queue; private bool running;
public MessageQueue(int maxSize = 1000) { queue = new BlockingCollection<Message>(maxSize); running = false; }
public void Produce(Message message) { // Producer: Add message to queue if (queue.TryAdd(message)) { Console.WriteLine($"Produced: OrderId={message.OrderId}, Amount={message.Amount}"); } else { Console.WriteLine("Queue full! Message rejected."); } }
public void Consume(Action<Message> handler) { // Consumer: Process messages from queue while (running || !queue.IsCompleted) { try { Message message = queue.Take(); handler(message); } catch (InvalidOperationException) { // Queue completed break; } } }
public void StartConsumer(Action<Message> handler, int consumerId) { // Start a consumer task Task.Run(() => { Console.WriteLine($"Consumer {consumerId} started"); Consume(handler); }); }
public void Start() { running = true; }
public void Stop() { running = false; queue.CompleteAdding(); }}
// Usageclass Program { static void ProcessOrder(Message message) { Console.WriteLine($"Processing order: {message.OrderId}"); // Process order... Thread.Sleep(1000); Console.WriteLine($"Order {message.OrderId} processed"); }
static void Main() { MessageQueue queue = new MessageQueue(100); queue.Start();
// Start multiple consumers queue.StartConsumer(ProcessOrder, 1); queue.StartConsumer(ProcessOrder, 2);
// Producer sends messages for (int i = 0; i < 10; i++) { queue.Produce(new Message { OrderId = i, Amount = 100 + i }); Thread.Sleep(100); }
Thread.Sleep(5000); queue.Stop(); }}One publisher, multiple subscribers. Each subscriber gets copy of message.
Characteristics:
Use cases:
from typing import List, Callable, Dict, Anyfrom threading import Lockimport threading
class Topic: """Pub-sub topic"""
def __init__(self, name: str): self.name = name self.subscribers: List[Callable] = [] self.lock = Lock()
def subscribe(self, handler: Callable): """Subscribe to topic""" with self.lock: self.subscribers.append(handler) print(f"Subscriber added to {self.name}. Total: {len(self.subscribers)}")
def unsubscribe(self, handler: Callable): """Unsubscribe from topic""" with self.lock: if handler in self.subscribers: self.subscribers.remove(handler)
def publish(self, message: Dict[str, Any]): """Publish message to all subscribers""" with self.lock: subscribers = self.subscribers.copy()
# Notify all subscribers (asynchronously) for subscriber in subscribers: try: # Run in separate thread to avoid blocking threading.Thread( target=subscriber, args=(message,), daemon=True ).start() except Exception as e: print(f"Error notifying subscriber: {e}")
class PubSubBroker: """Pub-sub message broker"""
def __init__(self): self.topics: Dict[str, Topic] = {} self.lock = Lock()
def get_topic(self, name: str) -> Topic: """Get or create topic""" with self.lock: if name not in self.topics: self.topics[name] = Topic(name) return self.topics[name]
def publish(self, topic_name: str, message: Dict[str, Any]): """Publish message to topic""" topic = self.get_topic(topic_name) topic.publish(message)
def subscribe(self, topic_name: str, handler: Callable): """Subscribe to topic""" topic = self.get_topic(topic_name) topic.subscribe(handler)
# Usagebroker = PubSubBroker()
# Subscribersdef email_handler(message): print(f"Email Service: Sending email for {message['event']}")
def sms_handler(message): print(f"SMS Service: Sending SMS for {message['event']}")
def analytics_handler(message): print(f"Analytics Service: Recording {message['event']}")
# Subscribe to 'user.created' topicbroker.subscribe('user.created', email_handler)broker.subscribe('user.created', sms_handler)broker.subscribe('user.created', analytics_handler)
# Publisher publishes eventbroker.publish('user.created', { 'event': 'user.created', 'user_id': 123,})import java.util.*;import java.util.concurrent.CopyOnWriteArrayList;import java.util.function.Consumer;
class Topic { private final String name; private final List<Consumer<Message>> subscribers = new CopyOnWriteArrayList<>();
public Topic(String name) { this.name = name; }
public void subscribe(Consumer<Message> handler) { subscribers.add(handler); System.out.println("Subscriber added to " + name + ". Total: " + subscribers.size()); }
public void unsubscribe(Consumer<Message> handler) { subscribers.remove(handler); }
public void publish(Message message) { // Notify all subscribers subscribers.forEach(subscriber -> { try { subscriber.accept(message); } catch (Exception e) { System.err.println("Error notifying subscriber: " + e.getMessage()); } }); }}
class PubSubBroker { private final Map<String, Topic> topics = new HashMap<>();
public synchronized Topic getTopic(String name) { return topics.computeIfAbsent(name, Topic::new); }
public void publish(String topicName, Message message) { Topic topic = getTopic(topicName); topic.publish(message); }
public void subscribe(String topicName, Consumer<Message> handler) { Topic topic = getTopic(topicName); topic.subscribe(handler); }}
// UsagePubSubBroker broker = new PubSubBroker();
// SubscribersConsumer<Message> emailHandler = message -> System.out.println("Email Service: Sending email for " + message.getEvent());
Consumer<Message> smsHandler = message -> System.out.println("SMS Service: Sending SMS for " + message.getEvent());
Consumer<Message> analyticsHandler = message -> System.out.println("Analytics Service: Recording " + message.getEvent());
// Subscribe to 'user.created' topicbroker.subscribe("user.created", emailHandler);broker.subscribe("user.created", smsHandler);broker.subscribe("user.created", analyticsHandler);
// Publisher publishes eventimport { EventEmitter } from 'events';
interface Message { event: string; userId?: number; email?: string;}
type MessageHandler = (message: Message) => void;
class Topic { private name: string; private subscribers: MessageHandler[] = []; private eventEmitter: EventEmitter;
constructor(name: string) { this.name = name; this.eventEmitter = new EventEmitter(); }
subscribe(handler: MessageHandler): void { // Subscribe to topic this.subscribers.push(handler); console.log(`Subscriber added to ${this.name}. Total: ${this.subscribers.length}`); }
unsubscribe(handler: MessageHandler): void { // Unsubscribe from topic const index = this.subscribers.indexOf(handler); if (index > -1) { this.subscribers.splice(index, 1); } }
publish(message: Message): void { // Publish message to all subscribers this.subscribers.forEach(subscriber => { try { // Run asynchronously to avoid blocking setImmediate(() => subscriber(message)); } catch (error) { console.error(`Error notifying subscriber: ${error}`); } }); }}
class PubSubBroker { private topics: Map<string, Topic> = new Map();
getTopic(name: string): Topic { // Get or create topic if (!this.topics.has(name)) { this.topics.set(name, new Topic(name)); } return this.topics.get(name)!; }
publish(topicName: string, message: Message): void { // Publish message to topic const topic = this.getTopic(topicName); topic.publish(message); }
subscribe(topicName: string, handler: MessageHandler): void { // Subscribe to topic const topic = this.getTopic(topicName); topic.subscribe(handler); }}
// Usageconst broker = new PubSubBroker();
// Subscribersconst emailHandler: MessageHandler = (message) => { console.log(`Email Service: Sending email for ${message.event}`);};
const smsHandler: MessageHandler = (message) => { console.log(`SMS Service: Sending SMS for ${message.event}`);};
const analyticsHandler: MessageHandler = (message) => { console.log(`Analytics Service: Recording ${message.event}`);};
// Subscribe to 'user.created' topicbroker.subscribe('user.created', emailHandler);broker.subscribe('user.created', smsHandler);broker.subscribe('user.created', analyticsHandler);
// Publisher publishes eventbroker.publish('user.created', { event: 'user.created', userId: 123,});#include <string>#include <vector>#include <map>#include <functional>#include <mutex>#include <iostream>#include <thread>
struct Message { std::string event; int userId; std::string email;};
class Topic {private: std::string name; std::vector<std::function<void(const Message&)>> subscribers; std::mutex mutex;
public: Topic(const std::string& name) : name(name) {}
void subscribe(std::function<void(const Message&)> handler) { // Subscribe to topic std::lock_guard<std::mutex> lock(mutex); subscribers.push_back(handler); std::cout << "Subscriber added to " << name << ". Total: " << subscribers.size() << std::endl; }
void unsubscribe(std::function<void(const Message&)> handler) { // Unsubscribe from topic std::lock_guard<std::mutex> lock(mutex); // Note: Simplified - in real implementation, need better handler comparison }
void publish(const Message& message) { // Publish message to all subscribers std::vector<std::function<void(const Message&)>> subs; { std::lock_guard<std::mutex> lock(mutex); subs = subscribers; }
// Notify all subscribers (asynchronously) for (auto& subscriber : subs) { try { std::thread([subscriber, message]() { subscriber(message); }).detach(); } catch (const std::exception& e) { std::cerr << "Error notifying subscriber: " << e.what() << std::endl; } } }};
class PubSubBroker {private: std::map<std::string, Topic> topics; std::mutex mutex;
public: Topic& getTopic(const std::string& name) { // Get or create topic std::lock_guard<std::mutex> lock(mutex); return topics.emplace(name, Topic(name)).first->second; }
void publish(const std::string& topicName, const Message& message) { // Publish message to topic Topic& topic = getTopic(topicName); topic.publish(message); }
void subscribe(const std::string& topicName, std::function<void(const Message&)> handler) { // Subscribe to topic Topic& topic = getTopic(topicName); topic.subscribe(handler); }};
// Usagevoid emailHandler(const Message& message) { std::cout << "Email Service: Sending email for " << message.event << std::endl;}
void smsHandler(const Message& message) { std::cout << "SMS Service: Sending SMS for " << message.event << std::endl;}
void analyticsHandler(const Message& message) { std::cout << "Analytics Service: Recording " << message.event << std::endl;}
int main() { PubSubBroker broker;
// Subscribe to 'user.created' topic broker.subscribe("user.created", emailHandler); broker.subscribe("user.created", smsHandler); broker.subscribe("user.created", analyticsHandler);
// Publisher publishes event
std::this_thread::sleep_for(std::chrono::milliseconds(100)); return 0;}using System;using System.Collections.Generic;using System.Threading.Tasks;
public class Message { public string Event { get; set; } public int? UserId { get; set; } public string Email { get; set; }}
public class Topic { private readonly string name; private readonly List<Action<Message>> subscribers; private readonly object lockObject = new object();
public Topic(string name) { this.name = name; this.subscribers = new List<Action<Message>>(); }
public void Subscribe(Action<Message> handler) { // Subscribe to topic lock (lockObject) { subscribers.Add(handler); Console.WriteLine($"Subscriber added to {name}. Total: {subscribers.Count}"); } }
public void Unsubscribe(Action<Message> handler) { // Unsubscribe from topic lock (lockObject) { subscribers.Remove(handler); } }
public void Publish(Message message) { // Publish message to all subscribers List<Action<Message>> subs; lock (lockObject) { subs = new List<Action<Message>>(subscribers); }
// Notify all subscribers (asynchronously) foreach (var subscriber in subs) { try { Task.Run(() => subscriber(message)); } catch (Exception e) { Console.Error.WriteLine($"Error notifying subscriber: {e.Message}"); } } }}
public class PubSubBroker { private readonly Dictionary<string, Topic> topics; private readonly object lockObject = new object();
public PubSubBroker() { topics = new Dictionary<string, Topic>(); }
private Topic GetTopic(string name) { // Get or create topic lock (lockObject) { if (!topics.ContainsKey(name)) { topics[name] = new Topic(name); } return topics[name]; } }
public void Publish(string topicName, Message message) { // Publish message to topic Topic topic = GetTopic(topicName); topic.Publish(message); }
public void Subscribe(string topicName, Action<Message> handler) { // Subscribe to topic Topic topic = GetTopic(topicName); topic.Subscribe(handler); }}
// Usageclass Program { static void EmailHandler(Message message) { Console.WriteLine($"Email Service: Sending email for {message.Event}"); }
static void SmsHandler(Message message) { Console.WriteLine($"SMS Service: Sending SMS for {message.Event}"); }
static void AnalyticsHandler(Message message) { Console.WriteLine($"Analytics Service: Recording {message.Event}"); }
static void Main() { PubSubBroker broker = new PubSubBroker();
// Subscribe to 'user.created' topic broker.Subscribe("user.created", EmailHandler); broker.Subscribe("user.created", SmsHandler); broker.Subscribe("user.created", AnalyticsHandler);
// Publisher publishes event broker.Publish("user.created", new Message { Event = "user.created", UserId = 123, });
System.Threading.Thread.Sleep(100); }}Message may be lost, but never duplicated.
Characteristics:
Use when: Non-critical messages, metrics, logs
Message delivered at least once, may have duplicates.
Characteristics:
Use when: Critical messages, order processing, payments
Message delivered exactly once. Requires deduplication.
Characteristics:
Use when: Financial transactions, critical operations
Ordering ensures messages processed in sequence.
Example: User account balance updates
Message 1: Balance = 100Message 2: Balance = 150 (add 50)Message 3: Balance = 120 (subtract 30)Correct order: 100 → 150 → 120
Wrong order: 100 → 120 → 150 = 150 (wrong!)
1. Per-Partition Ordering (Kafka)
2. Per-Queue Ordering (RabbitMQ)
3. Global Ordering
At the code level, message queues translate to producer/consumer classes, message handlers, and acknowledgment logic.
from abc import ABC, abstractmethodfrom typing import Dict, Any, Optionalimport json
class MessageHandler(ABC): """Base message handler interface"""
@abstractmethod def handle(self, message: Dict[str, Any]) -> bool: """ Handle message. Returns True if successful. Should be idempotent for at-least-once delivery. """ pass
@abstractmethod def can_handle(self, message_type: str) -> bool: """Check if handler can process this message type""" pass
class OrderProcessor(MessageHandler): """Process order messages"""
def __init__(self, order_service): self.order_service = order_service self.processed_ids = set() # For idempotency
def can_handle(self, message_type: str) -> bool: return message_type == 'order.created'
def handle(self, message: Dict[str, Any]) -> bool: order_id = message.get('order_id')
# Idempotency check if order_id in self.processed_ids: print(f"Order {order_id} already processed. Skipping.") return True # Already processed, consider success
try: # Process order self.order_service.process_order(order_id, message)
# Mark as processed self.processed_ids.add(order_id)
return True except Exception as e: print(f"Error processing order {order_id}: {e}") return False # Return False to trigger retry
class MessageConsumer: """Consumer that routes messages to handlers"""
def __init__(self, queue, handlers: List[MessageHandler]): self.queue = queue self.handlers = handlers
def consume(self): """Consume messages from queue""" while True: try: message_data = self.queue.get(timeout=1) message = json.loads(message_data)
# Find handler handler = self.find_handler(message.get('type'))
if handler: # Process message success = handler.handle(message)
if success: # Acknowledge message self.queue.task_done() else: # Return to queue for retry self.queue.put(message_data) else: print(f"No handler for message type: {message.get('type')}") self.queue.task_done()
except Exception as e: print(f"Error consuming message: {e}")
def find_handler(self, message_type: str) -> Optional[MessageHandler]: """Find handler for message type""" for handler in self.handlers: if handler.can_handle(message_type): return handler return Noneimport java.util.*;
interface MessageHandler { boolean handle(Message message); boolean canHandle(String messageType);}
class OrderProcessor implements MessageHandler { private final OrderService orderService; private final Set<String> processedIds = new HashSet<>();
public OrderProcessor(OrderService orderService) { this.orderService = orderService; }
@Override public boolean canHandle(String messageType) { return "order.created".equals(messageType); }
@Override public boolean handle(Message message) { String orderId = message.getOrderId();
// Idempotency check synchronized (processedIds) { if (processedIds.contains(orderId)) { System.out.println("Order " + orderId + " already processed. Skipping."); return true; // Already processed } }
try { // Process order orderService.processOrder(orderId, message);
// Mark as processed synchronized (processedIds) { processedIds.add(orderId); }
return true; } catch (Exception e) { System.err.println("Error processing order " + orderId + ": " + e.getMessage()); return false; // Return false to trigger retry } }}
class MessageConsumer { private final BlockingQueue<String> queue; private final List<MessageHandler> handlers;
public MessageConsumer(BlockingQueue<String> queue, List<MessageHandler> handlers) { this.queue = queue; this.handlers = handlers; }
public void consume() { while (true) { try { String messageData = queue.take(); Message message = parseMessage(messageData);
// Find handler MessageHandler handler = findHandler(message.getType());
if (handler != null) { // Process message boolean success = handler.handle(message);
if (success) { // Message processed successfully // (Acknowledgment handled by queue) } else { // Return to queue for retry queue.put(messageData); } } else { System.err.println("No handler for message type: " + message.getType()); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } catch (Exception e) { System.err.println("Error consuming message: " + e.getMessage()); } } }
private MessageHandler findHandler(String messageType) { return handlers.stream() .filter(h -> h.canHandle(messageType)) .findFirst() .orElse(null); }}interface Message { type: string; orderId?: string; [key: string]: any;}
interface MessageHandler { handle(message: Message): boolean; canHandle(messageType: string): boolean;}
class OrderProcessor implements MessageHandler { private orderService: any; private processedIds: Set<string> = new Set();
constructor(orderService: any) { this.orderService = orderService; }
canHandle(messageType: string): boolean { return messageType === 'order.created'; }
handle(message: Message): boolean { const orderId = message.orderId; if (!orderId) return false;
// Idempotency check if (this.processedIds.has(orderId)) { console.log(`Order ${orderId} already processed. Skipping.`); return true; // Already processed, consider success }
try { // Process order this.orderService.processOrder(orderId, message);
// Mark as processed this.processedIds.add(orderId);
return true; } catch (error) { console.error(`Error processing order ${orderId}:`, error); return false; // Return false to trigger retry } }}
class MessageConsumer { private queue: any; private handlers: MessageHandler[];
constructor(queue: any, handlers: MessageHandler[]) { this.queue = queue; this.handlers = handlers; }
async consume(): Promise<void> { while (true) { try { const messageData = await this.queue.get(1000); // timeout 1s const message: Message = JSON.parse(messageData);
// Find handler const handler = this.findHandler(message.type);
if (handler) { // Process message const success = handler.handle(message);
if (success) { // Acknowledge message this.queue.taskDone(); } else { // Return to queue for retry this.queue.put(messageData); } } else { console.log(`No handler for message type: ${message.type}`); this.queue.taskDone(); } } catch (error) { console.error(`Error consuming message:`, error); } } }
private findHandler(messageType: string): MessageHandler | null { return this.handlers.find(h => h.canHandle(messageType)) || null; }}#include <string>#include <unordered_set>#include <vector>#include <functional>#include <iostream>
struct Message { std::string type; std::string orderId; // Add other fields as needed};
class MessageHandler {public: virtual ~MessageHandler() = default; virtual bool handle(const Message& message) = 0; virtual bool canHandle(const std::string& messageType) = 0;};
class OrderProcessor : public MessageHandler {private: void* orderService; // Simplified - use proper service type std::unordered_set<std::string> processedIds; std::mutex mutex;
public: OrderProcessor(void* orderService) : orderService(orderService) {}
bool canHandle(const std::string& messageType) override { return messageType == "order.created"; }
bool handle(const Message& message) override { // Idempotency check { std::lock_guard<std::mutex> lock(mutex); if (processedIds.find(message.orderId) != processedIds.end()) { std::cout << "Order " << message.orderId << " already processed. Skipping." << std::endl; return true; // Already processed } }
try { // Process order // orderService->processOrder(message.orderId, message);
// Mark as processed { std::lock_guard<std::mutex> lock(mutex); processedIds.insert(message.orderId); }
return true; } catch (const std::exception& e) { std::cerr << "Error processing order " << message.orderId << ": " << e.what() << std::endl; return false; // Return false to trigger retry } }};
class MessageConsumer {private: // Queue type - simplified std::vector<std::string> queue; std::vector<MessageHandler*> handlers;
public: MessageConsumer(std::vector<MessageHandler*> handlers) : handlers(handlers) {}
void consume() { while (true) { try { if (queue.empty()) { std::this_thread::sleep_for(std::chrono::milliseconds(100)); continue; }
std::string messageData = queue.front(); queue.erase(queue.begin());
// Parse message (simplified) Message message; // Parse messageData into message
// Find handler MessageHandler* handler = findHandler(message.type);
if (handler) { // Process message bool success = handler->handle(message);
if (!success) { // Return to queue for retry queue.push_back(messageData); } } else { std::cerr << "No handler for message type: " << message.type << std::endl; } } catch (const std::exception& e) { std::cerr << "Error consuming message: " << e.what() << std::endl; } } }
private: MessageHandler* findHandler(const std::string& messageType) { for (auto* handler : handlers) { if (handler->canHandle(messageType)) { return handler; } } return nullptr; }};using System;using System.Collections.Generic;using System.Threading;
public interface IMessageHandler { bool Handle(Message message); bool CanHandle(string messageType);}
public class Message { public string Type { get; set; } public string OrderId { get; set; }}
public class OrderProcessor : IMessageHandler { private readonly object orderService; private readonly HashSet<string> processedIds; private readonly object lockObject = new object();
public OrderProcessor(object orderService) { this.orderService = orderService; this.processedIds = new HashSet<string>(); }
public bool CanHandle(string messageType) { return messageType == "order.created"; }
public bool Handle(Message message) { // Idempotency check lock (lockObject) { if (processedIds.Contains(message.OrderId)) { Console.WriteLine($"Order {message.OrderId} already processed. Skipping."); return true; // Already processed } }
try { // Process order // orderService.ProcessOrder(message.OrderId, message);
// Mark as processed lock (lockObject) { processedIds.Add(message.OrderId); }
return true; } catch (Exception e) { Console.Error.WriteLine($"Error processing order {message.OrderId}: {e.Message}"); return false; // Return false to trigger retry } }}
public class MessageConsumer { private readonly Queue<string> queue; private readonly List<IMessageHandler> handlers; private readonly object lockObject = new object();
public MessageConsumer(Queue<string> queue, List<IMessageHandler> handlers) { this.queue = queue; this.handlers = handlers; }
public void Consume() { while (true) { try { string messageData = null; lock (lockObject) { if (queue.Count > 0) { messageData = queue.Dequeue(); } }
if (messageData == null) { Thread.Sleep(100); continue; }
// Parse message (simplified) Message message = ParseMessage(messageData);
// Find handler IMessageHandler handler = FindHandler(message.Type);
if (handler != null) { // Process message bool success = handler.Handle(message);
if (!success) { // Return to queue for retry lock (lockObject) { queue.Enqueue(messageData); } } } else { Console.Error.WriteLine($"No handler for message type: {message.Type}"); } } catch (Exception e) { Console.Error.WriteLine($"Error consuming message: {e.Message}"); } } }
private IMessageHandler FindHandler(string messageType) { foreach (var handler in handlers) { if (handler.CanHandle(messageType)) { return handler; } } return null; }
private Message ParseMessage(string messageData) { // Simplified - use proper JSON parsing return new Message { Type = "order.created", OrderId = "123" }; }}Producer-Consumer
One producer, multiple consumers. Each message processed once. Perfect for task distribution.
Pub-Sub
One publisher, multiple subscribers. Each gets copy. Perfect for event broadcasting.
Delivery Guarantees
At-least-once most common. Requires idempotent consumers. Exactly-once is hardest but most reliable.
Ordering Matters
Message ordering critical for state changes. Per-partition ordering balances order and parallelism.