🏗️ Architectures Orientées Événements
Introduction
L'Architecture Orientée Événements (Event-Driven Architecture - EDA) est un paradigme architectural où les composants du système communiquent via des événements. Cette approche permet de construire des systèmes hautement découplés, scalables et réactifs.
Concepts Fondamentaux
Qu'est-ce qu'un Événement ?
Un événement est un changement d'état significatif dans le système, capturé sous forme de message immuable contenant les détails de ce qui s'est passé.
// Exemple d'événement
const orderCreatedEvent = {
eventId: "550e8400-e29b-41d4-a716-446655440000",
eventType: "OrderCreated",
aggregateId: "order-123",
timestamp: "2024-01-15T10:30:00Z",
version: 1,
payload: {
orderId: "order-123",
customerId: "customer-456",
items: [
{ productId: "prod-789", quantity: 2, price: 29.99 }
],
totalAmount: 59.98
},
metadata: {
userId: "user-789",
correlationId: "corr-abc",
causationId: "cause-def"
}
};
🎨 Patterns Architecturaux
1. Event Sourcing
Stocker l'état comme une séquence d'événements :
class Order:
def __init__(self, order_id):
self.order_id = order_id
self.events = []
self.state = {}
def apply_event(self, event):
"""Appliquer un événement pour modifier l'état"""
if event['type'] == 'OrderCreated':
self.state = {
'id': event['payload']['orderId'],
'status': 'created',
'items': event['payload']['items'],
'total': event['payload']['total']
}
elif event['type'] == 'OrderPaid':
self.state['status'] = 'paid'
self.state['paymentId'] = event['payload']['paymentId']
elif event['type'] == 'OrderShipped':
self.state['status'] = 'shipped'
self.state['trackingNumber'] = event['payload']['trackingNumber']
self.events.append(event)
def replay_events(self, events):
"""Reconstruire l'état à partir des événements"""
for event in events:
self.apply_event(event)
return self.state
2. CQRS (Command Query Responsibility Segregation)
Séparer les opérations de lecture et d'écriture :
// Côté Command (Écriture)
class OrderCommandHandler {
async handleCreateOrder(command) {
// Validation
const order = new Order(command);
// Logique métier
await order.validate();
// Émission d'événements
const event = {
type: 'OrderCreated',
aggregateId: order.id,
payload: order.toJSON()
};
await this.eventStore.append(event);
await this.eventBus.publish(event);
return order.id;
}
}
// Côté Query (Lecture)
class OrderQueryHandler {
constructor(readDatabase) {
this.db = readDatabase;
}
async getOrderById(orderId) {
// Lecture optimisée depuis une vue matérialisée
return this.db.orders.findById(orderId);
}
async getOrdersByCustomer(customerId) {
// Requête complexe sur modèle dénormalisé
return this.db.orders.find({
customerId,
include: ['items', 'payments', 'shipments']
});
}
}
// Projection - Synchronisation Read Model
class OrderProjection {
async handleOrderCreated(event) {
await this.db.orders.insert({
id: event.payload.orderId,
customerId: event.payload.customerId,
status: 'created',
createdAt: event.timestamp
});
}
async handleOrderPaid(event) {
await this.db.orders.update(
{ id: event.aggregateId },
{ status: 'paid', paidAt: event.timestamp }
);
}
}
3. Event-Driven Microservices
Architecture de microservices communicant par événements :
# Architecture avec RabbitMQ
services:
order-service:
publishes:
- OrderCreated
- OrderCancelled
subscribes:
- PaymentProcessed
- InventoryReserved
payment-service:
publishes:
- PaymentProcessed
- PaymentFailed
subscribes:
- OrderCreated
inventory-service:
publishes:
- InventoryReserved
- InventoryReleased
subscribes:
- OrderCreated
- OrderCancelled
notification-service:
subscribes:
- OrderCreated
- PaymentProcessed
- OrderShipped
🔄 Saga Pattern
Gérer les transactions distribuées avec des sagas :
Choreography-based Saga
# Service de commande
class OrderService:
async def create_order(self, order_data):
# Créer la commande
order = Order.create(order_data)
await self.repository.save(order)
# Publier l'événement
await self.publish_event('OrderCreated', {
'orderId': order.id,
'items': order.items,
'customerId': order.customer_id
})
async def handle_payment_failed(self, event):
# Compensation : annuler la commande
order = await self.repository.get(event['orderId'])
order.cancel()
await self.repository.save(order)
await self.publish_event('OrderCancelled', {
'orderId': order.id
})
# Service de paiement
class PaymentService:
async def handle_order_created(self, event):
try:
payment = await self.process_payment(
event['customerId'],
event['totalAmount']
)
await self.publish_event('PaymentProcessed', {
'orderId': event['orderId'],
'paymentId': payment.id
})
except PaymentError as e:
await self.publish_event('PaymentFailed', {
'orderId': event['orderId'],
'reason': str(e)
})
# Service d'inventaire
class InventoryService:
async def handle_payment_processed(self, event):
try:
reservation = await self.reserve_items(event['items'])
await self.publish_event('InventoryReserved', {
'orderId': event['orderId'],
'reservationId': reservation.id
})
except InsufficientStock as e:
# Déclencher la compensation
await self.publish_event('InventoryFailed', {
'orderId': event['orderId'],
'reason': str(e)
})
async def handle_order_cancelled(self, event):
# Compensation : libérer l'inventaire
await self.release_reservation(event['orderId'])
Orchestration-based Saga
class OrderSagaOrchestrator {
async executeOrderSaga(orderData) {
const sagaId = generateId();
const saga = {
id: sagaId,
state: 'STARTED',
orderData,
compensations: []
};
try {
// Étape 1: Créer la commande
const order = await this.orderService.createOrder(orderData);
saga.orderId = order.id;
saga.compensations.push(() => this.orderService.cancelOrder(order.id));
// Étape 2: Traiter le paiement
const payment = await this.paymentService.processPayment({
orderId: order.id,
amount: orderData.total
});
saga.paymentId = payment.id;
saga.compensations.push(() => this.paymentService.refund(payment.id));
// Étape 3: Réserver l'inventaire
const reservation = await this.inventoryService.reserve({
orderId: order.id,
items: orderData.items
});
saga.reservationId = reservation.id;
saga.compensations.push(() => this.inventoryService.release(reservation.id));
// Étape 4: Créer l'expédition
const shipment = await this.shippingService.createShipment({
orderId: order.id,
address: orderData.shippingAddress
});
saga.state = 'COMPLETED';
await this.saveSaga(saga);
return { success: true, orderId: order.id };
} catch (error) {
// Exécuter les compensations dans l'ordre inverse
saga.state = 'COMPENSATING';
await this.saveSaga(saga);
for (const compensation of saga.compensations.reverse()) {
try {
await compensation();
} catch (compError) {
console.error('Compensation failed:', compError);
}
}
saga.state = 'FAILED';
await this.saveSaga(saga);
return { success: false, error: error.message };
}
}
}
📊 Event Streaming vs Event Messaging
Event Messaging (RabbitMQ)
# Point-to-point ou Pub/Sub
# Messages consommés et supprimés
class RabbitMQEventBus:
def publish(self, event):
self.channel.basic_publish(
exchange='events',
routing_key=f'{event.aggregate_type}.{event.event_type}',
body=json.dumps(event.to_dict())
)
def subscribe(self, event_pattern, handler):
queue = self.channel.queue_declare('', exclusive=True)
self.channel.queue_bind(
exchange='events',
queue=queue.method.queue,
routing_key=event_pattern
)
self.channel.basic_consume(
queue=queue.method.queue,
on_message_callback=handler
)
Event Streaming (Kafka-style)
# Log immutable d'événements
# Replay possible, multiple consumers
class EventStreamProcessor:
def __init__(self, stream_name):
self.stream = stream_name
self.offset = 0
def append(self, event):
# Append-only log
partition = self.get_partition(event.aggregate_id)
offset = self.store.append(partition, event)
return offset
def read_from(self, offset, batch_size=100):
# Lecture depuis un offset
events = self.store.read(self.stream, offset, batch_size)
return events
def replay(self, from_beginning=False):
# Rejouer tous les événements
offset = 0 if from_beginning else self.get_last_offset()
while True:
events = self.read_from(offset)
if not events:
break
for event in events:
yield event
offset = events[-1].offset + 1
🎯 Cas d'Usage : Système de E-Commerce
Architecture Complète
// Agrégat Order avec Event Sourcing
class OrderAggregate {
constructor() {
this.events = [];
this.version = 0;
}
// Commandes
createOrder(customerId, items) {
this.applyEvent({
type: 'OrderCreated',
payload: { customerId, items, orderId: generateId() }
});
}
approvePayment(paymentId) {
if (this.state.status !== 'pending_payment') {
throw new Error('Invalid state for payment approval');
}
this.applyEvent({
type: 'PaymentApproved',
payload: { paymentId }
});
}
shipOrder(trackingNumber) {
if (this.state.status !== 'paid') {
throw new Error('Cannot ship unpaid order');
}
this.applyEvent({
type: 'OrderShipped',
payload: { trackingNumber }
});
}
// Application des événements
applyEvent(event) {
switch(event.type) {
case 'OrderCreated':
this.state = {
id: event.payload.orderId,
customerId: event.payload.customerId,
items: event.payload.items,
status: 'pending_payment'
};
break;
case 'PaymentApproved':
this.state.status = 'paid';
this.state.paymentId = event.payload.paymentId;
break;
case 'OrderShipped':
this.state.status = 'shipped';
this.state.trackingNumber = event.payload.trackingNumber;
break;
}
this.events.push({
...event,
version: ++this.version,
timestamp: new Date().toISOString()
});
}
// Récupération depuis l'event store
static fromEvents(events) {
const aggregate = new OrderAggregate();
events.forEach(event => aggregate.applyEvent(event));
return aggregate;
}
}
Projections et Read Models
# Projection pour le tableau de bord
class DashboardProjection:
def __init__(self, db):
self.db = db
async def handle_order_created(self, event):
# Mise à jour des statistiques
await self.db.stats.increment('daily_orders')
await self.db.stats.add_to_set(
'active_customers',
event['customerId']
)
await self.db.stats.increment_by(
'daily_revenue',
event['totalAmount']
)
async def handle_order_shipped(self, event):
# Mise à jour du temps de traitement
order = await self.db.orders.get(event['orderId'])
processing_time = event['timestamp'] - order['createdAt']
await self.db.stats.push(
'processing_times',
processing_time.total_seconds()
)
# Projection pour la recherche
class SearchProjection:
def __init__(self, elasticsearch):
self.es = elasticsearch
async def handle_order_created(self, event):
# Indexer pour la recherche
await self.es.index(
index='orders',
id=event['orderId'],
body={
'orderId': event['orderId'],
'customerId': event['customerId'],
'customerName': event['customerName'],
'items': event['items'],
'status': 'created',
'createdAt': event['timestamp']
}
)
async def handle_order_updated(self, event):
# Mise à jour partielle
await self.es.update(
index='orders',
id=event['orderId'],
body={'doc': event['changes']}
)
🔍 Monitoring et Observabilité
Distributed Tracing
// Traçage des événements avec OpenTelemetry
const { trace } = require('@opentelemetry/api');
class TracedEventBus {
async publish(event) {
const tracer = trace.getTracer('event-bus');
const span = tracer.startSpan(`publish:${event.type}`);
try {
// Ajouter le contexte de trace à l'événement
event.metadata = {
...event.metadata,
traceId: span.spanContext().traceId,
spanId: span.spanContext().spanId
};
span.setAttributes({
'event.type': event.type,
'event.aggregate_id': event.aggregateId,
'event.version': event.version
});
await this.broker.publish(event);
span.setStatus({ code: SpanStatusCode.OK });
} catch (error) {
span.recordException(error);
span.setStatus({ code: SpanStatusCode.ERROR });
throw error;
} finally {
span.end();
}
}
}
Métriques Essentielles
# Prometheus metrics
metrics:
# Événements
- events_published_total{type, service}
- events_consumed_total{type, service, status}
- event_processing_duration_seconds{type, service}
# Sagas
- saga_started_total{type}
- saga_completed_total{type}
- saga_failed_total{type}
- saga_compensation_executed_total{type, step}
# Event Store
- event_store_append_duration_seconds
- event_store_size_bytes
- event_store_replay_duration_seconds
✅ Bonnes Pratiques
1. Design d'Événements
// BON : Événement descriptif et immuable
const goodEvent = {
eventId: uuid(),
eventType: 'CustomerEmailChanged',
aggregateId: 'customer-123',
timestamp: new Date().toISOString(),
version: 5,
payload: {
oldEmail: 'old@example.com',
newEmail: 'new@example.com'
}
};
// MAUVAIS : Événement impératif
const badEvent = {
type: 'UpdateCustomer',
data: {
customerId: 123,
email: 'new@example.com'
}
};
2. Versioning des Événements
class EventUpgrader:
def upgrade(self, event):
"""Upgrade d'événements vers la dernière version"""
if event['version'] == 1 and event['type'] == 'OrderCreated':
# V1 -> V2: Ajouter le champ shippingMethod
event['payload']['shippingMethod'] = 'standard'
event['version'] = 2
if event['version'] == 2 and event['type'] == 'OrderCreated':
# V2 -> V3: Restructurer les items
event['payload']['items'] = [
{
'sku': item['productId'],
'quantity': item['qty'],
'unitPrice': item['price']
}
for item in event['payload']['items']
]
event['version'] = 3
return event
3. Eventual Consistency
// Gestion de la cohérence éventuelle
class EventuallyConsistentView {
async getOrder(orderId, consistencyLevel = 'eventual') {
if (consistencyLevel === 'strong') {
// Attendre que tous les événements soient traités
await this.waitForConsistency(orderId);
}
// Lire depuis le read model
const order = await this.readModel.getOrder(orderId);
if (consistencyLevel === 'eventual') {
// Avertir de la possible inconsistance
order.metadata = {
consistencyWarning: 'Data may be up to 5 seconds old'
};
}
return order;
}
async waitForConsistency(orderId) {
const maxWait = 5000; // 5 secondes
const start = Date.now();
while (Date.now() - start < maxWait) {
const position = await this.eventStore.getPosition(orderId);
const projected = await this.readModel.getProjectedPosition(orderId);
if (projected >= position) {
return;
}
await sleep(100);
}
throw new Error('Consistency timeout');
}
}
🎓 Exercice Pratique
Implémentez un système de réservation avec Event Sourcing et CQRS :
// TODO: Implémenter
// 1. Aggregate Reservation avec événements
// 2. Command handlers pour créer/modifier/annuler
// 3. Projections pour différentes vues
// 4. Saga pour gérer le workflow complet
// 5. Tests avec event store in-memory
Les architectures orientées événements offrent flexibilité, scalabilité et résilience, mais demandent une compréhension approfondie de la cohérence éventuelle et de la gestion d'état distribué.