🔀 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
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
📻 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
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.