Проблема: мікросервіси ростуть, синхронні HTTP-запити перетворюють систему на павутину. Один збій — і весь стек летить. Подієво-орієнтована архітектура (EDA) — єдиний спосіб розв'язати цю зв'язність. Ми впровадили EDA для проєктів з навантаженням до 10 000 подій на секунду і знаємо, як уникнути типових помилок.
Клієнт з інтернет-магазином на 50 000 замовлень на добу зіткнувся з тим, що при пікових навантаженнях сервіс замовлень вичерпувався за таймаутом, блокуючи інвентаризацію та сповіщення. Після впровадження EDA з Apache Kafka та Outbox Pattern час відгуку знизився на 40%, а пропускна здатність зросла в 3 рази. Економія бюджету на підтримку склала 30% (понад $5 000 на місяць). Середня затримка обробки події — 10 мс.
Чому EDA краще за синхронну взаємодію?
Синхронний HTTP-виклик схожий на телефонну розмову: потрібна відповідь негайно. EDA працює як пошта: відправив листа і забув. Порівняємо за параметрами:
| Параметр | Синхронне | Асинхронне (EDA) |
|---|---|---|
| Залежність за часом | Клієнт чекає | Клієнт не блокується |
| Відмовостійкість | Збій руйнує ланцюжок | Черга ізолює збої |
| Навантаження | Пряма пропорція RPS | Буферизація + backpressure |
| Складність розробки | Простіше зрозуміти | Складніше налагодження (ідемпотентність, моніторинг) |
Як ми впроваджуємо EDA: стек та кейс
Ми використовуємо Apache Kafka як основний брокер завдяки високій пропускній здатності та довготривалому зберіганню.
Структура події
interface DomainEvent<T = unknown> { id: string; // UUID — для ідемпотентності type: string; // 'user.registered', 'order.placed' version: string; // '1.0' — для schema evolution source: string; // 'order-service' correlationId: string; // наскрізний ID через всі сервіси causationId?: string; // ID події, що стала причиною occurredAt: string; // ISO 8601 data: T; } // Конкретна подія interface OrderPlacedEvent extends DomainEvent<{ orderId: string; customerId: string; items: Array<{ productId: string; quantity: number; price: number }>; total: number; shippingAddress: Address; }> { type: 'order.placed'; } Apache Kafka — основний брокер
import { Kafka, Partitioners } from 'kafkajs'; const kafka = new Kafka({ clientId: 'order-service', brokers: process.env.KAFKA_BROKERS.split(',') }); // Продюсер const producer = kafka.producer({ createPartitioner: Partitioners.LegacyPartitioner }); async function publishOrderPlaced(order: Order): Promise<void> { await producer.send({ topic: 'order.events', messages: [{ key: order.id, // партиціонування за ID замовлення value: JSON.stringify({ id: uuidv4(), type: 'order.placed', version: '1.0', source: 'order-service', correlationId: context.correlationId, occurredAt: new Date().toISOString(), data: { orderId: order.id, customerId: order.customerId, items: order.items, total: order.total } } satisfies OrderPlacedEvent), headers: { 'content-type': 'application/json', 'schema-version': '1.0' } }] }); } // Консюмер — Inventory Service const consumer = kafka.consumer({ groupId: 'inventory-service' }); await consumer.subscribe({ topics: ['order.events'], fromBeginning: false }); await consumer.run({ eachMessage: async ({ topic, partition, message }) => { const event = JSON.parse(message.value.toString()) as DomainEvent; // Ідемпотентність: перевіряємо, чи не обробляли вже цю подію const processed = await idempotencyRepo.exists(event.id); if (processed) return; try { if (event.type === 'order.placed') { await inventoryService.reserveStock(event.data.orderId, event.data.items); } await idempotencyRepo.mark(event.id); } catch (error) { // Публікуємо в Dead Letter Topic для аналізу await deadLetterProducer.send({ topic: 'order.events.dlq', messages: [{ value: message.value, headers: { 'failure-reason': error.message } }] }); } } }); Outbox Pattern — гарантована доставка
Помилка: зберегти в БД і потім опублікувати в Kafka — ризик втрати події. Правильно — Transactional Outbox:
// В рамках однієї транзакції БД async function createOrder(dto: CreateOrderDto): Promise<Order> { return db.transaction(async (trx) => { // 1. Зберігаємо замовлення const order = await trx('orders').insert({ ...orderData }).returning('*'); // 2. Зберігаємо подію в outbox-таблицю (в тій же транзакції!) await trx('outbox_events').insert({ id: uuidv4(), aggregate_id: order.id, event_type: 'order.placed', payload: JSON.stringify(orderPlacedEvent), status: 'pending', created_at: new Date() }); return order; }); } Окремий Outbox Poller читає pending-події та публікує їх в Kafka:
// Cron job або background worker async function processOutbox(): Promise<void> { const events = await db('outbox_events') .where({ status: 'pending' }) .orderBy('created_at') .limit(100) .forUpdate() .skipLocked(); for (const event of events) { try { await kafka.producer.send({ topic: getTopicForEventType(event.event_type), messages: [{ key: event.aggregate_id, value: event.payload }] }); await db('outbox_events') .where({ id: event.id }) .update({ status: 'published', published_at: new Date() }); } catch { await db('outbox_events') .where({ id: event.id }) .update({ retry_count: db.raw('retry_count + 1') }); } } } Альтернатива — Debezium CDC: читає WAL PostgreSQL і публікує зміни в Kafka без коду.
Як гарантувати доставку без втрат?
Ключові прийоми:
- Transactional Outbox
- Ідемпотентність консюмерів (перевірка за event.id)
- Dead Letter Queue для повідомлень з помилками
- Моніторинг затримок через Prometheus + Grafana
Ми налаштовували систему, яка обробляє 1000 msg/sec без втрат при відмові одного вузла Kafka.
Ефективна координація: хореографія vs оркестрація
| Характеристика | Хореографія | Оркестрація |
|---|---|---|
| Координація | Сервіси реагують на події | Центральний оркестратор |
| Зв'язність | Низька | Середня |
| Видимість потоку | Складно відстежити | Явно в коді |
| Тестування | Складніше | Простіше |
Хореографія підходить для слабко зв'язаних сценаріїв, оркестрація — коли потрібна строга послідовність. Ми комбінуємо: всередині сервісу — оркестрація, між — хореографія.
Що таке Event Sourcing та CQRS?
При Event Sourcing всі зміни зберігаються як події, які публікуються і в event store, і в брокер. CQRS (Command Query Responsibility Segregation) розділяє операції запису та читання, часто використовуючи ті ж події для побудови проекцій. EDA та CQRS/ES чудово поєднуються, але вимагають осмисленого проектування.
Приклад production-конфігурації Kafka
Для надійної роботи в проді налаштовуйте реплікацію фактором 3, retention 7 днів, compaction для ключових топіків. Мінімум 3 брокери, моніторинг lag consumer через Burrow. Для cloud-середовищ — Confluent Cloud або AWS MSK.
Що входить в роботу з впровадження EDA?
При замовленні ми надаємо:
- Архітектурну документацію (діаграми, вибір брокера)
- Розробку продюсерів та консюмерів з Outbox Pattern
- Налаштування Kafka/RabbitMQ з persistency та реплікацією
- Реалізацію ідемпотентності та Dead Letter Queue
- Навантажувальне тестування до 1000 msg/sec без втрат
- Моніторинг та алертинг (latency, throughput)
- Навчання команди
Отримайте консультацію з впровадження EDA для вашого проєкту — ми оцінимо архітектуру та запропонуємо оптимальне рішення.
Скільки часу займає впровадження EDA?
| Етап | Термін |
|---|---|
| Один сценарій (продюсер + 2–3 консюмери) | 1–2 тижні |
| + Outbox Pattern, ідемпотентність, DLQ | +1 тиждень |
| Повноцінна EDA для 5–10 сервісів + моніторинг | 4–8 тижнів |
Вибір брокера залежить від навантаження та вимог. Kafka — для високого throughput та довгого зберігання. RabbitMQ — для складної маршрутизації. Redis Streams — для низької затримки. Google Pub/Sub — для хмари. У більшості проєктів використовуємо Kafka.
Впровадження EDA за 4 кроки
- Аналіз та проектування: виділяємо межі контекстів, визначаємо події.
- Налаштування брокера: розгортаємо Kafka з реплікацією, налаштовуємо топіки.
- Реалізація продюсерів та консюмерів: пишемо код з Outbox та ідемпотентністю.
- Моніторинг та налагодження: налаштовуємо метрики, DLQ, алерти.
Наша команда — 10+ років досвіду в розподілених системах, понад 20 проєктів з EDA. Замовте впровадження EDA під ключ — отримайте надійну та масштабовану архітектуру. Зв'яжіться з нами для оцінки вашого проєкту.







