🔗 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
🛠️ 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 !