💼 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.