Difficulty: Advanced
When would you use a message queue? Compare RabbitMQ and Kafka and explain common messaging patterns.
Message queues decouple producers from consumers, enabling asynchronous processing, load leveling, and fault tolerance.
When to use: - Async processing: sending emails, generating reports, processing images - Decoupling services: order service publishes event, inventory and notification services consume independently - Load leveling: absorb traffic spikes by buffering requests - Reliable delivery: guaranteed processing even if consumers are temporarily down
RabbitMQ vs Kafka: - RabbitMQ: Traditional message broker. Push-based. Messages deleted after consumption. Best for task queues, RPC. - Kafka: Distributed event log. Pull-based. Messages retained for configurable period. Best for event streaming, analytics, high throughput.
Patterns: Point-to-point (work queue), Pub/Sub (fan-out), Request/Reply (RPC).
// Producer: send email task to queue
const amqp = require('amqplib');
async function publishEmailTask(emailData) {
const conn = await amqp.connect('amqp://localhost');
const channel = await conn.createChannel();
await channel.assertQueue('email_queue', {
durable: true // Survives broker restart
});
channel.sendToQueue('email_queue',
Buffer.from(JSON.stringify(emailData)),
{ persistent: true } // Survives broker restart
);
}
// Consumer: process email tasks
async function startEmailWorker() {
const conn = await amqp.connect('amqp://localhost');
const channel = await conn.createChannel();
await channel.assertQueue('email_queue', { durable: true });
channel.prefetch(1); // Process one at a time
channel.consume('email_queue', async (msg) => {
const data = JSON.parse(msg.content.toString());
try {
await sendEmail(data);
channel.ack(msg); // Success - remove from queue
} catch (error) {
if (msg.fields.redelivered) {
channel.nack(msg, false, false); // Send to DLQ
} else {
channel.nack(msg, false, true); // Retry once
}
}
});
}
Work queues distribute tasks among multiple workers. Durable queues and persistent messages survive broker restarts. Manual acknowledgment ensures no messages are lost.
// Kafka producer: publish order events
const { Kafka } = require('kafkajs');
const kafka = new Kafka({ brokers: ['localhost:9092'] });
const producer = kafka.producer();
async function publishOrderEvent(order) {
await producer.send({
topic: 'order-events',
messages: [{
key: order.userId, // Same user -> same partition -> ordered
value: JSON.stringify({
type: 'ORDER_CREATED',
data: order,
timestamp: Date.now()
})
}]
});
}
// Kafka consumer: multiple consumer groups
// Group 1: Inventory service
const inventoryConsumer = kafka.consumer({ groupId: 'inventory-service' });
await inventoryConsumer.subscribe({ topic: 'order-events' });
await inventoryConsumer.run({
eachMessage: async ({ message }) => {
const event = JSON.parse(message.value);
if (event.type === 'ORDER_CREATED') {
await reserveInventory(event.data);
}
}
});
// Group 2: Notification service (reads SAME messages independently)
const notifConsumer = kafka.consumer({ groupId: 'notification-service' });
await notifConsumer.subscribe({ topic: 'order-events' });
await notifConsumer.run({
eachMessage: async ({ message }) => {
const event = JSON.parse(message.value);
await sendOrderConfirmation(event.data);
}
});
// Key difference from RabbitMQ:
// - Messages are NOT deleted after consumption
// - Multiple consumer groups read independently
// - Messages ordered within a partition (by key)
Kafka retains events as a log, allowing multiple consumer groups to process the same events independently. Partitioning by key ensures ordering for related events.
// Dead Letter Queue (DLQ) handles messages that fail
// repeatedly, preventing them from blocking the queue
// RabbitMQ DLQ setup
await channel.assertQueue('email_dlq', { durable: true });
await channel.assertQueue('email_queue', {
durable: true,
arguments: {
'x-dead-letter-exchange': '',
'x-dead-letter-routing-key': 'email_dlq',
'x-message-ttl': 300000 // 5 min TTL
}
});
// DLQ processor: alert and store for manual review
channel.consume('email_dlq', async (msg) => {
const data = JSON.parse(msg.content.toString());
const deathInfo = msg.properties.headers['x-death'];
await db.failedMessages.create({
queue: 'email_queue',
payload: data,
reason: deathInfo?.[0]?.reason || 'unknown',
failedAt: new Date(),
retryCount: deathInfo?.[0]?.count || 0
});
await alertOps(`Email delivery permanently failed for ${data.to}`);
channel.ack(msg);
});
// Manual retry from DLQ
async function retryFromDLQ(messageId) {
const failed = await db.failedMessages.findById(messageId);
await publishEmailTask(failed.payload);
await db.failedMessages.markRetried(messageId);
}
Dead letter queues capture messages that cannot be processed after retries. This prevents poison messages from blocking the main queue and provides a path for manual investigation and retry.
Message Queue, Pub/Sub, Kafka, RabbitMQ, Dead Letter Queue, Event Sourcing