Модуль 3.5. Событийная модель и доменная интеграция
Темы: события домена, схемы, совместная эволюция, outbox, гарантированная доставка. Артефакт: событийный каталог. Практика: разработать 5 доменных событий.
Системный аналитик (SA) задаёт общий язык фактов между системами: как выглядит событие, какие поля обязательны, как оно меняется со временем, какие гарантии доставки и порядка соблюдаются. В этом модуле вы получите:
- чек-лист проектирования событий;
- правила эволюции схем (совместимость, версии);
- рабочие конфигурации outbox + CDC;
- практикум: 5 событий e-commerce-домена, готовые схемы и рекомендации по ключам/партициям.
Что такое доменное событие (и что им не является)
- Событие домена — необратимый факт, который уже случился в домене: OrderCreated, PaymentCaptured, RefundCompleted.
- Не события: команды (ReserveStock), логи трассировки, технические метрики. Команда адресна и просит «сделать», событие — факт «случилось».
Принципы событийной модели
- Правда и минимум: только факты и минимальный набор атрибутов (без PII/секретов).
- Детерминированные ключи: eventId (UUID), aggregateId (например, orderId), correlationId/causationId.
- Повторяемость: допускаем повторы и перестановки → консъюмеры обязаны быть идемпотентными.
- Порядок/партиционирование: порядок гарантирован только внутри ключа партиции (orderId), между агрегатами порядка нет.
Конверт (envelope) события: стандарт для проекта
Рекомендуемый каркас (не payload, а обёртка). Любой продьюсер и консъюмер должны его понимать:
{
"eventId": "uuid",
"type": "payment.captured.v1",
"occurredAt": "2025-08-20T12:34:56Z",
"producer": "payments-service",
"correlationId": "uuid-or-traceparent",
"causationId": "uuid-of-parent-event-or-command",
"aggregateType": "Order",
"aggregateId": "uuid-of-order",
"schemaVersion": "1.0.0",
"payload": { /* доменные поля события */ }
}
Почему так:
- eventId — основа дедупликации у подписчиков;
- correlationId/causationId — сквозная трассировка бизнес-потока;
- aggregateId — partition key (гарантия локального порядка).
Схемы сообщений и эволюция (совместная)
Схема (Avro/Protobuf/JSON Schema)
- Для шины (Kafka/NATS/Kinesis) — предпочтительны Avro/Proto + Schema Registry.
- Схема описывает payload; envelope можно фиксировать отдельно.
Политики совместимости (рекомендуем)
- Backward как дефолт (новая схема читает старые данные).
- Full (backward+forward) — если много живых потребителей с разными версиями.
- Запрещаем: переименование/удаление поля без MAJOR; смену типа; изменение семантики.
Эволюция без боли
- Новые поля — optional с дефолтом;
- Устаревшие — помечаем deprecated, оставляем N релизов;
- Ломающие — новое имя события: payment.captured.v2 и параллельная публикация в течение окна миграции.
Гарантированная доставка на практике
- В реальных EDA достигается at-least-once + идемпотентные консъюмеры → «effect exactly-once».
- DLQ (dead-letter queue): правим яд, не меняем данные молча.
- Повторы/перестановки: консъюмер хранит processed_event(eventId, processedAt) с TTL ≥ retention.
Метрики прочности
- consumer_lag, dlq_size, event_dup_rate, out_of_order_count, publish_latency_p95.
Outbox + CDC — лекарство от dual-write
Проблема
Обновили БД и отправили событие — любой сбой между операциями создаст рассинхрон. Нужна атомарность.
Шаблон Outbox
- В одной транзакции с бизнес-записью пишем строку в outbox.
- Фоновый паблишер (или CDC) читает outbox и публикует в шину.
- Продьюсер помечает сообщение как отправленное и/или удаляет/компактиует.
DDL-эскиз (PostgreSQL)
CREATE TABLE outbox_events ( id uuid PRIMARY KEY, aggregate_type text NOT NULL, aggregate_id text NOT NULL, type text NOT NULL, -- e.g. payment.captured.v1 payload jsonb NOT NULL, headers jsonb NOT NULL, occurred_at timestamptz NOT NULL DEFAULT now(), published_at timestamptz, attempts int NOT NULL DEFAULT 0 ); CREATE INDEX ON outbox_events (published_at) WHERE published_at IS NULL;
Паблишер:
- читает батч непубликованных;
- публикует (с ключом aggregate_id);
- при успехе проставляет published_at;
- при ошибке — увеличивает attempts и откладывает по экспоненте;
- яд «ядовитых» — в DLQ + тикет.
CDC-вариант: Debezium/Streams → снижает накладные расходы, повышает устойчивость.
Паттерны порядка и дедупликации
- Partition key = бизнес-ключ агрегата (orderId, customerId).
- В payload держим монотонный version/sequence объекта → консъюмер игнорирует устаревшие (seq < current).
- Таблица/кэш processed_event(eventId) (TTL 7–30 дней) → идемпотентность при повторах.
Безопасность и приватность
- В событиях нет PII/секретов — только идентификаторы.
- Если нужно передать контекст (валюта/сумма) — только то, что безопасно и не под NDA.
- Подписывайте вебхуки (HMAC/MTLS), ограничивайте ACL публикации/подписки; шифруйте at rest/in flight.
Наблюдаемость событий
- В каждый publish — correlationId, causationId, метки schemaVersion.
- Дашборды: publish_latency_p95, topic_rps, consumer_lag, dlq_size, dedup_hits, out_of_order.
- Логи: структурированные (JSON), без PII.
Категории событий (как говорить с бизнесом)
- События жизненного цикла (OrderCreated, OrderPaid, OrderShipped).
- События изменения атрибутов (CustomerAddressChanged).
- События компенсаций/ошибок (RefundFailed).
- Снимки/сводки (редко): OrderSnapshotCreated — для поздних подписчиков.
Практические примеры: 5 доменных событий
Для домена «Заказы–Платежи–Возвраты». Покажу Avro-схемы payload, AsyncAPI-фрагменты и рекомендации по партициям.
order.created.v1
Назначение: факт создания заказа.
Partition key: orderId.
Потребители: склад/логистика, аналитика, email.
Avro payload (order-created.v1.avsc):
{
"type":"record","name":"OrderCreated","namespace":"events.orders.v1",
"fields":[
{"name":"orderId","type":"string"},
{"name":"customerId","type":"string"},
{"name":"currency","type":"string"},
{"name":"totalAmount","type":"string"},
{"name":"createdAt","type":{"type":"long","logicalType":"timestamp-millis"}},
{"name":"items","type":{"type":"array","items":{
"name":"Item","type":"record","fields":[
{"name":"productId","type":"string"},
{"name":"qty","type":"int"},
{"name":"price","type":"string"}
]}}}
]
}
payment.captured.v1
Назначение: деньги фактически списаны.
Partition key: orderId (сохранить локальный порядок заказа).
Потребители: отгрузка, уведомления, бухгалтерия.
Avro payload (payment-captured.v1.avsc):
{
"type":"record","name":"PaymentCaptured","namespace":"events.payments.v1",
"fields":[
{"name":"paymentId","type":"string"},
{"name":"orderId","type":"string"},
{"name":"amount","type":"string"},
{"name":"currency","type":"string"},
{"name":"capturedAt","type":{"type":"long","logicalType":"timestamp-millis"}},
{"name":"method","type":["null",{"type":"enum","name":"Method","symbols":["CARD","APPLE_PAY","GOOGLE_PAY"]}], "default": null}
]
}
order.shipped.v1
Назначение: заказ передан в доставку.
Partition key: orderId.
Потребители: трекинг, E-mail/SMS, аналитика.
Avro payload (order-shipped.v1.avsc):
{
"type":"record","name":"OrderShipped","namespace":"events.orders.v1",
"fields":[
{"name":"orderId","type":"string"},
{"name":"shipmentId","type":"string"},
{"name":"carrier","type":"string"},
{"name":"trackingNumber","type":"string"},
{"name":"shippedAt","type":{"type":"long","logicalType":"timestamp-millis"}}
]
}
refund.completed.v1
Назначение: возврат завершён.
Partition key: paymentId (или orderId, если все потребители привязаны к заказу).
Потребители: биллинг, аналитика, антифрод.
Avro payload (refund-completed.v1.avsc):
{
"type":"record","name":"RefundCompleted","namespace":"events.refunds.v1",
"fields":[
{"name":"refundId","type":"string"},
{"name":"paymentId","type":"string"},
{"name":"orderId","type":"string"},
{"name":"amount","type":"string"},
{"name":"currency","type":"string"},
{"name":"completedAt","type":{"type":"long","logicalType":"timestamp-millis"}},
{"name":"reason","type":["null",{"type":"enum","name":"RefundReason","symbols":["CUSTOMER_REQUEST","DUPLICATE","FRAUD_SUSPECTED"]}],"default":null}
]
}
ustomer.address.changed.v1
Назначение: у клиента сменился адрес (для аналитики/антифрода/маркетинга).
Partition key: customerId.
Потребители: CRM/DWH.
Avro payload (customer-address-changed.v1.avsc):
{
"type":"record","name":"CustomerAddressChanged","namespace":"events.customers.v1",
"fields":[
{"name":"customerId","type":"string"},
{"name":"addressId","type":"string"},
{"name":"changedAt","type":{"type":"long","logicalType":"timestamp-millis"}},
{"name":"country","type":"string"},
{"name":"city","type":"string"}
]
}
Заметьте: никаких PII (улица/дом/квартира) в событии — только агрегированные поля.
10.x. AsyncAPI-фрагмент (один топик в качестве примера)
asyncapi: '3.0.0'
info: { title: Commerce Events, version: 1.0.0 }
channels:
payments.captured.v1:
address: payments.captured.v1
messages:
PaymentCaptured:
$ref: '#/components/messages/PaymentCaptured'
components:
messages:
PaymentCaptured:
name: PaymentCaptured
payload:
$ref: 'avro://events.payments.v1.PaymentCaptured'
correlationId:
location: "$message.header#/correlationId"
headers:
type: object
properties:
eventId: { type: string, format: uuid }
correlationId: { type: string }
causationId: { type: string }
aggregateId: { type: string }
servers:
prod: { host: kafka01:9092, protocol: kafka }
Каталог событий (артефакт)
Шаблон страницы каталога: каждое событие = карточка.
# Event Catalog vX.Y.Z ## payments.captured.v1 Описание: Деньги списаны. Producers: payments-service Consumers: shipping, email, billing Key (partition): orderId Ordering: гарантируется внутри orderId Semantics: at-least-once, idempotent required Schema: avro://events.payments.v1.PaymentCaptured (compat: backward) Retention: 7 days (compact: off) Security: ACL (produce: payments-svc, consume: shipping|email|billing), no PII Observability: metrics publish_latency_p95, topic_rps, consumer_lag; log fields eventId/correlationId Change policy: N/N-1, deprecate window 90 days
Рядом храните:
- схемы (Avro/Proto) и AsyncAPI;
- пример payload;
- матрицу зависимости (кто слушает);
- политику версионирования и окно депрекейта.
Тестирование событий
- Contract tests: consumer-driven (Pact/Avro-compat), гейт в CI.
- Golden cases: набор примеров payload (валид/невалид, min/max).
- Replays: возможность поднять консюмера и прокрутить от offset (staging topic).
- Chaos: дубликаты, перестановки, «ядовитые» сообщения, задержки.
Типовые риски и как их гасить
- Dual-write (рассинхрон БД и шины). → Outbox + CDC.
- Ломающие изменения схем. → Registry + compat + open-review, параллельная публикация v1/v2.
- PII в событиях. → Политика «минимум полей», DLP-сканы, ревью схем.
- Нет идемпотентности у консюмера. → Таблица processed_event, естественные ключи, seq/version.
- Потеря порядка между агрегатами. → Коммуникация ожиданий: порядок только внутри partition key.
- DLQ без обработки. → Авто-тикеты, SLO по разбору DLQ.
- Бессрочные ретраи. → Классификация ошибок, лимиты попыток, backoff + jitter, quarantine-topic.
Вопрос–Ответ
Q: Почему «exactly-once» — это миф?
A: В распределённых системах нет абсолютной гарантии; достигаем эффект exactly-once через at-least-once + идемпотентность + дедуп.
Q: Как выбирать partition key?
A: По агрегату, для которого важен порядок (orderId, customerId). Если важно пересечение нескольких — пересмотрите дизайн (саги, согласования).
Q: Можно ли публиковать состояние (снимок) вместо события?
A: Да, для поздних подписчиков — как дополнение (compact topic). Но «истинные» факты — события.
Q: Что делать, если консъюмер «отстаёт»?
A: Масштабируйте группу (больше партиций/инстансов), оптимизируйте обработку, включите backpressure, следите за consumer_lag.
Q: Как долго хранить события?
A: Зависят от кейсов: 7–30 дней для интеграции, дольше — для аналитики/реигрывания. Важно согласовать с безопасностью и GDPR.
Практика: разработать 5 доменных событий (90–120 мин)
Задание: для вашего проекта опишите 5 событий (как в §10):
- Выберите имена и версии: order.created.v1, payment.captured.v1, order.shipped.v1, refund.completed.v1, customer.address.changed.v1.
- Определите partition key и ожидания порядка.
- Опишите Avro-схемы payload и envelope.
- Заполните карточки Событийного каталога (producers/consumers, retention, compat, security).
- Подготовьте golden payloads (валид/границы/невалид).
- Пропишите стратегию эволюции (какие поля могут добавляться, окна депрекейта).
- Сформируйте RTM-связки: UC/FR → события → подписчики → метрики.
Критерии зачёта:
- События минимальны и безопасны (без PII).
- Есть явный partition key и позиция по порядку.
- Схемы совместимы (backward), определена политика версий.
- Каталог заполнен (SLO, ACL, retention, observability).
- Идемпотентность консъюмеров предусмотрена (eventId/seq).
Шпаргалка (коротко)
- Событие = факт, команда = инструкция.
- Порядок только внутри partition key.
- At-least-once + идемпотентность → практический exactly-once.
- Outbox + CDC — стандарт, а не опция.
- Schema Registry + compat; новые поля — optional.
- Никаких PII в событиях.
- Каталог событий — обязательный артефакт: кто публикует/слушает, ключ, версия, SLO, ACL, retention.
- Метрики и логи: eventId, correlationId, latency, lag, DLQ.



