Exchanges Route
Exchanges route messages to queues based on routing keys. Direct, topic, fanout, headers.
RabbitMQ is a traditional message broker that implements the AMQP (Advanced Message Queuing Protocol) standard. Unlike Kafka’s log-based approach, RabbitMQ uses a more traditional broker model where messages are delivered to consumers and then removed from queues (unless configured otherwise).
RabbitMQ’s architecture revolves around four key components:
This design provides flexible routing capabilities that Kafka doesn’t offer natively. RabbitMQ excels at scenarios requiring complex routing logic, task distribution, and point-to-point messaging.
The Flow: Producer → Exchange (with routing key) → Binding (matches routing key) → Queue → Consumer
This design allows for flexible routing - you can route the same message to multiple queues, filter messages based on patterns, or broadcast to all queues, depending on the exchange type.
Routes to queue with matching routing key.
Use case: Point-to-point messaging, task queues
Routes based on pattern matching (wildcards).
Patterns:
* - Matches one word# - Matches zero or more wordsUse case: Categorized messages, event routing
Broadcasts to all bound queues (ignores routing key).
Use case: Pub-sub, notifications, cache invalidation
Routes based on message headers (ignores routing key).
Use case: Complex routing logic
import pikaimport json
class RabbitMQProducer: """RabbitMQ producer"""
def __init__(self, host='localhost'): self.connection = pika.BlockingConnection( pika.ConnectionParameters(host=host) ) self.channel = self.connection.channel()
def setup_exchange(self, exchange_name: str, exchange_type: str = 'direct'): """Declare exchange""" self.channel.exchange_declare( exchange=exchange_name, exchange_type=exchange_type, durable=True # Survive broker restart )
def publish(self, exchange: str, routing_key: str, message: dict): """Publish message""" self.channel.basic_publish( exchange=exchange, routing_key=routing_key, body=json.dumps(message), properties=pika.BasicProperties( delivery_mode=2, # Make message persistent content_type='application/json' ) ) print(f"Published to {exchange} with key {routing_key}")
def close(self): """Close connection""" self.connection.close()
# Usageproducer = RabbitMQProducer()producer.setup_exchange('orders', 'direct')producer.publish('orders', 'order.created', { 'order_id': 123, 'user_id': 456, 'amount': 99.99})producer.close()import com.rabbitmq.client.*;
public class RabbitMQProducer { private final Connection connection; private final Channel channel;
public RabbitMQProducer(String host) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(host); this.connection = factory.newConnection(); this.channel = connection.createChannel(); }
public void setupExchange(String exchangeName, String exchangeType) throws Exception { // Declare exchange channel.exchangeDeclare(exchangeName, exchangeType, true); // Durable }
public void publish(String exchange, String routingKey, String message) throws Exception { // Publish message channel.basicPublish( exchange, routingKey, MessageProperties.PERSISTENT_TEXT_PLAIN, // Make persistent message.getBytes() ); System.out.println("Published to " + exchange + " with key " + routingKey); }
public void close() throws Exception { channel.close(); connection.close(); }}
// UsageRabbitMQProducer producer = new RabbitMQProducer("localhost");producer.setupExchange("orders", "direct");producer.publish("orders", "order.created", "{\"order_id\":123,\"user_id\":456,\"amount\":99.99}");producer.close();import amqp from 'amqplib';
class RabbitMQProducer { private connection: amqp.Connection | null = null; private channel: amqp.Channel | null = null;
async connect(host: string = 'amqp://localhost'): Promise<void> { // Connect to RabbitMQ this.connection = await amqp.connect(host); this.channel = await this.connection.createChannel(); }
async setupExchange(exchangeName: string, exchangeType: string = 'direct'): Promise<void> { // Declare exchange if (!this.channel) throw new Error('Not connected');
await this.channel.assertExchange(exchangeName, exchangeType, { durable: true // Survive broker restart }); }
async publish(exchange: string, routingKey: string, message: object): Promise<void> { // Publish message if (!this.channel) throw new Error('Not connected');
await this.channel.publish( exchange, routingKey, Buffer.from(JSON.stringify(message)), { persistent: true, // Make message persistent contentType: 'application/json' } ); console.log(`Published to ${exchange} with key ${routingKey}`); }
async close(): Promise<void> { // Close connection if (this.channel) await this.channel.close(); if (this.connection) await this.connection.close(); }}
// Usageconst producer = new RabbitMQProducer();await producer.connect();await producer.setupExchange('orders', 'direct');await producer.publish('orders', 'order.created', { order_id: 123, user_id: 456, amount: 99.99});await producer.close();#include <amqpcpp.h>#include <amqpcpp/libev.h>#include <iostream>#include <nlohmann/json.hpp>
class RabbitMQProducer {private: AMQP::TcpConnection* connection; AMQP::TcpChannel* channel;
public: RabbitMQProducer(const std::string& host = "localhost") { // Create connection AMQP::TcpConnectionHandler handler; connection = new AMQP::TcpConnection(&handler, AMQP::Address(host)); channel = new AMQP::TcpChannel(connection);
// Setup exchange channel->declareExchange("orders", AMQP::direct) .onSuccess([]() { std::cout << "Exchange declared" << std::endl; }); }
void publish(const std::string& exchange, const std::string& routingKey, const nlohmann::json& message) { // Publish message std::string body = message.dump();
channel->publish(exchange, routingKey, body) .onSuccess([]() { std::cout << "Message published" << std::endl; }); }
~RabbitMQProducer() { delete channel; delete connection; }};
// Usageint main() { RabbitMQProducer producer("localhost");
nlohmann::json message = { {"order_id", 123}, {"user_id", 456}, {"amount", 99.99} };
producer.publish("orders", "order.created", message); return 0;}using RabbitMQ.Client;using System;using System.Text;using System.Text.Json;
public class RabbitMQProducer : IDisposable { private readonly IConnection connection; private readonly IModel channel;
public RabbitMQProducer(string hostName = "localhost") { var factory = new ConnectionFactory { HostName = hostName }; connection = factory.CreateConnection(); channel = connection.CreateModel(); }
public void SetupExchange(string exchangeName, string exchangeType = "direct") { // Declare exchange channel.ExchangeDeclare( exchange: exchangeName, type: exchangeType, durable: true // Survive broker restart ); }
public void Publish(string exchange, string routingKey, object message) { // Publish message var body = Encoding.UTF8.GetBytes(JsonSerializer.Serialize(message));
var properties = channel.CreateBasicProperties(); properties.Persistent = true; // Make message persistent properties.ContentType = "application/json";
channel.BasicPublish( exchange: exchange, routingKey: routingKey, basicProperties: properties, body: body );
Console.WriteLine($"Published to {exchange} with key {routingKey}"); }
public void Dispose() { channel?.Close(); connection?.Close(); }}
// Usageusing var producer = new RabbitMQProducer();producer.SetupExchange("orders", "direct");producer.Publish("orders", "order.created", new { order_id = 123, user_id = 456, amount = 99.99});import pikaimport json
class RabbitMQConsumer: """RabbitMQ consumer"""
def __init__(self, host='localhost'): self.connection = pika.BlockingConnection( pika.ConnectionParameters(host=host) ) self.channel = self.connection.channel()
def setup_queue(self, queue_name: str, durable: bool = True): """Declare queue""" self.channel.queue_declare( queue=queue_name, durable=durable # Survive broker restart )
def bind_queue(self, queue: str, exchange: str, routing_key: str): """Bind queue to exchange""" self.channel.queue_bind( queue=queue, exchange=exchange, routing_key=routing_key )
def consume(self, queue: str, handler, auto_ack: bool = False): """Consume messages""" def callback(ch, method, properties, body): try: message = json.loads(body) # Process message handler(message)
# Acknowledge message if not auto_ack: ch.basic_ack(delivery_tag=method.delivery_tag) except Exception as e: print(f"Error processing message: {e}") # Reject and requeue if not auto_ack: ch.basic_nack( delivery_tag=method.delivery_tag, requeue=True )
# Set prefetch (how many unacked messages per consumer) self.channel.basic_qos(prefetch_count=1)
# Start consuming self.channel.basic_consume( queue=queue, on_message_callback=callback, auto_ack=auto_ack )
print(f"Consuming from {queue}...") self.channel.start_consuming()
def close(self): """Close connection""" self.connection.close()
# Usagedef handle_order(message): print(f"Processing order: {message['order_id']}") # Process order...
consumer = RabbitMQConsumer()consumer.setup_queue('order-processor', durable=True)consumer.bind_queue('order-processor', 'orders', 'order.created')consumer.consume('order-processor', handle_order, auto_ack=False)import com.rabbitmq.client.*;
public class RabbitMQConsumer { private final Connection connection; private final Channel channel;
public RabbitMQConsumer(String host) throws Exception { ConnectionFactory factory = new ConnectionFactory(); factory.setHost(host); this.connection = factory.newConnection(); this.channel = connection.createChannel(); }
public void setupQueue(String queueName, boolean durable) throws Exception { // Declare queue channel.queueDeclare(queueName, durable, false, false, null); }
public void bindQueue(String queue, String exchange, String routingKey) throws Exception { // Bind queue to exchange channel.queueBind(queue, exchange, routingKey); }
public void consume(String queue, java.util.function.Consumer<String> handler, boolean autoAck) throws Exception { // Set prefetch channel.basicQos(1);
// Consumer callback DeliverCallback deliverCallback = (consumerTag, delivery) -> { try { String message = new String(delivery.getBody(), "UTF-8"); // Process message handler.accept(message);
// Acknowledge if (!autoAck) { channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false); } } catch (Exception e) { System.err.println("Error processing message: " + e.getMessage()); // Reject and requeue if (!autoAck) { channel.basicNack( delivery.getEnvelope().getDeliveryTag(), false, true // Requeue ); } } };
// Start consuming channel.basicConsume(queue, autoAck, deliverCallback, consumerTag -> {}); System.out.println("Consuming from " + queue + "..."); }
public void close() throws Exception { channel.close(); connection.close(); }}
// UsageRabbitMQConsumer consumer = new RabbitMQConsumer("localhost");consumer.setupQueue("order-processor", true);consumer.bindQueue("order-processor", "orders", "order.created");consumer.consume("order-processor", message -> { System.out.println("Processing: " + message);}, false);import amqp from 'amqplib';
class RabbitMQConsumer { private connection: amqp.Connection | null = null; private channel: amqp.Channel | null = null;
async connect(host: string = 'amqp://localhost'): Promise<void> { // Connect to RabbitMQ this.connection = await amqp.connect(host); this.channel = await this.connection.createChannel(); }
async setupQueue(queueName: string, durable: boolean = true): Promise<void> { // Declare queue if (!this.channel) throw new Error('Not connected');
await this.channel.assertQueue(queueName, { durable: durable // Survive broker restart }); }
async bindQueue(queue: string, exchange: string, routingKey: string): Promise<void> { // Bind queue to exchange if (!this.channel) throw new Error('Not connected');
await this.channel.bindQueue(queue, exchange, routingKey); }
async consume( queue: string, handler: (message: any) => Promise<void>, autoAck: boolean = false ): Promise<void> { // Consume messages if (!this.channel) throw new Error('Not connected');
// Set prefetch (how many unacked messages per consumer) await this.channel.prefetch(1);
await this.channel.consume(queue, async (msg) => { if (!msg) return;
try { const message = JSON.parse(msg.content.toString()); // Process message await handler(message);
// Acknowledge message if (!autoAck) { this.channel!.ack(msg); } } catch (error) { console.error(`Error processing message: ${error}`); // Reject and requeue if (!autoAck) { this.channel!.nack(msg, false, true); } } }, { noAck: autoAck });
console.log(`Consuming from ${queue}...`); }
async close(): Promise<void> { // Close connection if (this.channel) await this.channel.close(); if (this.connection) await this.connection.close(); }}
// Usageasync function handleOrder(message: any) { console.log(`Processing order: ${message.order_id}`); // Process order...}
const consumer = new RabbitMQConsumer();await consumer.connect();await consumer.setupQueue('order-processor', true);await consumer.bindQueue('order-processor', 'orders', 'order.created');await consumer.consume('order-processor', handleOrder, false);#include <amqpcpp.h>#include <amqpcpp/libev.h>#include <iostream>#include <nlohmann/json.hpp>
class RabbitMQConsumer {private: AMQP::TcpConnection* connection; AMQP::TcpChannel* channel;
public: RabbitMQConsumer(const std::string& host = "localhost") { // Create connection AMQP::TcpConnectionHandler handler; connection = new AMQP::TcpConnection(&handler, AMQP::Address(host)); channel = new AMQP::TcpChannel(connection); }
void setupQueue(const std::string& queueName, bool durable = true) { // Declare queue channel->declareQueue(queueName, AMQP::durable) .onSuccess([]() { std::cout << "Queue declared" << std::endl; }); }
void bindQueue(const std::string& queue, const std::string& exchange, const std::string& routingKey) { // Bind queue to exchange channel->bindQueue(exchange, queue, routingKey) .onSuccess([]() { std::cout << "Queue bound" << std::endl; }); }
void consume(const std::string& queue, std::function<void(const nlohmann::json&)> handler) { // Consume messages channel->setQos(1); // Prefetch count
channel->consume(queue) .onReceived([handler](const AMQP::Message& msg, uint64_t deliveryTag, bool redelivered) { try { nlohmann::json message = nlohmann::json::parse(msg.message()); // Process message handler(message);
// Acknowledge // channel->ack(deliveryTag); } catch (const std::exception& e) { std::cerr << "Error processing message: " << e.what() << std::endl; // Reject and requeue // channel->reject(deliveryTag, true); } });
std::cout << "Consuming from " << queue << "..." << std::endl; }
~RabbitMQConsumer() { delete channel; delete connection; }};
// Usagevoid handleOrder(const nlohmann::json& message) { std::cout << "Processing order: " << message["order_id"] << std::endl; // Process order...}
int main() { RabbitMQConsumer consumer("localhost"); consumer.setupQueue("order-processor", true); consumer.bindQueue("order-processor", "orders", "order.created"); consumer.consume("order-processor", handleOrder); return 0;}using RabbitMQ.Client;using RabbitMQ.Client.Events;using System;using System.Text;using System.Text.Json;
public class RabbitMQConsumer : IDisposable { private readonly IConnection connection; private readonly IModel channel;
public RabbitMQConsumer(string hostName = "localhost") { var factory = new ConnectionFactory { HostName = hostName }; connection = factory.CreateConnection(); channel = connection.CreateModel(); }
public void SetupQueue(string queueName, bool durable = true) { // Declare queue channel.QueueDeclare( queue: queueName, durable: durable, // Survive broker restart exclusive: false, autoDelete: false, arguments: null ); }
public void BindQueue(string queue, string exchange, string routingKey) { // Bind queue to exchange channel.QueueBind(queue, exchange, routingKey); }
public void Consume(string queue, Action<object> handler, bool autoAck = false) { // Set prefetch (how many unacked messages per consumer) channel.BasicQos(0, 1, false);
var consumer = new EventingBasicConsumer(channel); consumer.Received += (model, ea) => { try { var body = ea.Body.ToArray(); var message = JsonSerializer.Deserialize<object>( Encoding.UTF8.GetString(body) );
// Process message handler(message);
// Acknowledge message if (!autoAck) { channel.BasicAck(ea.DeliveryTag, false); } } catch (Exception e) { Console.Error.WriteLine($"Error processing message: {e.Message}"); // Reject and requeue if (!autoAck) { channel.BasicNack(ea.DeliveryTag, false, true); } } };
channel.BasicConsume( queue: queue, autoAck: autoAck, consumer: consumer );
Console.WriteLine($"Consuming from {queue}..."); }
public void Dispose() { channel?.Close(); connection?.Close(); }}
// Usagevoid HandleOrder(object message) { Console.WriteLine($"Processing order: {message}"); // Process order...}
using var consumer = new RabbitMQConsumer();consumer.SetupQueue("order-processor", true);consumer.BindQueue("order-processor", "orders", "order.created");consumer.Consume("order-processor", HandleOrder, false);Critical for reliable message processing.
# Message removed immediately when deliveredchannel.basic_consume(queue='orders', on_message_callback=callback, auto_ack=True)Problem: If consumer crashes, message lost!
# Message removed only after ackdef callback(ch, method, properties, body): process_message(body) ch.basic_ack(delivery_tag=method.delivery_tag) # Acknowledge
channel.basic_consume(queue='orders', on_message_callback=callback, auto_ack=False)Benefits:
| Feature | RabbitMQ | Kafka |
|---|---|---|
| Model | Traditional broker | Streaming platform |
| Message Retention | Removed after consumption | Retained (configurable) |
| Routing | Flexible (exchanges) | Simple (topics/partitions) |
| Ordering | Per queue | Per partition |
| Throughput | Good | Excellent |
| Use Case | Task queues, RPC | Event streaming, logs |
Choose RabbitMQ when:
Choose Kafka when:
Exchanges Route
Exchanges route messages to queues based on routing keys. Direct, topic, fanout, headers.
Manual Ack
Use manual acknowledgment for reliability. Auto-ack removes messages immediately (risky).
Durability
Make queues/exchanges/messages durable to survive broker restart. Critical for production.
Flexible Routing
RabbitMQ’s flexible routing (exchanges) makes it great for complex routing scenarios.