Message Queues (RabbitMQ, Kafka)

Difficulty: Advanced

Question

When would you use a message queue? Compare RabbitMQ and Kafka and explain common messaging patterns.

Answer

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).

Code examples

RabbitMQ Work Queue Pattern

// 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 Event Streaming

// 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 Pattern

// 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.

Key points

Concepts covered

Message Queue, Pub/Sub, Kafka, RabbitMQ, Dead Letter Queue, Event Sourcing