🔗 Intégration Microservices

L'intégration de RabbitMQ dans une architecture microservices représente bien plus qu'une simple couche de communication : c'est l'implémentation d'un système nerveux distribué qui permet aux services de collaborer de manière autonome et résiliente. Cette intégration transforme une collection de services isolés en un écosystème cohérent et évolutif.

🧠 Théorie des Systèmes Distribués Cohérents

Cohérence Éventuelle vs Cohérence Forte

Dans une architecture microservices traditionnelle avec bases de données partagées, maintenir la cohérence forte nécessite des transactions distribuées (2PC) qui sont coûteuses et fragiles. RabbitMQ permet d'implémenter la cohérence éventuelle (eventual consistency) qui est souvent suffisante pour la plupart des cas métier.

Pattern Saga : Au lieu d'une transaction globale, chaque service effectue sa transaction locale puis publie un événement. Les autres services réagissent en chaîne, créant une "saga" de transactions coordonnées mais autonomes. En cas d'échec, des événements de compensation permettent de revenir à un état cohérent.

Event Sourcing distribué : Chaque microservice devient un domaine métier autonome qui maintient son propre event store. Les événements cross-domain passent par RabbitMQ, créant un journal distribué des interactions inter-services.

Anti-Patterns à Éviter

Microservices trop bavards : Ne transformez pas RabbitMQ en un système de RPC déguisé. Les interactions synchrones fréquentes indiquent souvent une mauvaise découpe de domaines.

Event storms : Évitez les chaînes d'événements qui s'auto-alimentent. Un événement ne doit pas créer un nouveau événement qui en crée un autre, risquant des boucles infinies.

Shared databases via events : N'utilisez pas les événements pour synchroniser des bases de données partagées. Chaque service doit posséder ses propres données.

L'intégration de RabbitMQ dans une architecture microservices permet un découplage efficace et une communication asynchrone robuste entre les services, créant un système distribué qui évolue naturellement avec les besoins métier.

🏗️ Architecture Event-Driven

Rendu du diagramme en cours...

🛠️ Service Base Pattern

Base Microservice avec RabbitMQ

class BaseMicroservice {
    constructor(serviceName, exchangeName = 'microservices') {
        this.serviceName = serviceName;
        this.exchangeName = exchangeName;
        this.connection = null;
        this.channel = null;
        this.eventHandlers = new Map();
    }
    
    async initialize() {
        await this.connectRabbitMQ();
        await this.setupExchange();
        await this.setupServiceQueue();
        await this.registerEventHandlers();
        
        console.log(`✅ Service ${this.serviceName} initialisé`);
    }
    
    async connectRabbitMQ() {
        const connectionUrl = process.env.RABBITMQ_URL || 'amqp://admin:password123@localhost:5672';
        
        this.connection = await amqp.connect(connectionUrl);
        this.channel = await this.connection.createChannel();
        
        // Gestion des déconnexions
        this.connection.on('close', () => {
            console.log('🔌 Connexion RabbitMQ fermée - reconnexion...');
            this.reconnect();
        });
    }
    
    async setupExchange() {
        await this.channel.assertExchange(this.exchangeName, 'topic', {
            durable: true
        });
    }
    
    async setupServiceQueue() {
        const queueName = `${this.serviceName}_events`;
        
        await this.channel.assertQueue(queueName, {
            durable: true,
            arguments: {
                'x-dead-letter-exchange': 'dlx',
                'x-dead-letter-routing-key': 'failed'
            }
        });
        
        this.serviceQueue = queueName;
    }
    
    async publishEvent(eventType, data, routingKey = null) {
        const event = {
            id: this.generateEventId(),
            type: eventType,
            service: this.serviceName,
            timestamp: new Date().toISOString(),
            data: data
        };
        
        const finalRoutingKey = routingKey || `${this.serviceName}.${eventType}`;
        
        await this.channel.publish(
            this.exchangeName,
            finalRoutingKey,
            Buffer.from(JSON.stringify(event)),
            { 
                persistent: true,
                messageId: event.id
            }
        );
        
        console.log(`📤 Événement publié [${finalRoutingKey}]: ${event.id}`);
        return event;
    }
    
    async subscribeToEvents(patterns, handler) {
        for (const pattern of patterns) {
            await this.channel.bindQueue(this.serviceQueue, this.exchangeName, pattern);
            console.log(`📥 Abonné au pattern: ${pattern}`);
        }
        
        await this.channel.consume(this.serviceQueue, async (msg) => {
            if (msg) {
                try {
                    const event = JSON.parse(msg.content.toString());
                    const routingKey = msg.fields.routingKey;
                    
                    console.log(`📨 Événement reçu [${routingKey}]: ${event.id}`);
                    
                    await handler(event, routingKey);
                    this.channel.ack(msg);
                    
                } catch (error) {
                    console.error('❌ Erreur traitement événement:', error.message);
                    this.channel.nack(msg, false, false); // Dead letter
                }
            }
        }, { noAck: false });
    }
    
    generateEventId() {
        return `${this.serviceName}_${Date.now()}_${Math.random().toString(36).substr(2, 9)}`;
    }
    
    async reconnect() {
        try {
            await new Promise(resolve => setTimeout(resolve, 5000));
            await this.initialize();
        } catch (error) {
            console.error('❌ Reconnexion échouée:', error.message);
            setTimeout(() => this.reconnect(), 10000);
        }
    }
}

🛍️ Exemple: Service de Commandes

Order Service Implementation

class OrderService extends BaseMicroservice {
    constructor() {
        super('order-service');
        this.orderRepository = new OrderRepository();
    }
    
    async registerEventHandlers() {
        // Écouter les événements de paiement
        await this.subscribeToEvents([
            'payment-service.payment.completed',
            'payment-service.payment.failed'
        ], this.handlePaymentEvent.bind(this));
        
        // Écouter les événements d'inventaire
        await this.subscribeToEvents([
            'inventory-service.stock.reserved',
            'inventory-service.stock.insufficient'
        ], this.handleInventoryEvent.bind(this));
    }
    
    async createOrder(orderData) {
        try {
            // 1. Validation de la commande
            const validatedOrder = await this.validateOrder(orderData);
            
            // 2. Sauvegarde locale
            const order = await this.orderRepository.create(validatedOrder);
            
            // 3. Publication événement création
            await this.publishEvent('order.created', {
                orderId: order.id,
                customerId: order.customerId,
                items: order.items,
                totalAmount: order.totalAmount
            });
            
            // 4. Démarrage du workflow asynchrone
            await this.startOrderWorkflow(order);
            
            return order;
            
        } catch (error) {
            console.error('❌ Erreur création commande:', error.message);
            
            await this.publishEvent('order.creation.failed', {
                error: error.message,
                orderData: orderData
            });
            
            throw error;
        }
    }
    
    async startOrderWorkflow(order) {
        // Étape 1: Réservation stock
        await this.publishEvent('inventory.reserve.requested', {
            orderId: order.id,
            items: order.items
        }, 'inventory-service.reserve.request');
        
        // Mise à jour statut
        await this.orderRepository.updateStatus(order.id, 'awaiting_inventory');
    }
    
    async handlePaymentEvent(event, routingKey) {
        const { orderId, status, transactionId } = event.data;
        
        switch (event.type) {
            case 'payment.completed':
                await this.orderRepository.updateStatus(orderId, 'paid');
                await this.publishEvent('order.paid', {
                    orderId,
                    transactionId,
                    paidAt: event.timestamp
                });
                
                // Démarrage expédition
                await this.publishEvent('shipping.create.requested', {
                    orderId
                }, 'shipping-service.create.request');
                break;
                
            case 'payment.failed':
                await this.orderRepository.updateStatus(orderId, 'payment_failed');
                await this.publishEvent('order.payment.failed', {
                    orderId,
                    reason: event.data.error
                });
                break;
        }
    }
    
    async handleInventoryEvent(event, routingKey) {
        const { orderId } = event.data;
        
        switch (event.type) {
            case 'stock.reserved':
                await this.orderRepository.updateStatus(orderId, 'inventory_reserved');
                
                // Démarrer processus paiement
                await this.publishEvent('payment.process.requested', {
                    orderId: orderId,
                    amount: event.data.totalAmount
                }, 'payment-service.process.request');
                break;
                
            case 'stock.insufficient':
                await this.orderRepository.updateStatus(orderId, 'cancelled');
                await this.publishEvent('order.cancelled', {
                    orderId,
                    reason: 'insufficient_stock'
                });
                break;
        }
    }
}

💳 Payment Service

class PaymentService extends BaseMicroservice {
    constructor() {
        super('payment-service');
        this.paymentProcessor = new PaymentProcessor();
    }
    
    async registerEventHandlers() {
        await this.subscribeToEvents([
            'payment-service.process.request'
        ], this.handlePaymentRequest.bind(this));
    }
    
    async handlePaymentRequest(event, routingKey) {
        const { orderId, amount, paymentMethod } = event.data;
        
        try {
            console.log(`💳 Traitement paiement commande ${orderId}: ${amount}€`);
            
            // Simulation appel gateway de paiement
            const paymentResult = await this.paymentProcessor.processPayment({
                orderId,
                amount,
                method: paymentMethod
            });
            
            if (paymentResult.success) {
                await this.publishEvent('payment.completed', {
                    orderId,
                    transactionId: paymentResult.transactionId,
                    amount: amount,
                    processedAt: new Date().toISOString()
                });
            } else {
                await this.publishEvent('payment.failed', {
                    orderId,
                    error: paymentResult.error,
                    attemptedAt: new Date().toISOString()
                });
            }
            
        } catch (error) {
            console.error('❌ Erreur paiement:', error.message);
            
            await this.publishEvent('payment.failed', {
                orderId,
                error: error.message,
                systemError: true
            });
        }
    }
}

class PaymentProcessor {
    async processPayment(paymentData) {
        // Simulation traitement asynchrone
        await new Promise(resolve => setTimeout(resolve, 2000));
        
        // Simulation échec 10% du temps
        if (Math.random() < 0.1) {
            return {
                success: false,
                error: 'Payment gateway timeout'
            };
        }
        
        return {
            success: true,
            transactionId: `TXN_${Date.now()}_${Math.random().toString(36).substr(2, 9)}`
        };
    }
}

📦 Inventory Service

class InventoryService extends BaseMicroservice {
    constructor() {
        super('inventory-service');
        this.inventory = new InventoryManager();
    }
    
    async registerEventHandlers() {
        await this.subscribeToEvents([
            'inventory-service.reserve.request',
            'inventory-service.release.request'
        ], this.handleInventoryEvent.bind(this));
    }
    
    async handleInventoryEvent(event, routingKey) {
        const { orderId, items } = event.data;
        
        switch (event.type) {
            case 'reserve.requested':
                await this.handleReservationRequest(orderId, items);
                break;
                
            case 'release.requested':
                await this.handleReleaseRequest(orderId, items);
                break;
        }
    }
    
    async handleReservationRequest(orderId, items) {
        try {
            console.log(`📦 Réservation stock pour commande ${orderId}`);
            
            // Vérifier disponibilité de tous les items
            const availabilityCheck = await this.inventory.checkAvailability(items);
            
            if (availabilityCheck.available) {
                // Réserver les items
                const reservation = await this.inventory.reserveItems(orderId, items);
                
                await this.publishEvent('stock.reserved', {
                    orderId,
                    items: reservation.items,
                    reservationId: reservation.id,
                    reservedAt: new Date().toISOString()
                });
                
            } else {
                await this.publishEvent('stock.insufficient', {
                    orderId,
                    unavailableItems: availabilityCheck.unavailableItems,
                    checkedAt: new Date().toISOString()
                });
            }
            
        } catch (error) {
            console.error('❌ Erreur réservation stock:', error.message);
            
            await this.publishEvent('stock.reservation.failed', {
                orderId,
                error: error.message,
                systemError: true
            });
        }
    }
}

class InventoryManager {
    constructor() {
        this.stock = new Map([
            ['PROD001', { available: 100, reserved: 0 }],
            ['PROD002', { available: 50, reserved: 0 }],
            ['PROD003', { available: 25, reserved: 0 }]
        ]);
    }
    
    async checkAvailability(items) {
        const unavailableItems = [];
        
        for (const item of items) {
            const stockInfo = this.stock.get(item.productId);
            
            if (!stockInfo || stockInfo.available < item.quantity) {
                unavailableItems.push({
                    productId: item.productId,
                    requested: item.quantity,
                    available: stockInfo?.available || 0
                });
            }
        }
        
        return {
            available: unavailableItems.length === 0,
            unavailableItems
        };
    }
    
    async reserveItems(orderId, items) {
        const reservationId = `RES_${Date.now()}`;
        const reservedItems = [];
        
        for (const item of items) {
            const stockInfo = this.stock.get(item.productId);
            stockInfo.available -= item.quantity;
            stockInfo.reserved += item.quantity;
            
            reservedItems.push({
                ...item,
                reservationId
            });
        }
        
        return {
            id: reservationId,
            orderId,
            items: reservedItems
        };
    }
}

📧 Notification Service

class NotificationService extends BaseMicroservice {
    constructor() {
        super('notification-service');
        this.emailService = new EmailService();
        this.smsService = new SMSService();
        this.pushService = new PushService();
    }
    
    async registerEventHandlers() {
        // S'abonner à tous les événements nécessitant des notifications
        await this.subscribeToEvents([
            '*.order.created',
            '*.order.paid', 
            '*.order.shipped',
            '*.order.cancelled',
            '*.payment.failed',
            '*.user.registered'
        ], this.handleNotificationEvent.bind(this));
    }
    
    async handleNotificationEvent(event, routingKey) {
        console.log(`📨 Traitement notification: ${event.type}`);
        
        try {
            const notificationRules = await this.getNotificationRules(event.type);
            
            for (const rule of notificationRules) {
                await this.sendNotification(rule, event);
            }
            
        } catch (error) {
            console.error('❌ Erreur notification:', error.message);
        }
    }
    
    async getNotificationRules(eventType) {
        const rules = {
            'order.created': [
                { type: 'email', template: 'order_confirmation', priority: 'high' },
                { type: 'push', template: 'order_created', priority: 'normal' }
            ],
            'order.shipped': [
                { type: 'email', template: 'shipping_notification', priority: 'normal' },
                { type: 'sms', template: 'shipping_sms', priority: 'low' }
            ],
            'payment.failed': [
                { type: 'email', template: 'payment_failed', priority: 'high' },
                { type: 'push', template: 'payment_issue', priority: 'high' }
            ]
        };
        
        return rules[eventType] || [];
    }
    
    async sendNotification(rule, event) {
        const notification = {
            id: this.generateEventId(),
            type: rule.type,
            template: rule.template,
            priority: rule.priority,
            recipient: await this.getRecipient(event),
            data: this.prepareNotificationData(event),
            createdAt: new Date().toISOString()
        };
        
        switch (rule.type) {
            case 'email':
                await this.emailService.send(notification);
                break;
            case 'sms':
                await this.smsService.send(notification);
                break;
            case 'push':
                await this.pushService.send(notification);
                break;
        }
        
        // Log pour audit
        await this.publishEvent('notification.sent', {
            notificationId: notification.id,
            type: rule.type,
            template: rule.template,
            recipient: notification.recipient,
            originalEvent: event.id
        });
    }
}

🔄 Saga Pattern Implementation

Order Processing Saga

class OrderProcessingSaga extends BaseMicroservice {
    constructor() {
        super('order-saga');
        this.activeOrders = new Map();
    }
    
    async registerEventHandlers() {
        await this.subscribeToEvents([
            'order-service.order.created',
            'inventory-service.stock.*',
            'payment-service.payment.*', 
            'shipping-service.shipment.*'
        ], this.handleSagaEvent.bind(this));
    }
    
    async handleSagaEvent(event, routingKey) {
        const orderId = event.data.orderId;
        let orderSaga = this.activeOrders.get(orderId);
        
        if (!orderSaga && event.type === 'order.created') {
            orderSaga = this.createOrderSaga(orderId, event.data);
            this.activeOrders.set(orderId, orderSaga);
        }
        
        if (orderSaga) {
            await this.processNextStep(orderSaga, event);
        }
    }
    
    createOrderSaga(orderId, orderData) {
        return {
            orderId,
            status: 'started',
            steps: [
                { name: 'inventory_check', status: 'pending', service: 'inventory-service' },
                { name: 'payment_processing', status: 'pending', service: 'payment-service' },
                { name: 'shipping_creation', status: 'pending', service: 'shipping-service' }
            ],
            orderData,
            createdAt: new Date().toISOString(),
            compensations: [] // Actions de rollback
        };
    }
    
    async processNextStep(saga, event) {
        const currentStep = saga.steps.find(s => s.status === 'pending');
        
        if (!currentStep) {
            // Saga terminée
            await this.completeSaga(saga);
            return;
        }
        
        switch (currentStep.name) {
            case 'inventory_check':
                await this.handleInventoryStep(saga, event);
                break;
            case 'payment_processing':
                await this.handlePaymentStep(saga, event);
                break;
            case 'shipping_creation':
                await this.handleShippingStep(saga, event);
                break;
        }
    }
    
    async handleInventoryStep(saga, event) {
        if (event.type === 'stock.reserved') {
            saga.steps[0].status = 'completed';
            saga.compensations.push({
                action: 'release_inventory',
                data: { orderId: saga.orderId, reservationId: event.data.reservationId }
            });
            
            // Démarrer étape paiement
            await this.publishEvent('payment.process.requested', {
                orderId: saga.orderId,
                amount: saga.orderData.totalAmount
            }, 'payment-service.process.request');
            
        } else if (event.type === 'stock.insufficient') {
            await this.compensateSaga(saga, 'insufficient_inventory');
        }
    }
    
    async handlePaymentStep(saga, event) {
        if (event.type === 'payment.completed') {
            saga.steps[1].status = 'completed';
            
            // Démarrer expédition
            await this.publishEvent('shipment.create.requested', {
                orderId: saga.orderId
            }, 'shipping-service.create.request');
            
        } else if (event.type === 'payment.failed') {
            await this.compensateSaga(saga, 'payment_failed');
        }
    }
    
    async compensateSaga(saga, reason) {
        console.log(`🔄 Compensation saga ${saga.orderId}: ${reason}`);
        
        // Exécuter toutes les actions de compensation dans l'ordre inverse
        for (const compensation of saga.compensations.reverse()) {
            try {
                await this.executeCompensation(compensation);
            } catch (error) {
                console.error('❌ Erreur compensation:', error.message);
            }
        }
        
        saga.status = 'compensated';
        
        await this.publishEvent('order.cancelled', {
            orderId: saga.orderId,
            reason,
            compensatedAt: new Date().toISOString()
        });
        
        this.activeOrders.delete(saga.orderId);
    }
    
    async executeCompensation(compensation) {
        switch (compensation.action) {
            case 'release_inventory':
                await this.publishEvent('inventory.release.requested', 
                    compensation.data, 
                    'inventory-service.release.request'
                );
                break;
                
            case 'refund_payment':
                await this.publishEvent('payment.refund.requested',
                    compensation.data,
                    'payment-service.refund.request'
                );
                break;
        }
    }
}

🌐 API Gateway Integration

Event-Driven API Gateway

class EventDrivenGateway {
    constructor() {
        this.rabbitClient = new BaseMicroservice('api-gateway');
        this.pendingRequests = new Map();
        this.setupRoutes();
    }
    
    async initialize() {
        await this.rabbitClient.initialize();
        
        // S'abonner aux réponses des services
        await this.rabbitClient.subscribeToEvents([
            '*.response.*',
            '*.result.*'
        ], this.handleServiceResponse.bind(this));
    }
    
    setupRoutes() {
        const express = require('express');
        const app = express();
        
        app.use(express.json());
        
        // Route asynchrone pour créer une commande
        app.post('/api/orders', async (req, res) => {
            try {
                const requestId = this.generateRequestId();
                
                // Publier événement via RabbitMQ
                await this.rabbitClient.publishEvent('order.create.requested', {
                    requestId,
                    ...req.body
                }, 'order-service.create.request');
                
                // Stocker la requête en attente
                this.pendingRequests.set(requestId, {
                    res,
                    timeout: setTimeout(() => {
                        this.pendingRequests.delete(requestId);
                        res.status(408).json({ error: 'Request timeout' });
                    }, 30000) // 30 secondes timeout
                });
                
            } catch (error) {
                res.status(500).json({ error: error.message });
            }
        });
        
        // Route pour status de commande (temps réel via WebSocket)
        app.get('/api/orders/:orderId/status', async (req, res) => {
            const orderId = req.params.orderId;
            
            // Demander le statut via RabbitMQ
            await this.rabbitClient.publishEvent('order.status.requested', {
                orderId
            }, 'order-service.status.request');
            
            // WebSocket pour réponse temps réel
            res.json({ message: 'Status request sent, check WebSocket for updates' });
        });
        
        this.app = app;
    }
    
    async handleServiceResponse(event, routingKey) {
        const requestId = event.data.requestId;
        const pendingRequest = this.pendingRequests.get(requestId);
        
        if (pendingRequest) {
            clearTimeout(pendingRequest.timeout);
            
            // Répondre à la requête HTTP originale
            if (event.type.includes('success') || event.type.includes('completed')) {
                pendingRequest.res.json({
                    success: true,
                    data: event.data
                });
            } else {
                pendingRequest.res.status(400).json({
                    success: false,
                    error: event.data.error || 'Service error'
                });
            }
            
            this.pendingRequests.delete(requestId);
        }
    }
}

📊 Service Mesh Integration

Service Discovery via RabbitMQ

class ServiceDiscovery {
    constructor() {
        this.services = new Map();
        this.healthChecks = new Map();
        this.rabbitClient = new BaseMicroservice('service-discovery');
    }
    
    async initialize() {
        await this.rabbitClient.initialize();
        
        // Écouter les événements de lifecycle des services
        await this.rabbitClient.subscribeToEvents([
            '*.service.started',
            '*.service.stopped',
            '*.service.health.*'
        ], this.handleServiceEvent.bind(this));
        
        // Health check périodique
        setInterval(() => {
            this.performHealthChecks();
        }, 30000);
    }
    
    async registerService(serviceInfo) {
        const service = {
            ...serviceInfo,
            registeredAt: new Date().toISOString(),
            lastSeen: new Date().toISOString(),
            status: 'healthy'
        };
        
        this.services.set(serviceInfo.name, service);
        
        // Publier événement registration
        await this.rabbitClient.publishEvent('service.registered', service);
        
        console.log(`✅ Service registré: ${serviceInfo.name}`);
    }
    
    async handleServiceEvent(event, routingKey) {
        const serviceName = event.service || event.data.serviceName;
        
        switch (event.type) {
            case 'service.started':
                await this.registerService(event.data);
                break;
                
            case 'service.stopped':
                this.services.delete(serviceName);
                await this.rabbitClient.publishEvent('service.deregistered', { serviceName });
                break;
                
            case 'service.health.ok':
                if (this.services.has(serviceName)) {
                    this.services.get(serviceName).lastSeen = new Date().toISOString();
                    this.services.get(serviceName).status = 'healthy';
                }
                break;
        }
    }
    
    async performHealthChecks() {
        console.log('🏥 Health checks des services...');
        
        for (const [serviceName, service] of this.services) {
            const lastSeen = new Date(service.lastSeen);
            const timeSinceLastSeen = Date.now() - lastSeen.getTime();
            
            if (timeSinceLastSeen > 60000) { // 1 minute
                service.status = 'unhealthy';
                
                await this.rabbitClient.publishEvent('service.unhealthy', {
                    serviceName,
                    lastSeen: service.lastSeen,
                    timeSinceLastSeen
                });
            }
        }
    }
    
    getHealthyServices() {
        return Array.from(this.services.values())
            .filter(service => service.status === 'healthy');
    }
}

🔄 Circuit Breaker Pattern

Circuit Breaker pour Services

class ServiceCircuitBreaker {
    constructor(serviceName, thresholds = {}) {
        this.serviceName = serviceName;
        this.state = 'CLOSED'; // CLOSED, OPEN, HALF_OPEN
        this.failures = 0;
        this.lastFailTime = 0;
        
        this.thresholds = {
            failureThreshold: 5,
            timeout: 60000,
            ...thresholds
        };
        
        this.successCount = 0;
        this.requestCount = 0;
    }
    
    async callService(serviceMethod, data) {
        if (this.state === 'OPEN') {
            if (Date.now() - this.lastFailTime > this.thresholds.timeout) {
                this.state = 'HALF_OPEN';
                console.log(`🔄 Circuit ${this.serviceName}: HALF_OPEN`);
            } else {
                throw new Error(`Circuit breaker OPEN pour ${this.serviceName}`);
            }
        }
        
        try {
            this.requestCount++;
            const result = await this.executeServiceCall(serviceMethod, data);
            
            this.onSuccess();
            return result;
            
        } catch (error) {
            this.onFailure();
            throw error;
        }
    }
    
    async executeServiceCall(method, data) {
        const requestId = this.generateRequestId();
        
        // Publication requête RPC via RabbitMQ
        await this.publishServiceRequest(method, data, requestId);
        
        // Attendre réponse avec timeout
        return new Promise((resolve, reject) => {
            const timeout = setTimeout(() => {
                this.pendingRequests.delete(requestId);
                reject(new Error('Service call timeout'));
            }, 10000);
            
            this.pendingRequests.set(requestId, {
                resolve: (result) => {
                    clearTimeout(timeout);
                    resolve(result);
                },
                reject: (error) => {
                    clearTimeout(timeout);
                    reject(error);
                }
            });
        });
    }
    
    onSuccess() {
        this.failures = 0;
        
        if (this.state === 'HALF_OPEN') {
            this.successCount++;
            
            if (this.successCount >= 3) { // 3 succès pour fermer
                this.state = 'CLOSED';
                this.successCount = 0;
                console.log(`✅ Circuit ${this.serviceName}: FERMÉ`);
            }
        }
    }
    
    onFailure() {
        this.failures++;
        this.lastFailTime = Date.now();
        
        if (this.failures >= this.thresholds.failureThreshold) {
            this.state = 'OPEN';
            console.log(`🚨 Circuit ${this.serviceName}: OUVERT`);
        }
    }
}

📈 Monitoring Microservices

Service Health Dashboard

class MicroservicesMonitor {
    constructor() {
        this.serviceMetrics = new Map();
        this.rabbitClient = new BaseMicroservice('microservices-monitor');
    }
    
    async initialize() {
        await this.rabbitClient.initialize();
        
        // Écouter tous les événements pour métriques
        await this.rabbitClient.subscribeToEvents(['#'], this.collectMetrics.bind(this));
        
        // API pour dashboard
        this.startDashboardAPI();
    }
    
    async collectMetrics(event, routingKey) {
        const serviceName = event.service;
        
        if (!this.serviceMetrics.has(serviceName)) {
            this.serviceMetrics.set(serviceName, {
                name: serviceName,
                events_processed: 0,
                last_activity: null,
                error_count: 0,
                success_count: 0,
                average_response_time: 0
            });
        }
        
        const metrics = this.serviceMetrics.get(serviceName);
        metrics.events_processed++;
        metrics.last_activity = new Date().toISOString();
        
        if (event.type.includes('failed') || event.type.includes('error')) {
            metrics.error_count++;
        } else {
            metrics.success_count++;
        }
        
        // Calcul taux de succès
        metrics.success_rate = (metrics.success_count / (metrics.success_count + metrics.error_count)) * 100;
    }
    
    startDashboardAPI() {
        const express = require('express');
        const app = express();
        
        app.get('/api/services', (req, res) => {
            const services = Array.from(this.serviceMetrics.values());
            
            res.json({
                total_services: services.length,
                healthy_services: services.filter(s => this.isServiceHealthy(s)).length,
                services: services
            });
        });
        
        app.get('/api/service/:name', (req, res) => {
            const service = this.serviceMetrics.get(req.params.name);
            
            if (service) {
                res.json(service);
            } else {
                res.status(404).json({ error: 'Service not found' });
            }
        });
        
        app.listen(3001, () => {
            console.log('📊 Microservices dashboard API running on port 3001');
        });
    }
    
    isServiceHealthy(service) {
        const lastActivity = new Date(service.last_activity);
        const timeSinceActivity = Date.now() - lastActivity.getTime();
        
        return timeSinceActivity < 300000 && service.success_rate > 95; // 5 min et >95% succès
    }
}

🔧 Configuration par Environment

Multi-Environment Setup

class EnvironmentConfig {
    static getConfig(environment) {
        const configs = {
            development: {
                rabbitmq: {
                    url: 'amqp://admin:password123@localhost:5672',
                    vhost: 'development',
                    exchange: 'dev_microservices',
                    retry_attempts: 3,
                    timeout: 5000
                },
                services: {
                    circuit_breaker: {
                        failureThreshold: 3,
                        timeout: 30000
                    }
                }
            },
            
            staging: {
                rabbitmq: {
                    url: process.env.STAGING_RABBITMQ_URL,
                    vhost: 'staging',
                    exchange: 'staging_microservices', 
                    retry_attempts: 5,
                    timeout: 10000
                },
                services: {
                    circuit_breaker: {
                        failureThreshold: 5,
                        timeout: 60000
                    }
                }
            },
            
            production: {
                rabbitmq: {
                    url: process.env.PROD_RABBITMQ_URL,
                    vhost: 'production',
                    exchange: 'prod_microservices',
                    retry_attempts: 10,
                    timeout: 15000,
                    ssl: true
                },
                services: {
                    circuit_breaker: {
                        failureThreshold: 10,
                        timeout: 120000
                    }
                }
            }
        };
        
        return configs[environment] || configs.development;
    }
}

⚡ Performance Optimization

Connection Pooling pour Microservices

class MicroserviceConnectionPool {
    constructor(serviceName, maxConnections = 5) {
        this.serviceName = serviceName;
        this.maxConnections = maxConnections;
        this.connections = [];
        this.activeConnections = 0;
        this.connectionQueue = [];
    }
    
    async getConnection() {
        if (this.connections.length > 0) {
            return this.connections.pop();
        }
        
        if (this.activeConnections < this.maxConnections) {
            return await this.createConnection();
        }
        
        // Attendre qu'une connexion se libère
        return new Promise((resolve) => {
            this.connectionQueue.push(resolve);
        });
    }
    
    async createConnection() {
        const connection = await amqp.connect(process.env.RABBITMQ_URL);
        const channel = await connection.createChannel();
        
        this.activeConnections++;
        
        connection.on('close', () => {
            this.activeConnections--;
        });
        
        return { connection, channel };
    }
    
    releaseConnection(connectionInfo) {
        if (this.connectionQueue.length > 0) {
            const waitingResolver = this.connectionQueue.shift();
            waitingResolver(connectionInfo);
        } else {
            this.connections.push(connectionInfo);
        }
    }
}

L'intégration microservices avec RabbitMQ permet de créer des systèmes distribués résilients et hautement scalables !

📝 Testez vos connaissances !

Répondez à 10 questions pour valider ce cours