Событийная архитектура: брокеры сообщений, гарантии доставки и паттерны
Event-Driven Architecture на практике: RabbitMQ или Kafka, формат событий, идемпотентные обработчики, Transactional Outbox, порядок и версионирование событий, Saga и CQRS, мониторинг асинхронных систем.
Оформление заказа в интернет-магазине запускает десяток действий: списать остаток, провести платёж, отправить письмо, начислить бонусы, передать заказ в доставку и учётную систему, обновить аналитику. Если сервис заказов вызывает каждого участника синхронно и ждёт ответа, сбой почтового сервиса или медленная учётная система ломают покупку. Событийная архитектура меняет схему: сервис заказов сообщает «заказ создан», а каждый заинтересованный сервис реагирует сам и в своём темпе. В статье — когда события действительно нужны, чем брокер очередей отличается от журнала событий, как не потерять и не обработать дважды сообщение, как развивать формат событий и какие паттерны — Saga, CQRS, Event Sourcing — стоит применять, а какие лучше отложить. Когда синхронные вызовы становятся проблемой Синхронный вызов прост и понятен: отправил запрос, получил ответ, продолжил работу. Проблемы начинаются, когда таких вызовов много и они выстраиваются в цепочки. Связанность. Сервис заказов знает обо всех, кого нужно уведомить. Новый получатель — это изменение и релиз сервиса заказов. Каскадные отказы. Недоступность любого участника цепочки превращается в ошибку для пользователя, даже если его функция второстепенна. Суммарная задержка. Время ответа складывается из времени всех вызовов, а самый медленный участник задаёт темп всей операции. Пиковые нагрузки. Всплеск заказов одновременно бьёт по всем зависимым сервисам, включая те, которым некуда спешить. События снимают эти проблемы для действий, результат которых не нужен отправителю немедленно. Там, где ответ нужен сейчас — проверить остаток перед оплатой, рассчитать цену, — синхронный вызов остаётся правильным выбором. Событие, команда и сообщение Путаница в терминах приводит к путанице в архитектуре, поэтому договоримся о словах. Событие — факт, который уже произошёл: OrderPlaced , PaymentCaptured . Называется глаголом в прошедшем времени, не адресовано конкретному получателю и не может быть отменено — только скомпенсировано новым событием. Команда — просьба выполнить действие: CapturePayment , SendInvoice . У неё есть конкретный адресат, и она может быть отклонена. Сообщение — транспортная единица, в которой передаётся событие или команда. Отправитель события не знает, кто его обработает, и не зависит от результата. Если сервису важно, что действие выполнено, — это команда, и ей нужен ответ или подтверждающее событие. Четыре значения «событийной архитектуры» Мартин Фаулер в докладе «The Many Meanings of Event-Driven Architecture» на конференции GOTO в 2017 году показал, что под этим названием смешивают четыре разных паттерна с разными свойствами. Уведомление о событии. В сообщении минимум данных — «заказ 123 создан». Получатель при необходимости сам запрашивает подробности. Слабая связанность, но дополнительные вызовы и зависимость от доступности источника. Передача состояния в событии. Событие несёт все нужные данные: состав заказа, адрес, сумму. Получатели хранят у себя копию и не обращаются к источнику. Выше устойчивость, но данные дублируются и становятся согласованными с задержкой. Event Sourcing. Журнал событий — основной источник данных, а текущее состояние вычисляется из него. CQRS. Модели записи и чтения разделены, и модель чтения часто строится из событий. Первые два паттерна — основа большинства систем. Последние два — сложные техники для отдельных частей системы, и путать их с «просто отправляем события в брокер» не стоит. Брокер очередей или журнал событий Инфраструктура для событий делится на два класса, и выбор между ними определяет возможности системы. Критерий Брокер очередей (RabbitMQ) Журнал событий (Apache Kafka) Модель Сообщение доставляется в очередь и удаляется после подтверждения Сообщения хранятся в журнале заданное время, потребители читают со своей позиции Маршрутизация Гибкая: по ключам, шаблонам, заголовкам По темам и разделам Повторное чтение Нет, после подтверждения сообщение удалено (кроме потоков RabbitMQ) Да, можно перечитать историю с любой позиции Порядок В пределах очереди при одном потребителе В пределах раздела, по ключу сообщения Типичные задачи Фоновые задачи, команды, интеграции, маршрутизация Потоки событий, аналитика, репликация данных между системами Эксплуатация Проще на старте Сложнее, выше требования к ресурсам Обе системы активно развиваются. В Apache Kafka 4.0, вышедшей в марте 2025 года, полностью удалена поддержка ZooKeeper: метаданные кластера управляются только встроенным механизмом KRaft, что упрощает развёртывание. Прямое обновление на 4.0 возможно только с кластеров, уже работающих в режиме KRaft. В RabbitMQ 4.0 удалено зеркалирование классических очередей, объявленное устаревшим тремя годами ранее; для отказоустойчивости используются кворумные очереди и потоки. Если инфраструктура размещена в российском облаке, разумно рассмотреть управляемые варианты — например, Managed Service for Apache Kafka и Message Queue в Yandex Cloud: обновления, резервирование и мониторинг берёт на себя провайдер. Для большинства продуктовых задач среднего масштаба RabbitMQ или управляемая очередь закрывают потребности с меньшими затратами. Kafka оправдана, когда события нужно хранить и перечитывать, потоков много и они большие, или события должны получать аналитические системы. Как оформить событие Структура события — это публичный контракт между командами, поэтому её стоит стандартизировать сразу. Удобная основа — спецификация CloudEvents, проект CNCF, который задаёт обязательные атрибуты: идентификатор, источник, тип, время. { "specversion": "1.0", "id": "8f3c2a9e-6f1d-4d3b-9a57-2c1e0b7d4f10", "source": "/services/orders", "type": "ru.example.orders.order-placed.v1", "time": "2026-09-01T10:30:00Z", "datacontenttype": "application/json", "subject": "order-100245", "correlationid": "req-5b1e7c", "data": { "orderId": "order-100245", "customerId": "cust-3310", "items": [{ "sku": "SKU-1", "quantity": 2, "price": "1490.00" }], "total": "2980.00", "currency": "RUB" } } Уникальный идентификатор нужен для защиты от повторной обработки. Версия в типе позволяет выпускать несовместимые изменения параллельно со старым форматом. Идентификатор корреляции связывает событие с исходным запросом пользователя и нужен для трассировки. Денежные суммы передаются строкой или в минимальных единицах, а не числом с плавающей точкой. Персональные данные — только если без них получатель не справится: события хранятся и копируются, и каждая копия подпадает под требования к защите данных. Публикация и потребление: пример на RabbitMQ Пример на Node.js с библиотекой amqplib . Публикатор использует канал с подтверждениями, чтобы знать, что брокер действительно принял сообщение: const amqp = require('amqplib'); const { randomUUID } = require('node:crypto'); async function createPublisher(url) { const connection = await amqp.connect(url); const channel = await connection.createConfirmChannel(); await channel.assertExchange('events', 'topic', { durable: true }); return async function publish(type, data, { correlationId } = {}) { const event = { specversion: '1.0', id: randomUUID(), source: '/services/orders', type, time: new Date().toISOString(), correlationid: correlationId, data, }; channel.publish('events', type, Buffer.from(JSON.stringify(event)), { persistent: true, contentType: 'application/json', messageId: event.id, }); await channel.waitForConfirms(); // ошибка, если брокер не принял сообщение return event.id; }; } Потребитель ограничивает число сообщений в работе, подтверждает обработку явно и отправляет сообщения, которые обработать не удалось, в очередь «мёртвых писем» через отдельный обменник: async function startConsumer(url, handle) { const connection = await amqp.connect(url); const channel = await connection.createChannel(); await channel.assertExchange('events', 'topic', { durable: true }); await channel.assertExchange('events.dlx', 'fanout', { durable: true }); await channel.assertQueue('loyalty.order-placed.dlq', { durable: true }); await channel.bindQueue('loyalty.order-placed.dlq', 'events.dlx', ''); await channel.assertQueue('loyalty.order-placed', { durable: true, arguments: { 'x-queue-type': 'quorum', 'x-dead-letter-exchange': 'events.dlx', 'x-delivery-limit': 5, // после 5 неудачных доставок — в DLQ }, }); await channel.bindQueue('loyalty.order-placed', 'events', 'ru.example.orders.order-placed.v1'); await channel.prefetch(10); await channel.consume('loyalty.order-placed', async (msg) = { if (!msg) return; try { const event = JSON.parse(msg.content.toString()); await handle(event); channel.ack(msg); } catch (err) { console.error('failed to handle event', err); channel.nack(msg, false, true); // вернуть в очередь; лимит доставок отправит его в DLQ } }); } Очередь мёртвых писем без мониторинга бесполезна: на рост её длины должен срабатывать алерт, а у команды — быть понятный порядок разбора и повторной отправки. Гарантии доставки В распределённой системе сообщение может потеряться или прийти дважды: сеть обрывается между обработкой и подтверждением, потребитель падает, брокер переключается на реплику. Возможные гарантии: Не более одного раза. Сообщение может потеряться, но не повторится. Подходит для телеметрии, где потеря одной точки некритична. Не менее одного раза. Сообщение не потеряется, но может прийти повторно. Стандарт для бизнес-событий. Ровно один раз. Достижимо в ограниченных рамках — например, транзакции Kafka при чтении и записи внутри Kafka. Как только обработка затрагивает внешнюю систему — базу данных, платёжный шлюз, почту, — гарантия превращается в сочетание доставки «не менее одного раза» и идемпотентного обработчика. Практический вывод: проектируйте каждый обработчик так, будто любое событие придёт дважды. Идемпотентный обработчик Надёжный способ — фиксировать идентификатор обработанного события в той же транзакции, что и бизнес-изменение. Если событие пришло повторно, вставка упрётся в уникальный ключ, и изменение не выполнится второй раз. async function handleOrderPlaced(pool, event) { const client = await pool.connect(); try { await client.query('BEGIN'); const inserted = await client.query( 'INSERT INTO processed_events (consumer, event_id) VALUES ($1, $2) ON CONFLICT DO NOTHING', ['loyalty', event.id] ); if (inserted.rowCount === 0) { // уже обработано await client.query('ROLLBACK'); return; } await client.query( 'INSERT INTO bonus_accruals (customer_id, order_id, amount) VALUES ($1, $2, $3)', [event.data.customerId, event.data.orderId, calculateBonus(event.data.total)] ); await client.query('COMMIT'); } catch (err) { await client.query('ROLLBACK'); throw err; } finally { client.release(); } } Таблица processed_events с уникальным ключом по паре «потребитель — идентификатор события» периодически очищается от записей старше максимального срока, в течение которого возможен повтор. Transactional Outbox: как не потерять событие при публикации Классическая ошибка: сервис сохраняет заказ в базу, а затем публикует событие в брокер. Если процесс упадёт между этими шагами, заказ есть, а события нет — бонусы не начислятся, доставка не узнает о заказе. Если поменять шаги местами, получится событие о заказе, которого нет. Паттерн Transactional Outbox решает проблему: событие записывается в таблицу исходящих сообщений в той же транзакции, что и бизнес-данные, а отдельный процесс публикует его в брокер. -- В одной транзакции с созданием заказа BEGIN; INSERT INTO orders (id, customer_id, total, status) VALUES ('order-100245', 'cust-3310', 2980.00, 'placed'); INSERT INTO outbox (id, aggregate_id, type, payload, created_at) VALUES (gen_random_uuid(), 'order-100245', 'ru.example.orders.order-placed.v1', $1, now()); COMMIT; -- Процесс-ретранслятор забирает пачку неопубликованных событий SELECT id, type, payload FROM outbox WHERE published_at IS NULL ORDER BY created_at LIMIT 100 FOR UPDATE SKIP LOCKED; -- ...публикует в брокер с подтверждением и отмечает published_at = now() Вместо опроса таблицы можно читать журнал изменений базы данных с помощью Debezium и публиковать события из него. Доставка по-прежнему «н