📢 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

Rendu du diagramme en cours...

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 !

📝 Testez vos connaissances !

Répondez à 10 questions pour valider ce cours