Проблема: мікросервіси ростуть, синхронні 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 під ключ — отримайте надійну та масштабовану архітектуру. Зв'яжіться з нами для оцінки вашого проєкту.







