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







