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







