📢 Publish/Subscribe (Pub/Sub)
Le pattern Publish/Subscribe est l'architecture fondamentale des systèmes événementiels modernes. Il transforme radicalement la façon dont nous concevons les interactions entre services, passant d'un modèle de commande direct à un modèle de diffusion d'événements qui favorise l'évolution architecturale et la réactivité.
🧠 Théorie des Systèmes Événementiels
Paradigme de l'Event-Driven Architecture
Le Pub/Sub n'est pas simplement un pattern de messagerie, c'est la fondation conceptuelle de l'Event-Driven Architecture (EDA). Dans ce paradigme, l'état du système est représenté par une séquence d'événements immuables plutôt que par des snapshots mutables.
Inversion du contrôle : Au lieu que le producteur décide qui doit être informé d'un changement, il se contente de déclarer que le changement a eu lieu. Les consommateurs décident eux-mêmes s'ils sont intéressés par ce type d'événement.
Émergence comportementale : Des comportements complexes émergent de l'interaction de règles simples. Un événement "OrderPlaced" peut déclencher automatiquement des processus de paiement, de mise à jour d'inventaire, de notification client, d'analytics, sans que le service de commande ait conscience de ces dépendances.
Théorie de l'Observateur Distribué
Pub/Sub généralise le pattern Observer à l'échelle distribuée. Chaque consommateur devient un observateur qui peut s'abonner et se désabonner dynamiquement aux événements qui l'intéressent. Cette généralisation permet :
Évolution en runtime : De nouveaux observateurs peuvent apparaître sans redéploiement Résilience : La panne d'un observateur n'affecte pas les autres Testabilité : Chaque observateur peut être testé indépendamment Auditabilité : Tous les événements métier passent par le même canal central
Le pattern Publish/Subscribe permet de diffuser un message à plusieurs consumers simultanément, créant un système nerveux distribué où l'information circule naturellement vers tous les composants intéressés. C'est l'inverse des Work Queues où un message n'est consommé que par un seul worker.
🎯 Concept Principal
Un exchange de type fanout route les messages vers toutes les queues qui lui sont liées.
🔄 Types d'Exchanges
1. Fanout Exchange
Diffuse vers toutes les queues liées :
// Publisher
channel.assertExchange('news', 'fanout', { durable: false });
channel.publish('news', '', Buffer.from('Breaking news!'));
2. Direct Exchange
Route selon une routing key exacte :
// Publisher avec routing key
channel.assertExchange('logs', 'direct', { durable: false });
channel.publish('logs', 'error', Buffer.from('Error message'));
channel.publish('logs', 'info', Buffer.from('Info message'));
3. Topic Exchange
Route selon des patterns de routing key :
// Publisher avec patterns
channel.assertExchange('events', 'topic', { durable: false });
channel.publish('events', 'user.login.success', Buffer.from(data));
channel.publish('events', 'user.register.failed', Buffer.from(data));
🛠️ Implémentation Pub/Sub Simple
Publisher (Émet des événements)
const amqp = require('amqplib/callback_api');
function publishNews(message) {
amqp.connect('amqp://localhost', function(error0, connection) {
if (error0) throw error0;
connection.createChannel(function(error1, channel) {
if (error1) throw error1;
const exchange = 'news_broadcast';
channel.assertExchange(exchange, 'fanout', {
durable: false
});
channel.publish(exchange, '', Buffer.from(message));
console.log(`[x] Actualité diffusée: ${message}`);
});
setTimeout(() => connection.close(), 500);
});
}
// Publication d'actualités
publishNews('Nouvelle version de RabbitMQ disponible !');
publishNews('Maintenance planifiée ce soir 22h-23h');
Subscriber (Reçoit tous les événements)
const amqp = require('amqplib/callback_api');
function subscribeToNews(consumerName) {
amqp.connect('amqp://localhost', function(error0, connection) {
if (error0) throw error0;
connection.createChannel(function(error1, channel) {
if (error1) throw error1;
const exchange = 'news_broadcast';
channel.assertExchange(exchange, 'fanout', {
durable: false
});
// Queue temporaire unique pour chaque consumer
channel.assertQueue('', {
exclusive: true
}, function(error2, q) {
if (error2) throw error2;
console.log(`[*] ${consumerName} en attente d'actualités`);
// Liaison queue -> exchange
channel.bindQueue(q.queue, exchange, '');
channel.consume(q.queue, function(msg) {
if (msg.content) {
console.log(`[${consumerName}] Reçu: ${msg.content.toString()}`);
}
}, {
noAck: true
});
});
});
});
}
// Plusieurs subscribers
subscribeToNews('Mobile App');
subscribeToNews('Web Dashboard');
subscribeToNews('Email Service');
🎯 Routing avec Direct Exchange
Publisher avec niveaux de logs
function publishLog(level, message) {
const exchange = 'direct_logs';
channel.assertExchange(exchange, 'direct', { durable: false });
channel.publish(exchange, level, Buffer.from(message));
console.log(`[x] Log ${level}: ${message}`);
}
publishLog('info', 'Application démarrée');
publishLog('error', 'Erreur de connexion base de données');
publishLog('warning', 'Mémoire faible');
Subscriber pour niveaux spécifiques
function subscribeToLogs(severities, consumerName) {
const exchange = 'direct_logs';
channel.assertExchange(exchange, 'direct', { durable: false });
channel.assertQueue('', { exclusive: true }, function(error2, q) {
console.log(`[*] ${consumerName} écoute: ${severities.join(', ')}`);
// Liaison pour chaque niveau de sévérité
severities.forEach(function(severity) {
channel.bindQueue(q.queue, exchange, severity);
});
channel.consume(q.queue, function(msg) {
console.log(`[${consumerName}] ${msg.fields.routingKey}: ${msg.content.toString()}`);
}, { noAck: true });
});
}
// Différents subscribers pour différents logs
subscribeToLogs(['error'], 'Error Handler');
subscribeToLogs(['info', 'warning'], 'General Monitor');
subscribeToLogs(['error', 'warning'], 'Alert System');
🏷️ Topic Exchange et Patterns
Publisher avec topics hiérarchiques
function publishEvent(category, action, status, data) {
const exchange = 'topic_events';
const routingKey = `${category}.${action}.${status}`;
channel.assertExchange(exchange, 'topic', { durable: false });
channel.publish(exchange, routingKey, Buffer.from(JSON.stringify(data)));
console.log(`[x] Événement: ${routingKey}`);
}
// Événements utilisateur
publishEvent('user', 'login', 'success', { userId: 123 });
publishEvent('user', 'register', 'failed', { reason: 'email_exists' });
// Événements système
publishEvent('system', 'backup', 'completed', { size: '2.5GB' });
publishEvent('system', 'update', 'started', { version: '1.2.0' });
Subscriber avec patterns
function subscribeToTopics(patterns, consumerName) {
const exchange = 'topic_events';
channel.assertExchange(exchange, 'topic', { durable: false });
channel.assertQueue('', { exclusive: true }, function(error2, q) {
console.log(`[*] ${consumerName} patterns: ${patterns.join(', ')}`);
patterns.forEach(function(pattern) {
channel.bindQueue(q.queue, exchange, pattern);
});
channel.consume(q.queue, function(msg) {
console.log(`[${consumerName}] ${msg.fields.routingKey}: ${msg.content.toString()}`);
}, { noAck: true });
});
}
// Différents abonnements
subscribeToTopics(['user.*'], 'User Service'); // Tous événements user
subscribeToTopics(['*.login.*'], 'Login Monitor'); // Tous les logins
subscribeToTopics(['user.*.failed'], 'Failure Tracker'); // Tous échecs user
subscribeToTopics(['#'], 'Full Logger'); // TOUS les événements
Patterns de Routing Key
| Pattern | Signification | Exemples |
|---------|---------------|----------|
| * | Un seul mot | user.* → user.login, user.logout |
| # | Zero ou plus de mots | user.# → user, user.login.success |
| exact | Correspondance exacte | user.login → seulement user.login |
🔧 Exemple Complet : Système de Notifications
Service de notifications
class NotificationService {
constructor() {
this.exchange = 'notifications';
this.setupConnection();
}
async setupConnection() {
this.connection = await amqp.connect('amqp://localhost');
this.channel = await this.connection.createChannel();
await this.channel.assertExchange(this.exchange, 'topic', {
durable: true
});
}
async sendUserNotification(userId, type, data) {
const routingKey = `user.${userId}.${type}`;
const message = JSON.stringify({ userId, type, data, timestamp: new Date() });
this.channel.publish(this.exchange, routingKey, Buffer.from(message), {
persistent: true
});
console.log(`[x] Notification envoyée: ${routingKey}`);
}
async sendGlobalNotification(type, data) {
const routingKey = `global.${type}`;
const message = JSON.stringify({ type, data, timestamp: new Date() });
this.channel.publish(this.exchange, routingKey, Buffer.from(message), {
persistent: true
});
}
}
const notificationService = new NotificationService();
// Notifications personnalisées
notificationService.sendUserNotification(123, 'order_shipped', {
orderId: 'ORD-001',
trackingNumber: 'TN123456'
});
// Notifications globales
notificationService.sendGlobalNotification('maintenance', {
start: '2024-01-15T22:00:00Z',
duration: '1h'
});
Service Email
function setupEmailService() {
const exchange = 'notifications';
channel.assertExchange(exchange, 'topic', { durable: true });
channel.assertQueue('email_notifications', { durable: true }, function(error, q) {
// S'abonner aux notifications qui nécessitent un email
const patterns = ['user.*.order_shipped', 'user.*.password_reset', 'global.maintenance'];
patterns.forEach(pattern => {
channel.bindQueue(q.queue, exchange, pattern);
});
channel.consume(q.queue, async function(msg) {
const notification = JSON.parse(msg.content.toString());
const routingKey = msg.fields.routingKey;
console.log(`[Email] Traitement: ${routingKey}`);
try {
await sendEmail(notification);
channel.ack(msg);
} catch (error) {
console.error('Erreur envoi email:', error);
channel.nack(msg, false, true); // Retry
}
}, { noAck: false });
});
}
📱 Avantages du Pub/Sub
- Découplage : Publishers et subscribers ne se connaissent pas
- Scalabilité : Ajout facile de nouveaux subscribers
- Flexibilité : Routing sophistiqué avec topics
- Résilience : Panne d'un subscriber n'affecte pas les autres
🎪 Cas d'Usage Avancés
Event Sourcing
publishEvent('order', 'created', 'success', orderData);
publishEvent('order', 'paid', 'success', paymentData);
publishEvent('order', 'shipped', 'success', shippingData);
Microservices Communication
// Service A notifie Service B et C
publishEvent('inventory', 'stock_updated', 'success', { productId, newStock });
Real-time Updates
// Mise à jour en temps réel des interfaces
publishEvent('dashboard', 'metric_updated', 'info', metricsData);
Le pattern Pub/Sub est essentiel pour construire des architectures event-driven robustes et scalables !