Event-Driven Architecture: асинхронная коммуникация
Event sourcing, CQRS, message queues и eventual consistency
Синхронные REST API создают coupling между сервисами. Event-Driven Architecture решает проблемы масштабирования и резильентности. Amazon, Netflix, Uber построены на событиях. Event sourcing, CQRS, message queues — разберем паттерны асинхронной архитектуры. Проблемы синхронной архитектуры Tight coupling — сервисы зависят друг от друга Cascading failures — падение одного ломает все Performance bottlenecks — ожидание ответа Scaling challenges — сложно масштабировать части системы Event-Driven Architecture Principles: Асинхронность — не ждем ответа Loose coupling — сервисы независимы Event immutability — события неизменяемы Eventual consistency — согласованность со временем Event Structure: { "eventId": "evt_123456", "eventType": "OrderCreated", "timestamp": "2024-01-15T10:30:00Z", "aggregateId": "order_789", "version": 1, "payload": { "orderId": "order_789", "userId": "user_456", "items": [ { "productId": "prod_1", "quantity": 2, "price": 29.99 } ], "total": 59.98 }, "metadata": { "correlationId": "req_abc", "causationId": "evt_prev", "userId": "user_456" } } Message Broker: RabbitMQ Publisher: const amqp = require('amqplib'); class EventPublisher { constructor() { this.connection = null; this.channel = null; } async connect() { this.connection = await amqp.connect('amqp://localhost'); this.channel = await this.connection.createChannel(); // Declare exchange await this.channel.assertExchange('events', 'topic', { durable: true }); } async publish(eventType, payload) { const event = { eventId: crypto.randomUUID(), eventType, timestamp: new Date().toISOString(), payload }; await this.channel.publish( 'events', eventType, // Routing key Buffer.from(JSON.stringify(event)), { persistent: true } ); console.log('Published event:', eventType); } async close() { await this.channel.close(); await this.connection.close(); } } // Usage const publisher = new EventPublisher(); await publisher.connect(); await publisher.publish('OrderCreated', { orderId: 'order_789', userId: 'user_456', total: 59.98 }); Consumer: class EventConsumer { constructor(queueName) { this.queueName = queueName; this.handlers = new Map(); } async connect() { this.connection = await amqp.connect('amqp://localhost'); this.channel = await this.connection.createChannel(); await this.channel.assertExchange('events', 'topic', { durable: true }); await this.channel.assertQueue(this.queueName, { durable: true }); } subscribe(eventType, handler) { this.handlers.set(eventType, handler); // Bind queue to exchange with routing key this.channel.bindQueue(this.queueName, 'events', eventType); } async start() { await this.channel.consume(this.queueName, async (msg) => { if (!msg) return; const event = JSON.parse(msg.content.toString()); const handler = this.handlers.get(event.eventType); if (handler) { try { await handler(event); this.channel.ack(msg); } catch (error) { console.error('Handler error:', error); // Requeue or send to DLQ this.channel.nack(msg, false, false); } } else { this.channel.ack(msg); } }); } } // Usage const consumer = new EventConsumer('email-service'); await consumer.connect(); consumer.subscribe('OrderCreated', async (event) => { console.log('Sending email for order:', event.payload.orderId); await sendOrderConfirmationEmail(event.payload); }); await consumer.start(); Event Sourcing Храним все изменения как события вместо текущего состояния. Aggregate: class Order { constructor() { this.id = null; this.userId = null; this.items = []; this.status = 'pending'; this.total = 0; this.version = 0; } // Command handlers create(userId, items) { if (this.id) throw new Error('Order already created'); return new OrderCreated({ orderId: crypto.randomUUID(), userId, items, total: items.reduce((sum, i) => sum + i.price * i.quantity, 0) }); } approve() { if (this.status !== 'pending') { throw new Error('Cannot approve order'); } return new OrderApproved({ orderId: this.id }); } cancel(reason) { if (this.status === 'shipped') { throw new Error('Cannot cancel shipped order'); } return new OrderCancelled({ orderId: this.id, reason }); } // Event handlers (apply events to rebuild state) applyOrderCreated(event) { this.id = event.payload.orderId; this.userId = event.payload.userId; this.items = event.payload.items; this.total = event.payload.total; this.status = 'pending'; this.version++; } applyOrderApproved(event) { this.status = 'approved'; this.version++; } applyOrderCancelled(event) { this.status = 'cancelled'; this.version++; } // Rebuild from events static fromEvents(events) { const order = new Order(); events.forEach(event => { const methodName = `apply${event.eventType}`; if (order[methodName]) { order[methodName](event); } }); return order; } } Event Store: class EventStore { constructor(db) { this.db = db; } async appendEvents(aggregateId, events, expectedVersion) { const transaction = await this.db.transaction(); try { // Optimistic concurrency check const currentVersion = await transaction('events') .where({ aggregateId }) .max('version'); if (currentVersion !== expectedVersion) { throw new Error('Concurrency conflict'); } // Insert events for (const event of events) { await transaction('events').insert({ eventId: event.eventId, aggregateId, eventType: event.eventType, payload: JSON.stringify(event.payload), version: event.version, timestamp: event.timestamp }); } await transaction.commit(); } catch (error) { await transaction.rollback(); throw error; } } async getEvents(aggregateId, fromVersion = 0) { const rows = await this.db('events') .where({ aggregateId }) .where('version', '>', fromVersion) .orderBy('version'); return rows.map(row => ({ eventId: row.eventId, eventType: row.eventType, payload: JSON.parse(row.payload), version: row.version, timestamp: row.timestamp })); } } CQRS (Command Query Responsibility Segregation) Разделение чтения и записи. Command Side (Write): class OrderCommandHandler { constructor(eventStore, eventPublisher) { this.eventStore = eventStore; this.eventPublisher = eventPublisher; } async handleCreateOrder(command) { const order = new Order(); const event = order.create(command.userId, command.items); await this.eventStore.appendEvents(event.payload.orderId, [event], 0); await this.eventPublisher.publish(event.eventType, event.payload); return { orderId: event.payload.orderId }; } async handleApproveOrder(command) { // Load order from events const events = await this.eventStore.getEvents(command.orderId); const order = Order.fromEvents(events); const event = order.approve(); await this.eventStore.appendEvents( command.orderId, [event], order.version ); await this.eventPublisher.publish(event.eventType, event.payload); } } Query Side (Read): // Projection (denormalized read model) class OrderProjection { constructor(db) { this.db = db; } async handleOrderCreated(event) { await this.db('orders_view').insert({ orderId: event.payload.orderId, userId: event.payload.userId, status: 'pending', total: event.payload.total, createdAt: event.timestamp }); } async handleOrderApproved(event) { await this.db('orders_view') .where({ orderId: event.payload.orderId }) .update({ status: 'approved' }); } async handleOrderCancelled(event) { await this.db('orders_view') .where({ orderId: event.payload.orderId }) .update({ status: 'cancelled' }); } } // Query service class OrderQueryService { constructor(db) { this.db = db; } async getOrder(orderId) { return await this.db('orders_view') .where({ orderId }) .first(); } async getUserOrders(userId) { return await this.db('orders_view') .where({ userId }) .orderBy('createdAt', 'desc'); } } Saga Pattern Distributed transaction через события. class OrderSaga { constructor(eventPublisher) { this.eventPublisher = eventPublisher; } async handleOrderCreated(event) { const { orderId, userId, total } = event.payload; try { // Step 1: Reserve inventory await this.eventPublisher.publish('ReserveInventory', { orderId, items: event.payload.items }); } catch (error) { await this.eventPublisher.publish('OrderFailed', { orderId, reason: 'Inventory reservation failed' }); } } async handleInventoryReserved(event) { const { orderId } = event.payload; try { // Step 2: Process payment await this.eventPublisher.publish('ProcessPayment', { orderId, amount: event.payload.total }); } catch (error) { // Compensating transaction await this.eventPublisher.publish('ReleaseInventory', { orderId }); } } async handlePaymentProcessed(event) { // Step 3: Complete order await this.eventPublisher.publish('CompleteOrder', { orderId: event.payload.orderId }); } } Dead Letter Queue class DeadLetterQueue { async handle(failedMessage, error, retryCount) { if (retryCount { this.retry(failedMessage); }, delay); } else { // Send to DLQ await this.channel.sendToQueue('dlq', failedMessage.content, { headers: { 'x-error': error.message, 'x-retry-count': retryCount, 'x-failed-at': new Date().toISOString() } }); // Alert on-call await alert.notify('Message sent to DLQ', { message: failedMessage, error }); } } } Idempotency class IdempotentEventHandler { constructor(db) { this.db = db; } async handle(event) { // Check if already processed const exists = await this.db('processed_events') .where({ eventId: event.eventId }) .first(); if (exists) { console.log('Event already processed, skipping'); return; } // Process event await this.processEvent(event); // Mark as processed await this.db('processed_events').insert({ eventId: event.eventId, processedAt: new Date() }); } } Best Practices Event naming — прошедшее время (OrderCreated, не CreateOrder) Immutable events — никогда не изменяйте Schema evolution — версионирование событий Idempotency — обрабатывайте события дважды безопасно Ordering — не полагайтесь на порядок между aggregates Monitoring — tracking lag, DLQ size Заключение: Event-Driven Architecture — ключ к масштабируемым системам. Event sourcing дает audit trail и time travel. CQRS оптимизирует read и write. Saga для distributed transactions. Message brokers для асинхронности. Начните с простых событий. Добавляйте event sourcing где нужна история. CQRS для сложных query. Eventual consistency — компромисс за масштабирование.