🎯 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
📢 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.