🔀 Types d'Exchanges

Introduction

Les exchanges sont le cœur du système de routing de RabbitMQ. Ils déterminent comment les messages sont acheminés vers les queues. RabbitMQ propose quatre types d'exchanges, chacun avec sa logique de routage spécifique.

🎯 Direct Exchange

Le Direct Exchange route les messages vers les queues dont la routing key correspond exactement à celle du message.

Principe de Fonctionnement

Rendu du diagramme en cours...

Implémentation Pratique

import pika
import json
from enum import Enum

class LogLevel(Enum):
    INFO = "info"
    WARNING = "warning"
    ERROR = "error"
    CRITICAL = "critical"

class DirectExchangeLogger:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        self.setup()
    
    def setup(self):
        # Déclarer le direct exchange
        self.channel.exchange_declare(
            exchange='logs_direct',
            exchange_type='direct',
            durable=True
        )
        
        # Queues pour chaque niveau
        log_queues = {
            'logs.info': LogLevel.INFO.value,
            'logs.warning': LogLevel.WARNING.value,
            'logs.error': LogLevel.ERROR.value,
            'logs.critical': LogLevel.CRITICAL.value,
            'logs.admin': None  # Recevra error + critical
        }
        
        for queue_name, routing_key in log_queues.items():
            self.channel.queue_declare(queue=queue_name, durable=True)
            
            if routing_key:
                # Binding standard
                self.channel.queue_bind(
                    exchange='logs_direct',
                    queue=queue_name,
                    routing_key=routing_key
                )
            else:
                # Admin queue reçoit les erreurs critiques
                for critical_level in [LogLevel.ERROR.value, LogLevel.CRITICAL.value]:
                    self.channel.queue_bind(
                        exchange='logs_direct',
                        queue=queue_name,
                        routing_key=critical_level
                    )
    
    def log(self, level: LogLevel, message: str, context=None):
        """Envoyer un log avec le niveau approprié"""
        log_entry = {
            'level': level.value,
            'message': message,
            'timestamp': datetime.now().isoformat(),
            'service': 'logger-service',
            'context': context or {}
        }
        
        self.channel.basic_publish(
            exchange='logs_direct',
            routing_key=level.value,  # Routing exacte
            body=json.dumps(log_entry),
            properties=pika.BasicProperties(
                delivery_mode=2,
                priority=self.get_priority(level)
            )
        )
    
    def get_priority(self, level):
        priorities = {
            LogLevel.INFO: 1,
            LogLevel.WARNING: 5,
            LogLevel.ERROR: 8,
            LogLevel.CRITICAL: 10
        }
        return priorities[level]

# Usage
logger = DirectExchangeLogger()
logger.log(LogLevel.INFO, "User logged in", {'user_id': 123})
logger.log(LogLevel.ERROR, "Database connection failed", {'db': 'postgres'})
logger.log(LogLevel.CRITICAL, "System out of memory", {'mem_usage': '95%'})

🌐 Topic Exchange

Le Topic Exchange utilise des patterns avec wildcards pour router les messages.

Wildcards

  • * : exactly one word
  • # : zero or more words
class TopicExchangeRouter:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        self.setup_topic_routing()
    
    def setup_topic_routing(self):
        # Topic exchange
        self.channel.exchange_declare(
            exchange='events_topic',
            exchange_type='topic',
            durable=True
        )
        
        # Scenarios de routing complexe
        routing_scenarios = {
            # Analytics - tous les événements
            'analytics.all_events': '#',
            
            # Notifications - tous les completed
            'notifications.completed': '*.*.completed',
            
            # Orders - tous les événements order
            'processing.orders': 'order.#',
            
            # Payments EU - paiements EU seulement
            'payments.eu': 'payment.eu.*',
            
            # Errors - toutes les erreurs
            'monitoring.errors': '#.error',
            
            # User activity - activité utilisateur
            'user.activity': 'user.*.activity',
            
            # Critical alerts - alertes par région
            'alerts.us': 'alert.us.#',
            'alerts.eu': 'alert.eu.#'
        }
        
        for queue_name, pattern in routing_scenarios.items():
            self.channel.queue_declare(queue=queue_name, durable=True)
            self.channel.queue_bind(
                exchange='events_topic',
                queue=queue_name,
                routing_key=pattern
            )
    
    def publish_event(self, routing_key, event_data):
        """Publier avec routing key structurée"""
        # Format: entity.region.action
        # Ex: order.eu.completed, user.us.activity, alert.asia.critical
        
        message = {
            'routing_key': routing_key,
            'timestamp': datetime.now().isoformat(),
            'data': event_data
        }
        
        self.channel.basic_publish(
            exchange='events_topic',
            routing_key=routing_key,
            body=json.dumps(message)
        )
        
        print(f"Published: {routing_key}")

# Examples d'utilisation
router = TopicExchangeRouter()

# Ces messages iront dans différentes queues selon les patterns
router.publish_event('order.eu.completed', {'order_id': 123})
router.publish_event('user.us.activity', {'user_id': 456, 'action': 'login'})
router.publish_event('payment.eu.processed', {'payment_id': 789})
router.publish_event('alert.us.critical', {'message': 'Server down'})
router.publish_event('system.monitoring.error', {'error': 'Connection timeout'})

Routing Matrix

Rendu du diagramme en cours...

📻 Fanout Exchange

Le Fanout Exchange broadcasts tous les messages vers toutes les queues liées, ignorant la routing key.

// notification-system.js
const amqp = require('amqplib');

class FanoutNotificationSystem {
  constructor() {
    this.connection = null;
    this.channel = null;
  }
  
  async setup() {
    this.connection = await amqp.connect('amqp://localhost');
    this.channel = await this.connection.createChannel();
    
    // Fanout exchange pour broadcast
    await this.channel.assertExchange('notifications_fanout', 'fanout', {
      durable: true
    });
    
    // Chaque service crée sa propre queue exclusive
    const services = ['email', 'sms', 'push', 'slack', 'webhook'];
    
    for (const service of services) {
      const queueName = `notifications.${service}`;
      
      await this.channel.assertQueue(queueName, {
        durable: false,  // Temporary queues pour notifications
        exclusive: false
      });
      
      // Bind sans routing key (fanout l'ignore)
      await this.channel.bindQueue(queueName, 'notifications_fanout');
      
      console.log(`${service} service queue bound to fanout`);
    }
  }
  
  async broadcastNotification(notification) {
    /**
     * Broadcaster une notification à tous les services
     */
    const message = {
      id: require('uuid').v4(),
      type: 'broadcast',
      title: notification.title,
      content: notification.content,
      priority: notification.priority || 'normal',
      timestamp: new Date().toISOString(),
      metadata: notification.metadata || {}
    };
    
    // Publier sans routing key (fanout broadcast tout)
    await this.channel.publish(
      'notifications_fanout',
      '', // Routing key ignorée
      Buffer.from(JSON.stringify(message)),
      {
        persistent: true,
        priority: message.priority === 'high' ? 8 : 3
      }
    );
    
    console.log(`📢 Broadcast sent: ${message.title}`);
    return message.id;
  }
  
  async subscribeService(serviceName, handler) {
    /**
     * Subscriber pour un service spécifique
     */
    const queueName = `notifications.${serviceName}`;
    
    await this.channel.consume(queueName, async (msg) => {
      if (!msg) return;
      
      try {
        const notification = JSON.parse(msg.content.toString());
        console.log(`${serviceName} processing:`, notification.title);
        
        // Traitement spécialisé
        await handler(notification);
        
        this.channel.ack(msg);
      } catch (error) {
        console.error(`${serviceName} error:`, error);
        this.channel.nack(msg, false, false); // Send to DLX
      }
    });
  }
}

// Service handlers
const emailHandler = async (notification) => {
  console.log(`📧 Sending email: ${notification.title}`);
  // await emailService.send(...)
};

const smsHandler = async (notification) => {
  if (notification.priority === 'high') {
    console.log(`📱 Sending SMS: ${notification.title}`);
    // await smsService.send(...)
  }
};

const slackHandler = async (notification) => {
  console.log(`💬 Posting to Slack: ${notification.title}`);
  // await slackService.postMessage(...)
};

// Usage
const notificationSystem = new FanoutNotificationSystem();
await notificationSystem.setup();

// Subscribers
notificationSystem.subscribeService('email', emailHandler);
notificationSystem.subscribeService('sms', smsHandler);
notificationSystem.subscribeService('slack', slackHandler);

// Broadcasting
await notificationSystem.broadcastNotification({
  title: 'System Maintenance',
  content: 'Scheduled maintenance tonight 2-4 AM',
  priority: 'high',
  metadata: { maintenance_type: 'database' }
});

🏷️ Headers Exchange

Le Headers Exchange route basé sur les attributs des headers plutôt que sur la routing key.

class HeadersExchangeProcessor:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        self.setup_headers_routing()
    
    def setup_headers_routing(self):
        # Headers exchange
        self.channel.exchange_declare(
            exchange='document_processor',
            exchange_type='headers',
            durable=True
        )
        
        # Queues avec matching sur headers
        processor_queues = [
            # PDF processor - format=pdf ET type=document
            {
                'queue': 'processor.pdf',
                'headers': {'format': 'pdf', 'type': 'document'},
                'match': 'all'
            },
            # Image processor - format=jpg OU format=png
            {
                'queue': 'processor.image',
                'headers': {'format': 'jpg'},
                'match': 'any'
            },
            {
                'queue': 'processor.image_png',
                'headers': {'format': 'png'},
                'match': 'any'  
            },
            # Urgent processor - urgent=true (peu importe le format)
            {
                'queue': 'processor.urgent',
                'headers': {'urgent': 'true'},
                'match': 'any'
            },
            # Archive processor - tous les documents anciens
            {
                'queue': 'processor.archive',
                'headers': {'age': 'old', 'type': 'document'},
                'match': 'all'
            }
        ]
        
        for processor in processor_queues:
            # Créer la queue
            self.channel.queue_declare(
                queue=processor['queue'],
                durable=True
            )
            
            # Bind avec headers matching
            bind_headers = {
                'x-match': processor['match'],  # 'all' ou 'any'
                **processor['headers']
            }
            
            self.channel.queue_bind(
                exchange='document_processor',
                queue=processor['queue'],
                routing_key='',  # Ignoré pour headers exchange
                arguments=bind_headers
            )
    
    def submit_document(self, document):
        """Soumettre un document pour traitement"""
        
        headers = {
            'format': document['format'].lower(),
            'type': document['type'],
            'size': str(document['size']),
            'urgent': 'true' if document.get('urgent') else 'false'
        }
        
        # Ajouter header age si document ancien
        if document.get('created_days_ago', 0) > 30:
            headers['age'] = 'old'
        
        message = {
            'document_id': document['id'],
            'url': document['url'],
            'metadata': document.get('metadata', {})
        }
        
        # Publier avec headers (routing key vide)
        self.channel.basic_publish(
            exchange='document_processor',
            routing_key='',  # Ignoré
            body=json.dumps(message),
            properties=pika.BasicProperties(
                headers=headers,
                delivery_mode=2
            )
        )
        
        print(f"Document submitted: {document['id']} with headers: {headers}")

# Exemple d'usage
processor = HeadersExchangeProcessor()

# Ces documents iront vers différents processors
documents = [
    {
        'id': 'doc1',
        'format': 'PDF',
        'type': 'document',
        'size': 1024000,
        'url': 'http://example.com/doc1.pdf'
    },
    {
        'id': 'img1',
        'format': 'jpg',
        'type': 'image',
        'size': 512000,
        'urgent': True,
        'url': 'http://example.com/img1.jpg'
    },
    {
        'id': 'old_doc',
        'format': 'pdf',
        'type': 'document',
        'size': 2048000,
        'created_days_ago': 45,
        'url': 'http://example.com/old.pdf'
    }
]

for doc in documents:
    processor.submit_document(doc)

Matching Logic

# Exemples de matching
matching_examples = [
    {
        'message_headers': {'format': 'pdf', 'type': 'document', 'urgent': 'true'},
        'binding_headers': {'x-match': 'all', 'format': 'pdf', 'type': 'document'},
        'matches': True  # Tous les headers requis sont présents
    },
    {
        'message_headers': {'format': 'jpg', 'size': 'large'},
        'binding_headers': {'x-match': 'any', 'format': 'jpg', 'format': 'png'},
        'matches': True  # Au moins un header match
    },
    {
        'message_headers': {'format': 'doc', 'urgent': 'false'},
        'binding_headers': {'x-match': 'all', 'format': 'pdf', 'urgent': 'true'},
        'matches': False  # format et urgent ne matchent pas
    }
]

🎪 Exchange par Défaut

RabbitMQ fournit un exchange par défaut (empty string) qui route directement vers les queues.

// Default exchange - routing direct vers queue
class DefaultExchangeExample {
  async simpleQueueExample() {
    const connection = await amqp.connect('amqp://localhost');
    const channel = await connection.createChannel();
    
    const queueName = 'simple_queue';
    await channel.assertQueue(queueName, { durable: true });
    
    // Publier vers l'exchange par défaut
    // routing_key = nom de la queue de destination
    await channel.sendToQueue(
      queueName,
      Buffer.from('Hello Simple Queue!')
    );
    
    // Équivalent à:
    await channel.publish(
      '',           // Exchange par défaut (empty string)
      queueName,    // Routing key = nom de queue
      Buffer.from('Hello Simple Queue!')
    );
    
    console.log(`Message sent to ${queueName}`);
    await connection.close();
  }
}

🔧 Exchange Alternatif

Configuration d'exchanges de fallback :

class AlternateExchangeSetup:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
    
    def setup_with_alternate(self):
        # Exchange de fallback pour messages non-routables
        self.channel.exchange_declare(
            exchange='unrouted_messages',
            exchange_type='fanout',
            durable=True
        )
        
        # Queue pour messages non-routés
        self.channel.queue_declare(
            queue='unrouted.messages',
            durable=True
        )
        
        self.channel.queue_bind(
            exchange='unrouted_messages',
            queue='unrouted.messages'
        )
        
        # Exchange principal avec alternate
        self.channel.exchange_declare(
            exchange='main_exchange',
            exchange_type='direct',
            durable=True,
            arguments={
                'alternate-exchange': 'unrouted_messages'
            }
        )
        
        # Queue normale
        self.channel.queue_declare(queue='normal_queue', durable=True)
        self.channel.queue_bind(
            exchange='main_exchange',
            queue='normal_queue',
            routing_key='valid'
        )
    
    def test_routing(self):
        # Message routable
        self.channel.basic_publish(
            exchange='main_exchange',
            routing_key='valid',  # Ira vers normal_queue
            body='Routable message'
        )
        
        # Message non-routable
        self.channel.basic_publish(
            exchange='main_exchange',
            routing_key='invalid',  # Ira vers alternate exchange
            body='Unroutable message'
        )

🏗️ Exchange Hierarchy

Création d'une hiérarchie complexe d'exchanges :

class ExchangeHierarchy:
    def __init__(self):
        self.connection = pika.BlockingConnection(
            pika.ConnectionParameters('localhost')
        )
        self.channel = self.connection.channel()
        self.setup_hierarchy()
    
    def setup_hierarchy(self):
        """
        Hiérarchie: Main → Regional → Service-specific
        """
        
        # Level 1: Main entry point
        self.channel.exchange_declare(
            exchange='main.events',
            exchange_type='topic',
            durable=True
        )
        
        # Level 2: Regional routing
        regions = ['us', 'eu', 'asia']
        for region in regions:
            self.channel.exchange_declare(
                exchange=f'regional.{region}',
                exchange_type='topic',
                durable=True
            )
            
            # Route from main to regional
            self.channel.exchange_bind(
                destination=f'regional.{region}',
                source='main.events',
                routing_key=f'*.{region}.*'
            )
        
        # Level 3: Service-specific exchanges
        services = ['orders', 'payments', 'users']
        
        for region in regions:
            for service in services:
                exchange_name = f'{service}.{region}'
                
                self.channel.exchange_declare(
                    exchange=exchange_name,
                    exchange_type='direct',
                    durable=True
                )
                
                # Route from regional to service
                self.channel.exchange_bind(
                    destination=exchange_name,
                    source=f'regional.{region}',
                    routing_key=f'{service}.{region}.*'
                )
        
        # Final queues
        for region in regions:
            for service in services:
                queue_name = f'{service}.{region}.processor'
                
                self.channel.queue_declare(
                    queue=queue_name,
                    durable=True
                )
                
                self.channel.queue_bind(
                    exchange=f'{service}.{region}',
                    queue=queue_name,
                    routing_key=f'{service}.process'
                )
    
    def publish_hierarchical(self, service, region, action, data):
        """Publier dans la hiérarchie"""
        routing_key = f'{service}.{region}.{action}'
        
        self.channel.basic_publish(
            exchange='main.events',  # Point d'entrée
            routing_key=routing_key,
            body=json.dumps(data)
        )
        
        print(f"Published to hierarchy: {routing_key}")

# Test de la hiérarchie
hierarchy = ExchangeHierarchy()

# Ces messages suivront la hiérarchie d'exchanges
hierarchy.publish_hierarchical('orders', 'eu', 'process', {'order_id': 123})
hierarchy.publish_hierarchical('payments', 'us', 'validate', {'payment_id': 456})

📊 Comparaison et Choix d'Exchange

Rendu du diagramme en cours...

Matrice de Décision

| Criteria | Direct | Topic | Fanout | Headers | |----------|--------|--------|--------|---------| | Complexité | Faible | Moyenne | Très faible | Élevée | | Performance | Élevée | Bonne | Très élevée | Moyenne | | Flexibilité | Faible | Élevée | Nulle | Très élevée | | Use Cases | Commands | Events | Broadcast | Content filtering |

✅ Bonnes Pratiques par Exchange

Direct Exchange

# DO
routing_keys = ['user.create', 'user.update', 'user.delete']

# DON'T 
routing_keys = ['userCreate', 'user_update', 'USER-DELETE']  # Inconsistent

Topic Exchange

// DO - Hiérarchie claire
'entity.region.action'
'order.eu.created'
'payment.us.failed'

// DON'T - Trop granulaire
'order.eu.paris.district15.created'

Fanout Exchange

# DO - Pour broadcast réel
scenarios = ['cache_invalidation', 'system_announcements', 'real_time_updates']

# DON'T - Pour routing sélectif
# Utiliser Topic ou Direct à la place

Headers Exchange

// DO - Pour filtrage complexe
headers = {
  content_type: 'application/pdf',
  file_size: 'large',
  processing_priority: 'high',
  geographic_region: 'GDPR'
};

// DON'T - Pour routing simple
// headers = { 'simple': 'routing' }  // Utiliser Direct

🎯 Exercice Pratique

Implémentez un système de processing de fichiers multimedia :

# TODO: Créer un système avec :
# 1. Headers exchange pour router par type/format
# 2. Topic exchange pour événements processing
# 3. Fanout pour notifications de completion
# 4. Direct pour commands spécialisées
# 5. Tests avec différents types de fichiers

Maîtriser les différents types d'exchanges vous permet de construire des architectures de messaging sophistiquées et performantes.

📝 Testez vos connaissances !

Répondez à 10 questions pour valider ce cours