MFormations
Modern Backend Engineering

Chapitre 13

Chapitre 13 — Message Queues

Chapitre 13 — Message Queues

Cours — Message Queues

1. Introduction aux Message Queues

Pourquoi les files d'attente ?

  • Découplage : producteur et consommateur ne se connaissent pas
  • Buffer : lissage des pics de charge
  • Async : traitement différé des tâches longues
  • Scalabilité : multiples consommateurs
  • Résilience : messages persistés en cas de panne

Patterns de messagerie

PatternDescriptionCas d'usage
Point-to-Point1 producteur → 1 consommateurTask queue
Pub/Sub1 producteur → N consommateursNotifications
Request/ReplyMessage avec réponseRPC
Competing ConsumersN consommateurs sur 1 queueLoad balancing

2. RabbitMQ (AMQP)

Concepts AMQP

Publisher → Exchange → Binding → Queue → Consumer

Types d'exchanges

TypeRoutage
Directrouting_key exacte
Topicrouting_key avec pattern (topic.#, topic.*)
Fanoutbroadcast à toutes les queues
Headersbasé sur les en-têtes

Exemple Node.js

import amqp from 'amqplib';

async function setup() {
  const conn = await amqp.connect('amqp://localhost');
  const channel = await conn.createChannel();

  // Declare exchange
  await channel.assertExchange('orders', 'topic', { durable: true });

  // Declare queue
  await channel.assertQueue('order.created', { durable: true });

  // Bind queue to exchange
  await channel.bindQueue('order.created', 'orders', 'order.created.*');

  // Publish
  channel.publish('orders', 'order.created.new', Buffer.from(JSON.stringify({
    orderId: 123,
    userId: 456,
    amount: 99.99
  })), {
    persistent: true,
    contentType: 'application/json'
  });

  // Consume
  channel.consume('order.created', (msg) => {
    if (msg) {
      const data = JSON.parse(msg.content.toString());
      console.log('Processing order:', data.orderId);
      channel.ack(msg);
    }
  });
}

Dead Letter Queue (DLQ)

await channel.assertQueue('orders.dlq', { durable: true });
await channel.assertQueue('orders.main', {
  durable: true,
  deadLetterExchange: '',
  deadLetterRoutingKey: 'orders.dlq',
  maxLength: 1000,
  messageTtl: 60000, // 1 minute
});

3. Kafka

Architecture

Topic (partition 0) ──▶ Consumer Group A
Topic (partition 1) ──▶ Consumer Group A
Topic (partition 2) ──▶ Consumer Group B

Concepts clés

  • Topic : catégorie de messages
  • Partition : division d'un topic (ordre garanti dans une partition)
  • Offset : position du message dans la partition
  • Consumer Group : groupe de consommateurs (load balancing)
  • Broker : serveur Kafka
  • Replication : copies des partitions (HA)

Producer

import { Kafka } from 'kafkajs';

const kafka = new Kafka({
  clientId: 'order-service',
  brokers: ['localhost:9092'],
});

const producer = kafka.producer({
  allowAutoTopicCreation: true,
  transactionTimeout: 30000,
});

await producer.connect();

await producer.send({
  topic: 'orders',
  messages: [
    {
      key: 'order-123',
      value: JSON.stringify({ orderId: 123, userId: 456 }),
      headers: { 'event-type': 'order.created' },
    },
  ],
});

await producer.disconnect();

Consumer

const consumer = kafka.consumer({ groupId: 'order-processor' });

await consumer.connect();
await consumer.subscribe({ topic: 'orders', fromBeginning: true });

await consumer.run({
  eachMessage: async ({ topic, partition, message }) => {
    console.log({
      topic,
      partition,
      offset: message.offset,
      key: message.key.toString(),
      value: message.value.toString(),
    });

    // Traitement
    await processOrder(JSON.parse(message.value.toString()));

    // Auto-commit (ou manuel avec enableAutoCommit: false)
  },
  autoCommitInterval: 5000, // commit toutes les 5s
});

Kafka Streams (traitement temps réel)

const { StreamsConfig } = require('kafkajs');

// Stream processing avec kafka.js ou utilisez kafka-streams
// Exemple conceptuel :
// 1. Lire depuis un topic source
// 2. Transformer les données
// 3. Écrire dans un topic destination

4. Redis Pub/Sub

import Redis from 'ioredis';

const publisher = new Redis();
const subscriber = new Redis();

// Subscriber
subscriber.subscribe('notifications', (err, count) => {
  console.log(`Subscribed to ${count} channels`);
});

subscriber.on('message', (channel, message) => {
  console.log(`Received ${message} from ${channel}`);
});

// Publisher
setInterval(() => {
  publisher.publish('notifications', JSON.stringify({
    type: 'alert',
    message: 'Server CPU > 80%',
    timestamp: Date.now(),
  }));
}, 5000);

Redis Streams (persistant)

// Add to stream
await publisher.xadd('mystream', '*', 'key', 'value', 'temperature', 25);

// Read from stream
const results = await subscriber.xread(
  'BLOCK', 5000,
  'STREAMS', 'mystream', '$'
);

5. Architecture Event-Driven

Event-Driven Microservices

[Order Service] ──▶ order.created ──▶ [Payment Service]
                      │
                      ├──▶ [Inventory Service]
                      │
                      └──▶ [Notification Service]

Avantages

  • Découplage fort entre services
  • Évolutivité indépendante
  • Résilience (un service down n'impacte pas les autres)
  • Traçabilité (event sourcing)

Inconvénients

  • Complexité (eventual consistency)
  • Debug plus difficile
  • Nécessite de la monitoring

6. Dead Letter Queues (DLQ)

Quand utiliser une DLQ ?

  • Message mal formaté
  • Échec de traitement après N tentatives
  • Timeout de traitement
  • Message en dehors des limites (taille, TTL)
async function consumeWithRetry(channel, queue, dlq, maxRetries = 3) {
  await channel.assertQueue(queue, {
    deadLetterExchange: '',
    deadLetterRoutingKey: dlq,
  });

  channel.consume(queue, async (msg) => {
    try {
      await processMessage(msg);
      channel.ack(msg);
    } catch (error) {
      const retryCount = (msg.properties.headers['x-retry-count'] || 0) + 1;

      if (retryCount <= maxRetries) {
        // Retry with delay
        channel.publish(
          '',
          queue,
          msg.content,
          {
            headers: { 'x-retry-count': retryCount, 'x-delay': retryCount * 5000 },
            persistent: true,
          }
        );
        channel.ack(msg);
      } else {
        // Send to DLQ
        channel.reject(msg, false); // false = ne pas requeue
      }
    }
  });
}

7. Message Ordering et Partitioning

Garantir l'ordre

  • Kafka : ordre garanti dans une partition
  • RabbitMQ : ordre garanti dans une queue (avec 1 consommateur)

Partition key

// Kafka : même clé = même partition = ordre garanti
await producer.send({
  topic: 'orders',
  messages: [{ key: `user-${userId}`, value: payload }],
});

Problèmes d'ordre

  • Plusieurs consommateurs sur une queue = ordre non garanti
  • Solution : router par clé de partition

8. Idempotent Consumers

Principe

Un message peut être traité plusieurs fois → le résultat doit être le même.

Comment implémenter l'idempotence

async function processEvent(event) {
  // 1. Vérifier si déjà traité
  const alreadyProcessed = await redis.exists(`processed:${event.id}`);
  if (alreadyProcessed) {
    return; // Déjà traité
  }

  // 2. Marquer comme en cours
  await redis.set(`processing:${event.id}`, Date.now(), 'EX', 30);

  // 3. Traiter (transactionnel)
  await db.transaction(async (tx) => {
    // Insérer le résultat
    await tx.query('INSERT INTO events (id, data) VALUES ($1, $2)', [event.id, event.data]);
    // Marquer comme traité
    await redis.set(`processed:${event.id}`, '1', 'EX', 86400);
  });

  // 4. Nettoyer
  await redis.del(`processing:${event.id}`);
}

At-Least-Once vs Exactly-Once

GuaranteeDescriptionCoût
At-most-oncePeut perdre des messagesFaible
At-least-oncePeut dupliquerMoyen
Exactly-oncePas de perte, pas de duplicationÉlevé

9. Comparaison des Message Brokers

CritèreRabbitMQKafkaRedis Pub/Sub
ModèleQueue + ExchangeLog (partition)Pub/Sub
PersistanceOuiOuiNon (ou Streams)
Ordre1 queue1 partitionNon garanti
Throughput10k msg/s1M+ msg/s100k msg/s
RétentionAck-basedTime-basedEn mémoire
Cas d'usageTasks, RPCEvents, logs, streamsReal-time, cache
ComplexitéMoyenneÉlevéeFaible

10. Patterns Avancés

Event Sourcing

Stockage de l'état comme séquence d'événements.

  • Avantage : audit trail complet, rejouer les events
  • Inconvénient : complexité, stockage

CQRS (Command Query Responsibility Segregation)

Séparation des commandes (écritures) et requêtes (lectures).

  • Commande : validation, mise à jour, publication d'event
  • Requête : lecture optimisée (matérialisation)

Saga Pattern

Pour les transactions distribuées.

  • Choreography : chaque service publie/reagit aux events
  • Orchestration : un orchestrateur coordonne

Outbox Pattern

Pour garantir la fiabilité des événements.

1. Écrire dans la BDD + outbox (même transaction)
2. Processus lit l'outbox et publie les events
3. Supprime de l'outbox après publication