🎯 Patterns de Messagerie

Introduction

Les patterns de messagerie sont des solutions éprouvées à des problèmes récurrents dans les systèmes de communication asynchrone. RabbitMQ supporte nativement plusieurs de ces patterns grâce à son modèle AMQP flexible.

📬 Work Queues (Task Queues)

Concept

Distribuer des tâches entre plusieurs workers pour paralléliser le traitement :

# Producer - Envoie des tâches
import pika
import json

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# Déclarer une queue durable
channel.queue_declare(queue='task_queue', durable=True)

def send_task(task_data):
    message = json.dumps(task_data)
    
    channel.basic_publish(
        exchange='',
        routing_key='task_queue',
        body=message,
        properties=pika.BasicProperties(
            delivery_mode=2,  # Message persistant
        )
    )
    print(f"Task sent: {task_data}")

# Envoyer plusieurs tâches
for i in range(10):
    send_task({'task_id': i, 'type': 'process_image', 'url': f'image_{i}.jpg'})
# Consumer - Worker qui traite les tâches
def worker():
    connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
    channel = connection.channel()
    
    channel.queue_declare(queue='task_queue', durable=True)
    
    # Fair dispatch - un message à la fois par worker
    channel.basic_qos(prefetch_count=1)
    
    def process_task(ch, method, properties, body):
        task = json.loads(body)
        print(f"Processing task: {task}")
        
        # Simulation de traitement
        time.sleep(task.get('complexity', 1))
        
        print(f"Task completed: {task['task_id']}")
        ch.basic_ack(delivery_tag=method.delivery_tag)
    
    channel.basic_consume(queue='task_queue', on_message_callback=process_task)
    
    print('Worker started. Waiting for tasks...')
    channel.start_consuming()

Load Balancing et Scalabilité

# docker-compose.yml - Scaling workers
version: '3.8'
services:
  rabbitmq:
    image: rabbitmq:3-management
    ports:
      - "5672:5672"
      - "15672:15672"
  
  worker:
    build: ./worker
    deploy:
      replicas: 5  # 5 workers en parallèle
    environment:
      RABBITMQ_HOST: rabbitmq
    depends_on:
      - rabbitmq
Rendu du diagramme en cours...

📢 Publish/Subscribe

Fanout Exchange

Broadcasting de messages à tous les consommateurs :

// Publisher
const amqp = require('amqplib');

async function publishNews() {
  const connection = await amqp.connect('amqp://localhost');
  const channel = await connection.createChannel();
  
  // Déclarer un fanout exchange
  await channel.assertExchange('news', 'fanout', { durable: false });
  
  const news = {
    id: Date.now(),
    title: 'Breaking News',
    content: 'Important announcement...',
    timestamp: new Date().toISOString()
  };
  
  // Publier vers l'exchange (pas de routing key nécessaire)
  channel.publish('news', '', Buffer.from(JSON.stringify(news)));
  
  console.log('News published:', news);
  
  await channel.close();
  await connection.close();
}
// Subscriber - Multiple instances
async function subscribeToNews(subscriberName) {
  const connection = await amqp.connect('amqp://localhost');
  const channel = await connection.createChannel();
  
  await channel.assertExchange('news', 'fanout', { durable: false });
  
  // Queue exclusive et temporaire
  const q = await channel.assertQueue('', { exclusive: true });
  
  // Bind la queue à l'exchange
  await channel.bindQueue(q.queue, 'news', '');
  
  console.log(`${subscriberName} waiting for news...`);
  
  channel.consume(q.queue, (msg) => {
    const news = JSON.parse(msg.content.toString());
    console.log(`${subscriberName} received:`, news);
    
    // Traitement spécifique au subscriber
    if (subscriberName === 'EmailService') {
      sendEmailNotification(news);
    } else if (subscriberName === 'PushService') {
      sendPushNotification(news);
    } else if (subscriberName === 'SMSService') {
      sendSMSAlert(news);
    }
    
    channel.ack(msg);
  });
}

// Lancer plusieurs subscribers
subscribeToNews('EmailService');
subscribeToNews('PushService');
subscribeToNews('SMSService');
subscribeToNews('AnalyticsService');

🔀 Routing Pattern

Direct Exchange

Routing basé sur une clé exacte :

# Router - Direct Exchange
class LogRouter:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        
        # Direct exchange pour le routing
        self.channel.exchange_declare(
            exchange='logs_direct',
            exchange_type='direct'
        )
    
    def emit_log(self, severity, message):
        self.channel.basic_publish(
            exchange='logs_direct',
            routing_key=severity,  # 'info', 'warning', 'error'
            body=message
        )
        print(f"[{severity}] {message}")

# Consumers avec filtrage
def subscribe_to_logs(severities):
    connection = pika.BlockingConnection(
        pika.ConnectionParameters('localhost')
    )
    channel = connection.channel()
    
    channel.exchange_declare(exchange='logs_direct', exchange_type='direct')
    
    result = channel.queue_declare(queue='', exclusive=True)
    queue_name = result.method.queue
    
    # Bind pour chaque sévérité souhaitée
    for severity in severities:
        channel.queue_bind(
            exchange='logs_direct',
            queue=queue_name,
            routing_key=severity
        )
    
    def callback(ch, method, properties, body):
        print(f"[{method.routing_key}] Received: {body.decode()}")
        
        # Traitement selon la sévérité
        if method.routing_key == 'error':
            send_alert_to_admin(body)
        elif method.routing_key == 'warning':
            log_to_monitoring_system(body)
    
    channel.basic_consume(
        queue=queue_name,
        on_message_callback=callback,
        auto_ack=True
    )
    
    print(f'Listening for: {severities}')
    channel.start_consuming()

# Consumer 1 : Erreurs seulement
subscribe_to_logs(['error'])

# Consumer 2 : Warnings et erreurs
subscribe_to_logs(['warning', 'error'])

# Consumer 3 : Tout
subscribe_to_logs(['info', 'warning', 'error'])

Topic Exchange

Routing avec patterns et wildcards :

// Topic-based routing
class EventRouter {
  constructor() {
    this.exchangeName = 'events';
  }
  
  async setup() {
    this.connection = await amqp.connect('amqp://localhost');
    this.channel = await this.connection.createChannel();
    
    await this.channel.assertExchange(this.exchangeName, 'topic', {
      durable: true
    });
  }
  
  async publishEvent(routingKey, event) {
    // Routing keys: service.entity.action
    // Ex: order.payment.completed, user.profile.updated
    
    await this.channel.publish(
      this.exchangeName,
      routingKey,
      Buffer.from(JSON.stringify(event))
    );
    
    console.log(`Event published: ${routingKey}`);
  }
  
  async subscribe(pattern, handler) {
    const q = await this.channel.assertQueue('', { exclusive: true });
    
    // Patterns:
    // * = exactement un mot
    // # = zéro ou plusieurs mots
    // order.*.completed = tous les événements completed d'order
    // user.# = tous les événements user
    // #.error = toutes les erreurs
    
    await this.channel.bindQueue(q.queue, this.exchangeName, pattern);
    
    await this.channel.consume(q.queue, async (msg) => {
      const event = JSON.parse(msg.content.toString());
      console.log(`Pattern [${pattern}] matched: ${msg.fields.routingKey}`);
      
      await handler(event, msg.fields.routingKey);
      this.channel.ack(msg);
    });
  }
}

// Usage
const router = new EventRouter();
await router.setup();

// Publishers
await router.publishEvent('order.payment.completed', { orderId: 123, amount: 99.99 });
await router.publishEvent('order.shipping.dispatched', { orderId: 123, trackingNo: 'ABC' });
await router.publishEvent('user.profile.updated', { userId: 456, changes: {...} });
await router.publishEvent('system.monitoring.error', { service: 'api', error: '...' });

// Subscribers avec patterns
// Analytics - tous les événements order
await router.subscribe('order.#', async (event, key) => {
  await analytics.track(key, event);
});

// Notification - tous les événements completed
await router.subscribe('*.*.completed', async (event, key) => {
  await notificationService.send(`${key}: ${event}`);
});

// Error handler - toutes les erreurs
await router.subscribe('#.error', async (event, key) => {
  await errorHandler.process(event);
});

🔄 Request/Reply Pattern (RPC)

Implémentation de RPC sur RabbitMQ :

# RPC Server
class RPCServer:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        self.channel.queue_declare(queue='rpc_queue')
        
    def fibonacci(self, n):
        """Fonction à exposer via RPC"""
        if n <= 1:
            return n
        return self.fibonacci(n-1) + self.fibonacci(n-2)
    
    def on_request(self, ch, method, props, body):
        n = int(body)
        print(f"Calculating fibonacci({n})")
        
        # Calculer la réponse
        response = self.fibonacci(n)
        
        # Envoyer la réponse à la queue de reply
        ch.basic_publish(
            exchange='',
            routing_key=props.reply_to,
            properties=pika.BasicProperties(
                correlation_id=props.correlation_id
            ),
            body=str(response)
        )
        
        ch.basic_ack(delivery_tag=method.delivery_tag)
    
    def start(self):
        self.channel.basic_qos(prefetch_count=1)
        self.channel.basic_consume(
            queue='rpc_queue',
            on_message_callback=self.on_request
        )
        
        print("RPC Server started. Awaiting requests...")
        self.channel.start_consuming()
# RPC Client
import uuid

class RPCClient:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        
        # Queue pour recevoir les réponses
        result = self.channel.queue_declare(queue='', exclusive=True)
        self.callback_queue = result.method.queue
        
        self.channel.basic_consume(
            queue=self.callback_queue,
            on_message_callback=self.on_response,
            auto_ack=True
        )
        
        self.response = None
        self.corr_id = None
    
    def on_response(self, ch, method, props, body):
        if self.corr_id == props.correlation_id:
            self.response = body
    
    def call(self, n, timeout=30):
        self.response = None
        self.corr_id = str(uuid.uuid4())
        
        # Envoyer la requête
        self.channel.basic_publish(
            exchange='',
            routing_key='rpc_queue',
            properties=pika.BasicProperties(
                reply_to=self.callback_queue,
                correlation_id=self.corr_id,
                expiration=str(timeout * 1000)  # TTL en ms
            ),
            body=str(n)
        )
        
        # Attendre la réponse
        start_time = time.time()
        while self.response is None:
            self.connection.process_data_events(time_limit=1)
            
            if time.time() - start_time > timeout:
                raise TimeoutError(f"RPC call timed out after {timeout}s")
        
        return int(self.response)

# Usage
client = RPCClient()
print(f"fibonacci(30) = {client.call(30)}")

🔐 Headers Exchange

Routing basé sur les headers :

// Headers-based routing
async function setupHeadersExchange() {
  const connection = await amqp.connect('amqp://localhost');
  const channel = await connection.createChannel();
  
  await channel.assertExchange('headers_exchange', 'headers', {
    durable: true
  });
  
  // Publisher
  async function publishWithHeaders(message, headers) {
    await channel.publish(
      'headers_exchange',
      '', // Pas de routing key
      Buffer.from(JSON.stringify(message)),
      { headers }
    );
  }
  
  // Subscriber avec match sur headers
  async function subscribeWithHeaders(matchHeaders, matchType = 'all') {
    const q = await channel.assertQueue('', { exclusive: true });
    
    // x-match: 'all' = tous les headers doivent matcher
    // x-match: 'any' = au moins un header doit matcher
    await channel.bindQueue(q.queue, 'headers_exchange', '', {
      'x-match': matchType,
      ...matchHeaders
    });
    
    await channel.consume(q.queue, (msg) => {
      console.log('Matched message:', {
        headers: msg.properties.headers,
        body: JSON.parse(msg.content.toString())
      });
      channel.ack(msg);
    });
  }
  
  // Exemples d'utilisation
  
  // Publisher envoie avec différents headers
  await publishWithHeaders(
    { order: 'Order 1' },
    { format: 'pdf', type: 'invoice', urgent: true }
  );
  
  await publishWithHeaders(
    { order: 'Order 2' },
    { format: 'xml', type: 'report', urgent: false }
  );
  
  // Subscriber 1: Match tous les headers
  await subscribeWithHeaders(
    { format: 'pdf', type: 'invoice' },
    'all'
  );
  
  // Subscriber 2: Match au moins un header
  await subscribeWithHeaders(
    { urgent: true },
    'any'
  );
}

💀 Dead Letter Exchange (DLX)

Gestion des messages en échec :

class DeadLetterSetup:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
    
    def setup_dlx(self):
        # Dead Letter Exchange
        self.channel.exchange_declare(
            exchange='dlx',
            exchange_type='direct'
        )
        
        # Dead Letter Queue
        self.channel.queue_declare(
            queue='dead_letter_queue',
            durable=True
        )
        
        self.channel.queue_bind(
            exchange='dlx',
            queue='dead_letter_queue',
            routing_key='failed'
        )
        
        # Main Queue avec DLX configuré
        self.channel.queue_declare(
            queue='main_queue',
            durable=True,
            arguments={
                'x-dead-letter-exchange': 'dlx',
                'x-dead-letter-routing-key': 'failed',
                'x-message-ttl': 60000,  # TTL de 60 secondes
                'x-max-retries': 3  # Custom header pour retry count
            }
        )
    
    def process_with_dlx(self):
        def callback(ch, method, properties, body):
            try:
                # Traitement qui peut échouer
                data = json.loads(body)
                
                if random.random() < 0.3:  # 30% de chance d'échec
                    raise Exception("Random processing error")
                
                print(f"Successfully processed: {data}")
                ch.basic_ack(delivery_tag=method.delivery_tag)
                
            except Exception as e:
                print(f"Error processing: {e}")
                
                # Vérifier le nombre de tentatives
                retry_count = 0
                if properties.headers:
                    retry_count = properties.headers.get('x-retry-count', 0)
                
                if retry_count < 3:
                    # Republier avec compteur incrémenté
                    ch.basic_publish(
                        exchange='',
                        routing_key='main_queue',
                        body=body,
                        properties=pika.BasicProperties(
                            headers={'x-retry-count': retry_count + 1}
                        )
                    )
                    print(f"Retrying... Attempt {retry_count + 1}")
                else:
                    print(f"Max retries reached. Message sent to DLX")
                    # Rejeter pour envoyer vers DLX
                    ch.basic_nack(
                        delivery_tag=method.delivery_tag,
                        requeue=False
                    )
        
        self.channel.basic_consume(
            queue='main_queue',
            on_message_callback=callback
        )
        
        self.channel.start_consuming()

🏪 Pattern Saga avec Compensation

// Saga Orchestrator avec RabbitMQ
class SagaOrchestrator {
  constructor(channel) {
    this.channel = channel;
    this.sagas = new Map();
  }
  
  async startSaga(sagaId, steps) {
    const saga = {
      id: sagaId,
      steps: steps,
      currentStep: 0,
      state: 'RUNNING',
      compensations: []
    };
    
    this.sagas.set(sagaId, saga);
    await this.executeNextStep(sagaId);
  }
  
  async executeNextStep(sagaId) {
    const saga = this.sagas.get(sagaId);
    
    if (saga.currentStep >= saga.steps.length) {
      saga.state = 'COMPLETED';
      await this.publishSagaEvent(sagaId, 'SAGA_COMPLETED');
      return;
    }
    
    const step = saga.steps[saga.currentStep];
    
    // Publier la commande pour cette étape
    await this.channel.publish(
      'saga_exchange',
      step.command,
      Buffer.from(JSON.stringify({
        sagaId,
        step: saga.currentStep,
        payload: step.payload
      }))
    );
    
    // Attendre la réponse (timeout après 30s)
    setTimeout(() => this.handleTimeout(sagaId), 30000);
  }
  
  async handleStepSuccess(sagaId, compensation) {
    const saga = this.sagas.get(sagaId);
    
    // Enregistrer la compensation pour cette étape
    if (compensation) {
      saga.compensations.push(compensation);
    }
    
    saga.currentStep++;
    await this.executeNextStep(sagaId);
  }
  
  async handleStepFailure(sagaId, error) {
    const saga = this.sagas.get(sagaId);
    saga.state = 'COMPENSATING';
    
    console.log(`Saga ${sagaId} failed at step ${saga.currentStep}: ${error}`);
    
    // Exécuter les compensations dans l'ordre inverse
    for (const compensation of saga.compensations.reverse()) {
      await this.channel.publish(
        'saga_exchange',
        compensation.command,
        Buffer.from(JSON.stringify(compensation.payload))
      );
    }
    
    saga.state = 'FAILED';
    await this.publishSagaEvent(sagaId, 'SAGA_FAILED');
  }
  
  async handleTimeout(sagaId) {
    const saga = this.sagas.get(sagaId);
    if (saga.state === 'RUNNING') {
      await this.handleStepFailure(sagaId, 'Timeout');
    }
  }
}

// Exemple d'utilisation pour une commande e-commerce
const orderSaga = {
  id: 'order-123',
  steps: [
    {
      command: 'reserve.inventory',
      payload: { items: [...] },
      compensation: {
        command: 'release.inventory',
        payload: { reservationId: null }
      }
    },
    {
      command: 'process.payment',
      payload: { amount: 100 },
      compensation: {
        command: 'refund.payment',
        payload: { paymentId: null }
      }
    },
    {
      command: 'create.shipment',
      payload: { address: {...} },
      compensation: {
        command: 'cancel.shipment',
        payload: { shipmentId: null }
      }
    }
  ]
};

📊 Comparaison des Patterns

| Pattern | Use Case | Avantages | Inconvénients | |---------|----------|-----------|---------------| | Work Queue | Traitement parallèle | Scalabilité simple | Pas de broadcast | | Pub/Sub | Broadcasting | Découplage total | Tous reçoivent tout | | Routing | Filtrage par critère | Flexibilité | Configuration complexe | | RPC | Appels synchrones | Familiar pattern | Couplage temporel | | Topic | Routing complexe | Très flexible | Peut être complexe | | DLX | Gestion d'erreurs | Robustesse | Configuration additionnelle |

🎯 Exercice Pratique

Implémentez un système de notification multi-canal :

// TODO: Implémenter
// 1. Publisher qui envoie des notifications
// 2. Router avec topic exchange (email.*, sms.*, push.*)
// 3. Consumers pour chaque canal
// 4. DLX pour les notifications échouées
// 5. Monitoring et métriques

Les patterns de messagerie sont les briques de base pour construire des architectures distribuées robustes et scalables avec RabbitMQ.

📝 Testez vos connaissances !

Répondez à 10 questions pour valider ce cours