💼 Cas d'Usage et Bonnes Pratiques

Introduction

RabbitMQ excelle dans de nombreux scénarios réels. Cette leçon explore des cas d'usage concrets avec des implémentations complètes et les bonnes pratiques associées.

🛒 Cas d'Usage 1: E-Commerce Pipeline

Architecture Globale

Rendu du diagramme en cours...

Implémentation Complète

# order_service.py
import pika
import json
import uuid
from datetime import datetime
from enum import Enum

class OrderStatus(Enum):
    CREATED = "created"
    PAYMENT_PENDING = "payment_pending"
    PAID = "paid"
    SHIPPED = "shipped"
    DELIVERED = "delivered"
    CANCELLED = "cancelled"

class ECommerceOrderService:
    def __init__(self, rabbitmq_url='amqp://localhost'):
        # Connexion avec pool
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url)
        )
        self.channel = self.connection.channel()
        self.setup_infrastructure()
        
    def setup_infrastructure(self):
        """Setup exchanges, queues et bindings"""
        
        # === EXCHANGES ===
        exchanges = [
            ('ecom.commands', 'direct'),
            ('ecom.events', 'topic'),
            ('ecom.notifications', 'fanout'),
            ('ecom.dlx', 'direct')
        ]
        
        for name, type in exchanges:
            self.channel.exchange_declare(
                exchange=name,
                exchange_type=type,
                durable=True
            )
        
        # === QUEUES ===
        queues = {
            # Processing queues
            'order.payment': {'priority': 10, 'ttl': 600000},
            'order.inventory': {'priority': 5, 'ttl': 300000},
            'order.shipping': {'priority': 3, 'ttl': 1800000},
            
            # Notification queues
            'notifications.email': {},
            'notifications.sms': {},
            'notifications.push': {},
            
            # Analytics
            'analytics.events': {'ttl': 86400000},  # 24h
            
            # Dead letters
            'dead.letters': {}
        }
        
        for queue_name, args in queues.items():
            queue_args = {
                'x-dead-letter-exchange': 'ecom.dlx',
                'x-dead-letter-routing-key': 'failed'
            }
            
            if 'priority' in args:
                queue_args['x-max-priority'] = args['priority']
            if 'ttl' in args:
                queue_args['x-message-ttl'] = args['ttl']
                
            self.channel.queue_declare(
                queue=queue_name,
                durable=True,
                arguments=queue_args
            )
        
        # === BINDINGS ===
        # Command routing
        command_bindings = [
            ('order.payment', 'payment.process'),
            ('order.inventory', 'inventory.reserve'),
            ('order.shipping', 'shipping.create')
        ]
        
        for queue, routing_key in command_bindings:
            self.channel.queue_bind(
                exchange='ecom.commands',
                queue=queue,
                routing_key=routing_key
            )
        
        # Event routing
        event_bindings = [
            ('analytics.events', 'order.#'),           # Tous les events order
            ('notifications.email', 'order.*.confirmed'),
            ('order.inventory', 'payment.completed'),
            ('order.shipping', 'inventory.reserved')
        ]
        
        for queue, pattern in event_bindings:
            self.channel.queue_bind(
                exchange='ecom.events',
                queue=queue,
                routing_key=pattern
            )
        
        # Notification fanout
        for ntype in ['email', 'sms', 'push']:
            self.channel.queue_bind(
                exchange='ecom.notifications',
                queue=f'notifications.{ntype}'
            )
        
        # DLX binding
        self.channel.queue_bind(
            exchange='ecom.dlx',
            queue='dead.letters',
            routing_key='#'
        )
    
    def process_order(self, order_data):
        """Processing principal d'une commande"""
        order_id = str(uuid.uuid4())
        
        order = {
            'id': order_id,
            'customer_id': order_data['customer_id'],
            'items': order_data['items'],
            'total_amount': order_data['total_amount'],
            'status': OrderStatus.CREATED.value,
            'created_at': datetime.now().isoformat(),
            'metadata': {
                'source': 'web',
                'correlation_id': str(uuid.uuid4())
            }
        }
        
        try:
            # 1. Sauvegarder en base (non montré ici)
            # db.orders.insert(order)
            
            # 2. Démarrer le workflow de paiement
            self.channel.basic_publish(
                exchange='ecom.commands',
                routing_key='payment.process',
                body=json.dumps({
                    'order_id': order_id,
                    'amount': order['total_amount'],
                    'customer_id': order['customer_id']
                }),
                properties=pika.BasicProperties(
                    correlation_id=order['metadata']['correlation_id'],
                    priority=5,
                    delivery_mode=2
                )
            )
            
            # 3. Publier événement de création
            self.publish_event('order.created.new', {
                'orderId': order_id,
                'customerId': order['customer_id'],
                'amount': order['total_amount']
            })
            
            return {'success': True, 'order_id': order_id}
            
        except Exception as e:
            # Publier événement d'échec
            self.publish_event('order.created.failed', {
                'error': str(e),
                'order_data': order_data
            })
            return {'success': False, 'error': str(e)}
    
    def publish_event(self, routing_key, event_data):
        """Helper pour publier des événements"""
        event = {
            'id': str(uuid.uuid4()),
            'type': routing_key,
            'timestamp': datetime.now().isoformat(),
            'data': event_data
        }
        
        self.channel.basic_publish(
            exchange='ecom.events',
            routing_key=routing_key,
            body=json.dumps(event),
            properties=pika.BasicProperties(delivery_mode=2)
        )

# payment_service.py
class PaymentService:
    def __init__(self, rabbitmq_url='amqp://localhost'):
        self.connection = pika.BlockingConnection(
            pika.URLParameters(rabbitmq_url)
        )
        self.channel = self.connection.channel()
        
    def start_consuming(self):
        """Démarrer la consommation des commandes de paiement"""
        self.channel.basic_qos(prefetch_count=5)  # Traiter 5 à la fois
        
        self.channel.basic_consume(
            queue='order.payment',
            on_message_callback=self.process_payment_command
        )
        
        print("Payment Service started. Awaiting commands...")
        self.channel.start_consuming()
    
    def process_payment_command(self, ch, method, properties, body):
        """Traiter une commande de paiement"""
        try:
            command = json.loads(body)
            order_id = command['order_id']
            amount = command['amount']
            
            print(f"Processing payment for order {order_id}: ${amount}")
            
            # Simulation du traitement paiement
            payment_result = self.process_external_payment(
                command['customer_id'], 
                amount
            )
            
            if payment_result['success']:
                # Publier succès
                self.publish_payment_event('payment.completed', {
                    'orderId': order_id,
                    'paymentId': payment_result['payment_id'],
                    'amount': amount
                })
                
                # Déclencher étape suivante
                ch.basic_publish(
                    exchange='ecom.commands',
                    routing_key='inventory.reserve',
                    body=json.dumps({
                        'order_id': order_id,
                        'items': command.get('items', [])
                    })
                )
            else:
                # Publier échec
                self.publish_payment_event('payment.failed', {
                    'orderId': order_id,
                    'error': payment_result['error']
                })
            
            ch.basic_ack(delivery_tag=method.delivery_tag)
            
        except Exception as e:
            print(f"Error processing payment: {e}")
            # Reject avec requeue pour retry
            ch.basic_nack(
                delivery_tag=method.delivery_tag,
                requeue=True
            )
    
    def publish_payment_event(self, event_type, data):
        """Publier événement de paiement"""
        event = {
            'type': event_type,
            'timestamp': datetime.now().isoformat(),
            'data': data
        }
        
        self.channel.basic_publish(
            exchange='ecom.events',
            routing_key=f'order.{event_type}',
            body=json.dumps(event)
        )
    
    def process_external_payment(self, customer_id, amount):
        """Simulation traitement paiement externe"""
        import random
        import time
        
        # Simulate API call delay
        time.sleep(0.5)
        
        # 90% de succès
        if random.random() < 0.9:
            return {
                'success': True,
                'payment_id': f'pay_{uuid.uuid4().hex[:8]}'
            }
        else:
            return {
                'success': False,
                'error': 'Payment gateway timeout'
            }

🤖 Cas d'Usage 2: IoT Data Processing

Architecture IoT Scalable

// iot-data-processor.js
const amqp = require('amqplib');
const { InfluxDB } = require('@influxdata/influxdb-client');

class IoTDataProcessor {
  constructor() {
    this.rabbitmq = null;
    this.influxdb = new InfluxDB({
      url: 'http://localhost:8086',
      token: process.env.INFLUX_TOKEN
    });
  }
  
  async setup() {
    // Connexion RabbitMQ
    this.connection = await amqp.connect('amqp://localhost');
    this.channel = await this.connection.createChannel();
    
    await this.setupInfrastructure();
  }
  
  async setupInfrastructure() {
    // Exchange pour données IoT avec routing par sensor type
    await this.channel.assertExchange('iot.data', 'topic', {
      durable: true
    });
    
    // Exchange pour alertes
    await this.channel.assertExchange('iot.alerts', 'direct', {
      durable: true
    });
    
    // Queues par type de processing
    const queues = {
      // Real-time processing
      'iot.temperature.realtime': {
        pattern: 'sensor.temperature.*',
        maxLength: 1000
      },
      'iot.humidity.realtime': {
        pattern: 'sensor.humidity.*',
        maxLength: 1000
      },
      
      // Batch processing
      'iot.batch.hourly': {
        pattern: 'sensor.#',
        ttl: 3600000  // 1 hour
      },
      
      // ML processing
      'iot.ml.anomaly': {
        pattern: 'sensor.*.anomaly',
        maxLength: 10000
      },
      
      // Alertes critiques
      'iot.alerts.critical': {
        pattern: 'critical',
        priority: 10
      }
    };
    
    for (const [queueName, config] of Object.entries(queues)) {
      const args = {};
      
      if (config.maxLength) {
        args['x-max-length'] = config.maxLength;
        args['x-overflow'] = 'drop-head';  // Drop oldest
      }
      
      if (config.ttl) {
        args['x-message-ttl'] = config.ttl;
      }
      
      if (config.priority) {
        args['x-max-priority'] = config.priority;
      }
      
      await this.channel.assertQueue(queueName, {
        durable: true,
        arguments: args
      });
      
      // Bind to appropriate exchange
      if (config.pattern) {
        const exchange = queueName.includes('alerts') ? 'iot.alerts' : 'iot.data';
        await this.channel.bindQueue(queueName, exchange, config.pattern);
      }
    }
  }
  
  // Simulateur de capteurs
  async simulateSensors() {
    const sensorTypes = ['temperature', 'humidity', 'pressure', 'vibration'];
    const locations = ['factory-1', 'factory-2', 'warehouse-a', 'warehouse-b'];
    
    setInterval(async () => {
      for (const location of locations) {
        for (const sensorType of sensorTypes) {
          const value = this.generateSensorValue(sensorType);
          const isAnomaly = this.detectAnomaly(sensorType, value);
          
          const sensorData = {
            sensorId: `${location}-${sensorType}-01`,
            type: sensorType,
            value: value,
            location: location,
            timestamp: new Date().toISOString(),
            metadata: {
              firmware: '1.2.3',
              battery: Math.random() * 100
            }
          };
          
          // Routing key: sensor.{type}.{location}
          let routingKey = `sensor.${sensorType}.${location}`;
          
          if (isAnomaly) {
            routingKey += '.anomaly';
            
            // Alerte critique si température > 80°C
            if (sensorType === 'temperature' && value > 80) {
              await this.channel.publish(
                'iot.alerts',
                'critical',
                Buffer.from(JSON.stringify({
                  alert: 'High temperature detected',
                  sensor: sensorData,
                  severity: 'critical'
                }))
              );
            }
          }
          
          // Publier la donnée
          await this.channel.publish(
            'iot.data',
            routingKey,
            Buffer.from(JSON.stringify(sensorData)),
            {
              timestamp: Date.now(),
              headers: {
                'sensor-type': sensorType,
                'location': location,
                'anomaly': isAnomaly
              }
            }
          );
        }
      }
    }, 5000); // Toutes les 5 secondes
  }
  
  generateSensorValue(type) {
    switch (type) {
      case 'temperature': return 20 + Math.random() * 60; // 20-80°C
      case 'humidity': return Math.random() * 100;        // 0-100%
      case 'pressure': return 900 + Math.random() * 200;  // 900-1100 hPa
      case 'vibration': return Math.random() * 10;        // 0-10 G
      default: return Math.random();
    }
  }
  
  detectAnomaly(type, value) {
    const thresholds = {
      temperature: { min: 10, max: 70 },
      humidity: { min: 5, max: 95 },
      pressure: { min: 950, max: 1050 },
      vibration: { min: 0, max: 8 }
    };
    
    const threshold = thresholds[type];
    return value < threshold.min || value > threshold.max;
  }
}

// Processeurs spécialisés
class TemperatureProcessor {
  constructor(channel) {
    this.channel = channel;
    this.influx = new InfluxDB(...);
  }
  
  async start() {
    this.channel.consume('iot.temperature.realtime', async (msg) => {
      const data = JSON.parse(msg.content.toString());
      
      // Stockage temps-réel en InfluxDB
      await this.influx.writePoint({
        measurement: 'temperature',
        tags: {
          sensor_id: data.sensorId,
          location: data.location
        },
        fields: {
          value: data.value,
          battery: data.metadata.battery
        },
        timestamp: new Date(data.timestamp)
      });
      
      // Calcul de moyennes glissantes
      const avg = await this.calculateMovingAverage(data.sensorId, 5);
      
      if (Math.abs(data.value - avg) > 10) {
        // Publier alerte de dérive
        await this.channel.publish(
          'iot.alerts',
          'temperature.drift',
          Buffer.from(JSON.stringify({
            sensor: data.sensorId,
            current: data.value,
            average: avg,
            deviation: Math.abs(data.value - avg)
          }))
        );
      }
      
      this.channel.ack(msg);
    });
  }
}

📧 Cas d'Usage 3: Système de Notification Multi-Canal

# notification_system.py
class NotificationSystem:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        self.setup_notification_system()
    
    def setup_notification_system(self):
        """Setup système de notification intelligent"""
        
        # Exchange avec routing par priorité et type
        self.channel.exchange_declare(
            exchange='notifications',
            exchange_type='topic',
            durable=True
        )
        
        # Queues par canal avec QoS différents
        channels = {
            'email': {'qos': 50, 'priority': 3},
            'sms': {'qos': 10, 'priority': 8},      # Limité et cher
            'push': {'qos': 100, 'priority': 5},
            'slack': {'qos': 20, 'priority': 6}
        }
        
        for channel_name, config in channels.items():
            # Queue normale
            self.channel.queue_declare(
                queue=f'notifications.{channel_name}',
                durable=True,
                arguments={
                    'x-max-priority': config['priority'],
                    'x-dead-letter-exchange': 'notifications.dlx'
                }
            )
            
            # Queue retry avec délai
            self.channel.queue_declare(
                queue=f'notifications.{channel_name}.retry',
                durable=True,
                arguments={
                    'x-message-ttl': 30000,  # 30s retry delay
                    'x-dead-letter-exchange': 'notifications',
                    'x-dead-letter-routing-key': f'{channel_name}.normal'
                }
            )
            
            # Bindings avec patterns
            patterns = [
                f'{channel_name}.#',           # Toutes les notif de ce canal
                f'*.{channel_name}.*',         # Pattern global pour le canal
                'urgent.#'                     # Toutes les urgentes
            ]
            
            for pattern in patterns:
                self.channel.queue_bind(
                    exchange='notifications',
                    queue=f'notifications.{channel_name}',
                    routing_key=pattern
                )
    
    def send_notification(self, channel_type, priority, recipient, content):
        """Envoyer notification avec routing intelligent"""
        
        notification = {
            'id': str(uuid.uuid4()),
            'channel': channel_type,
            'recipient': recipient,
            'content': content,
            'created_at': datetime.now().isoformat(),
            'attempts': 0
        }
        
        # Routing key basé sur priorité et canal
        priority_level = 'urgent' if priority > 7 else 'normal'
        routing_key = f'{priority_level}.{channel_type}.notification'
        
        self.channel.basic_publish(
            exchange='notifications',
            routing_key=routing_key,
            body=json.dumps(notification),
            properties=pika.BasicProperties(
                priority=priority,
                delivery_mode=2,
                headers={
                    'notification-type': channel_type,
                    'recipient-id': recipient.get('id'),
                    'campaign-id': content.get('campaign_id')
                }
            )
        )

# Processeur Email avec rate limiting
class EmailProcessor:
    def __init__(self, channel):
        self.channel = channel
        self.rate_limiter = TokenBucket(rate=100, capacity=1000)  # 100/sec
        
    async def start(self):
        # QoS pour limiter les messages en parallèle
        self.channel.basic_qos(prefetch_count=50)
        
        self.channel.basic_consume(
            queue='notifications.email',
            on_message_callback=self.process_email
        )
        
    def process_email(self, ch, method, properties, body):
        try:
            # Rate limiting
            if not self.rate_limiter.consume():
                # Trop de requêtes, remettre en queue avec délai
                ch.basic_nack(
                    delivery_tag=method.delivery_tag,
                    requeue=True
                )
                time.sleep(0.1)
                return
            
            notification = json.loads(body)
            
            # Envoyer l'email
            result = self.send_email(notification)
            
            if result['success']:
                ch.basic_ack(delivery_tag=method.delivery_tag)
            else:
                # Retry avec exponential backoff
                self.schedule_retry(notification, ch, method)
                
        except Exception as e:
            print(f"Email processing error: {e}")
            ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

🎮 Cas d'Usage 4: Gaming Matchmaking

# gaming_matchmaking.py
class GameMatchmaking:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        self.setup_matchmaking()
    
    def setup_matchmaking(self):
        # Exchange pour events de matchmaking
        self.channel.exchange_declare(
            exchange='matchmaking',
            exchange_type='topic',
            durable=True
        )
        
        # Queues par skill level et région
        skill_levels = ['bronze', 'silver', 'gold', 'diamond']
        regions = ['eu', 'na', 'asia']
        
        for skill in skill_levels:
            for region in regions:
                queue_name = f'matchmaking.{skill}.{region}'
                
                self.channel.queue_declare(
                    queue=queue_name,
                    durable=True,
                    arguments={
                        'x-message-ttl': 60000,  # 1 minute max wait
                        'x-max-length': 100      # Max 100 joueurs en attente
                    }
                )
                
                # Bind avec pattern
                self.channel.queue_bind(
                    exchange='matchmaking',
                    queue=queue_name,
                    routing_key=f'player.{skill}.{region}'
                )
        
        # Queue pour matches créés
        self.channel.queue_declare(
            queue='matches.created',
            durable=True
        )
        
        self.channel.queue_bind(
            exchange='matchmaking',
            queue='matches.created',
            routing_key='match.created'
        )
    
    def queue_player(self, player):
        """Mettre un joueur en file d'attente"""
        routing_key = f"player.{player['skill_level']}.{player['region']}"
        
        queue_entry = {
            'player_id': player['id'],
            'skill_rating': player['rating'],
            'preferred_game_mode': player['game_mode'],
            'queued_at': datetime.now().isoformat()
        }
        
        self.channel.basic_publish(
            exchange='matchmaking',
            routing_key=routing_key,
            body=json.dumps(queue_entry),
            properties=pika.BasicProperties(
                priority=player.get('premium', 0) * 5,  # Premium priority
                expiration='60000'  # Expire après 1 minute
            )
        )
    
    def start_matchmaking_service(self, skill_level, region):
        """Service de matchmaking pour un skill/region"""
        queue_name = f'matchmaking.{skill_level}.{region}'
        
        waiting_players = []
        
        def process_queue_entry(ch, method, properties, body):
            player_data = json.loads(body)
            waiting_players.append(player_data)
            
            # Essayer de créer un match
            if len(waiting_players) >= 2:  # Match 1v1
                match = self.create_match(waiting_players[:2])
                waiting_players.clear()
                
                # Publier match créé
                ch.basic_publish(
                    exchange='matchmaking',
                    routing_key='match.created',
                    body=json.dumps(match)
                )
            
            ch.basic_ack(delivery_tag=method.delivery_tag)
        
        self.channel.basic_consume(
            queue=queue_name,
            on_message_callback=process_queue_entry
        )
        
        print(f"Matchmaking service started for {skill_level}/{region}")
        self.channel.start_consuming()

🎯 Bonnes Pratiques Avancées

1. Pattern de Resilience

class ResilientPublisher:
    def __init__(self, connection_params):
        self.params = connection_params
        self.connection = None
        self.channel = None
        self.reconnect()
    
    def reconnect(self):
        """Reconnexion avec exponential backoff"""
        max_attempts = 5
        base_delay = 1
        
        for attempt in range(max_attempts):
            try:
                if self.connection and not self.connection.is_closed:
                    self.connection.close()
                
                self.connection = pika.BlockingConnection(self.params)
                self.channel = self.connection.channel()
                self.channel.confirm_delivery()  # Publisher confirms
                
                print("✅ Reconnected to RabbitMQ")
                return True
                
            except Exception as e:
                delay = base_delay * (2 ** attempt)
                print(f"❌ Reconnection attempt {attempt + 1} failed: {e}")
                print(f"⏳ Retrying in {delay} seconds...")
                time.sleep(delay)
        
        raise Exception("Failed to reconnect after maximum attempts")
    
    def publish_with_retry(self, exchange, routing_key, message, max_retries=3):
        """Publication avec retry automatique"""
        for attempt in range(max_retries):
            try:
                confirmed = self.channel.basic_publish(
                    exchange=exchange,
                    routing_key=routing_key,
                    body=json.dumps(message),
                    properties=pika.BasicProperties(delivery_mode=2),
                    mandatory=True
                )
                
                if confirmed:
                    return True
                    
            except (pika.exceptions.ConnectionClosed, 
                    pika.exceptions.ChannelClosed) as e:
                print(f"Connection lost, reconnecting... (attempt {attempt + 1})")
                self.reconnect()
                
            except Exception as e:
                print(f"Publish error: {e}")
                if attempt == max_retries - 1:
                    raise
                time.sleep(0.5 * (attempt + 1))
        
        return False

2. Message Versioning

// Version management pour backward compatibility
class VersionedMessage {
  static create(type, data, version = '1.0') {
    return {
      meta: {
        version: version,
        type: type,
        timestamp: new Date().toISOString(),
        id: require('uuid').v4()
      },
      payload: data
    };
  }
  
  static upgrade(message) {
    const { version, type } = message.meta;
    
    // Migration V1 -> V2
    if (version === '1.0' && type === 'OrderCreated') {
      message.payload.shippingAddress = message.payload.address;
      delete message.payload.address;
      message.meta.version = '1.1';
    }
    
    // Migration V1.1 -> V2.0
    if (version === '1.1' && type === 'OrderCreated') {
      message.payload.items = message.payload.items.map(item => ({
        ...item,
        sku: item.productId,
        unitPrice: item.price
      }));
      message.meta.version = '2.0';
    }
    
    return message;
  }
}

// Consumer avec upgrade automatique
class VersionedConsumer {
  async processMessage(rawMessage) {
    let message = JSON.parse(rawMessage.content.toString());
    
    // Upgrade vers la dernière version
    message = VersionedMessage.upgrade(message);
    
    // Traiter avec la version courante
    switch (message.meta.type) {
      case 'OrderCreated':
        await this.handleOrderCreated(message.payload);
        break;
      // ...
    }
  }
}

3. Circuit Breaker Pattern

import time
from enum import Enum

class CircuitState(Enum):
    CLOSED = "closed"
    OPEN = "open"
    HALF_OPEN = "half_open"

class CircuitBreaker:
    def __init__(self, failure_threshold=5, timeout=60):
        self.failure_threshold = failure_threshold
        self.timeout = timeout
        self.failure_count = 0
        self.last_failure = None
        self.state = CircuitState.CLOSED
    
    def call(self, func, *args, **kwargs):
        """Appeler une fonction avec circuit breaker"""
        
        if self.state == CircuitState.OPEN:
            if time.time() - self.last_failure > self.timeout:
                self.state = CircuitState.HALF_OPEN
                print("🔄 Circuit breaker: HALF_OPEN")
            else:
                raise Exception("Circuit breaker OPEN - call rejected")
        
        try:
            result = func(*args, **kwargs)
            
            if self.state == CircuitState.HALF_OPEN:
                self.state = CircuitState.CLOSED
                self.failure_count = 0
                print("✅ Circuit breaker: CLOSED")
            
            return result
            
        except Exception as e:
            self.failure_count += 1
            self.last_failure = time.time()
            
            if self.failure_count >= self.failure_threshold:
                self.state = CircuitState.OPEN
                print("🚫 Circuit breaker: OPEN")
            
            raise e

# Usage avec RabbitMQ
class ResilientConsumer:
    def __init__(self):
        self.channel = self.get_channel()
        self.external_api_breaker = CircuitBreaker(
            failure_threshold=3,
            timeout=30
        )
    
    def process_message(self, ch, method, properties, body):
        try:
            data = json.loads(body)
            
            # Appel externe avec circuit breaker
            result = self.external_api_breaker.call(
                self.call_external_api,
                data
            )
            
            ch.basic_ack(delivery_tag=method.delivery_tag)
            
        except Exception as e:
            print(f"Processing failed: {e}")
            
            # Si circuit ouvert, rejeter sans requeue
            if "Circuit breaker OPEN" in str(e):
                ch.basic_nack(
                    delivery_tag=method.delivery_tag,
                    requeue=False
                )
            else:
                # Autre erreur, requeue pour retry
                ch.basic_nack(
                    delivery_tag=method.delivery_tag,
                    requeue=True
                )

📊 Monitoring et Observabilité

Custom Metrics

from prometheus_client import Counter, Histogram, Gauge

class RabbitMQMetrics:
    def __init__(self):
        # Métriques métier
        self.messages_processed = Counter(
            'rabbitmq_messages_processed_total',
            'Total processed messages',
            ['service', 'queue', 'status']
        )
        
        self.processing_duration = Histogram(
            'rabbitmq_message_processing_duration_seconds',
            'Message processing duration',
            ['service', 'message_type']
        )
        
        self.queue_length = Gauge(
            'rabbitmq_queue_length',
            'Current queue length',
            ['queue']
        )
        
        self.consumer_lag = Gauge(
            'rabbitmq_consumer_lag_seconds',
            'Consumer lag',
            ['queue']
        )
    
    def record_processing(self, service, queue, status, duration):
        self.messages_processed.labels(
            service=service,
            queue=queue,
            status=status
        ).inc()
        
        if status == 'success':
            self.processing_duration.labels(
                service=service,
                message_type=queue
            ).observe(duration)

# Wrapper instrumenté
class InstrumentedConsumer:
    def __init__(self, service_name, metrics):
        self.service_name = service_name
        self.metrics = metrics
    
    def consume_with_metrics(self, queue, callback):
        def instrumented_callback(ch, method, properties, body):
            start_time = time.time()
            
            try:
                # Message processing
                callback(ch, method, properties, body)
                
                duration = time.time() - start_time
                self.metrics.record_processing(
                    self.service_name,
                    queue,
                    'success',
                    duration
                )
                
            except Exception as e:
                duration = time.time() - start_time
                self.metrics.record_processing(
                    self.service_name,
                    queue,
                    'error',
                    duration
                )
                raise e
        
        return instrumented_callback

✅ Checklist des Bonnes Pratiques

Architecture

  • ✅ Utiliser des exchanges typés appropriés
  • ✅ Nommer les resources de façon cohérente
  • ✅ Configurer TTL sur les messages
  • ✅ Implémenter Dead Letter Queues
  • ✅ Séparer par Virtual Hosts

Code

  • ✅ Gérer les reconnexions automatiques
  • ✅ Utiliser publisher confirms
  • ✅ Implémenter consumer acknowledgments
  • ✅ Ajouter circuit breakers
  • ✅ Instrumenter avec métriques

Opérations

  • ✅ Monitoring complet (queues, latence, erreurs)
  • ✅ Alerting sur métriques critiques
  • ✅ Backup des définitions
  • ✅ Documentation des flows
  • ✅ Tests de disaster recovery

🎯 Exercice Final

Implémentez un système complet de e-learning :

// TODO: Créer un système avec :
// 1. Inscription étudiants (events)
// 2. Progression cours (commands)
// 3. Notifications certificats (pub/sub)
// 4. Analytics engagement (streaming)
// 5. Système de retry et error handling
// 6. Monitoring complet

Ces cas d'usage démontrent la puissance et la flexibilité de RabbitMQ pour résoudre des problèmes complexes de communication dans des systèmes distribués.

📝 Testez vos connaissances !

Répondez à 10 questions pour valider ce cours