Domain-Driven Design Доменные события и интеграция контекстов
0%

Доменные события и интеграция контекстов

Доменные события и интеграция контекстов

В предыдущих статьях мы собирали модель внутри одной границы: выделили ограниченные контексты, построили агрегаты и спрятали хранение за репозиториями. Теперь — самое опасное место любой системы: связь между границами.

Заказ подтверждён — надо списать резерв на складе, выставить счёт, начислить бонусы, отправить письмо, обновить витрину аналитики. Если написать это в лоб, внутри метода подтверждения заказа появятся вызовы четырёх чужих сервисов, и контекст «Продажи» будет знать про биллинг, склад, почту и BI. Граница, которую мы так старательно проводили, испарится.

Доменные события позволяют одному контексту сообщить о факте, ничего не зная о том, кому этот факт нужен.


1. Интуиция: от вызова к факту

Сравним две формулировки одного требования. Императивная: «когда заказ подтверждают, вызови склад, потом биллинг, потом отправь письмо» — здесь Order начальник, который знает всех подчинённых, и добавление пятого потребителя правит код заказа. Событийная: «заказ был подтверждён», точка.

Второе — это язык бизнеса. Эксперт домена говорит именно так: «после того как платёж прошёл, мы…», «если груз задержан, то…». Прошедшее время в речи эксперта почти всегда указывает на доменное событие. Эванс добавил Domain Events в модель уже после первой книги; каноничный разбор — у Вернона в «Implementing Domain-Driven Design», гл. 8.

Численный аргумент. Пусть n контекстов интересуются событиями друг друга. При прямых вызовах каждый производитель знает каждого потребителя: до O(n²) связей, и каждое новое требование правит существующий код источника. При публикации событий источник знает только про свой контракт — O(n) связей, добавление потребителя не меняет ни строчки у производителя. Это не про скорость, а про стоимость изменения — тот же аргумент, что и в обзоре трека.

Цена, которую надо назвать сразу: вы меняете синхронную транзакционную согласованность на согласованность в конечном счёте плюс новый класс отказов — дубли, переупорядочивание, потерянные сообщения. Дальше вся статья про то, как платить эту цену осознанно.


2. Что такое доменное событие строго

Доменное событие — неизменяемый факт о том, что произошло в домене, значимый для эксперта предметной области. Из определения следуют жёсткие правила:

Свойство Почему
Имя в прошедшем времени: OrderConfirmed, PaymentFailed Факт случился, отменить нельзя — только компенсировать
Неизменяемость (immutable) Прошлое не редактируется; событие безопасно делить между обработчиками
occurred_at (часто и recorded_at) Порядок, отладка, метрика end-to-end задержки
event_id Дедупликация у потребителя
Идентификаторы, а не ссылки на объекты Событие переживёт процесс и сериализацию
Термин из единого языка Событие — часть единого языка, а не техническая деталь

Чем событие не является. Это не команда: SendEmail — команда, адресат обязан выполнить, отказ = ошибка; OrderConfirmed — факт, обработчиков может быть ноль, один или десять. Это не CRUD-нотификация: OrderRowUpdated бесполезно, потребитель не знает, что случилось, и вынужден догадываться по диффу. И это не транспорт: «событие» ≠ «сообщение в Kafka», Kafka — лишь способ доставки одного из видов событий.

Внутреннее событие ≠ интеграционное

Путаница здесь стоит командам месяцев:

  1. Доменное событие (внутреннее) — факт внутри одного контекста. Может содержать доменные типы (Money, OrderId), живёт в памяти процесса. Контракт приватный, меняется свободно.
  2. Интеграционное событие (внешнее) — факт, опубликованный в чужие контексты. Это публичный контракт: примитивные типы, версия, схема, обратная совместимость. На нём висят чужие команды.
  3. Команда — просьба что-то сделать, адресная, может быть отклонена.

Самая частая архитектурная ошибка — публиковать внутренние доменные события напрямую в Kafka: тогда переименование поля внутри агрегата ломает три чужих сервиса. Внутреннее событие нужно транслировать в интеграционное — это тот же antikorruption layer из карты контекстов, только на исходящей стороне.


3. Где рождается событие: внутри агрегата

Событие записывает тот, кто владеет инвариантом и совершает переход состояния, — агрегат. Не сервис приложения и не контроллер: иначе снова появится место, где состояние меняют, «забыв» про событие.

from dataclasses import dataclass, field
from datetime import datetime, timezone
from decimal import Decimal
from uuid import UUID, uuid4


@dataclass(frozen=True, slots=True)          # frozen => неизменяемость
class DomainEvent:
    event_id: UUID = field(default_factory=uuid4, kw_only=True)
    occurred_at: datetime = field(
        default_factory=lambda: datetime.now(timezone.utc), kw_only=True)


@dataclass(frozen=True, slots=True)
class OrderConfirmed(DomainEvent):
    order_id: str
    customer_id: str
    total: Decimal
    currency: str


class RecordsEvents:
    """Примесь «умею накапливать факты»."""

    def __init__(self) -> None:
        self._events: list[DomainEvent] = []

    def record(self, event: DomainEvent) -> None:
        self._events.append(event)

    def pull_events(self) -> list[DomainEvent]:
        """Забрать и очистить: факт уедет из агрегата ровно один раз."""
        events, self._events = self._events, []
        return events


class Order(RecordsEvents):
    def __init__(self, order_id: str, customer_id: str, total: Decimal) -> None:
        super().__init__()
        self.id, self.customer_id, self.total = order_id, customer_id, total
        self.status = "draft"

    def confirm(self) -> None:
        if self.status != "draft":
            raise ValueError(f"нельзя подтвердить заказ в статусе {self.status}")
        if self.total <= 0:
            raise ValueError("сумма заказа должна быть положительной")
        # 1) меняем состояние, 2) фиксируем факт — всегда в паре
        self.status = "confirmed"
        self.record(OrderConfirmed(order_id=self.id, customer_id=self.customer_id,
                                   total=self.total, currency="RUB"))

Ключевая деталь: событие записывается только вместе с успешным переходом. Нарушен инвариант — метод бросил исключение, события нет. Связка «состояние ↔ факт» становится неразрывной на уровне кода, а не дисциплины разработчика.

Событие не должно быть эхом сеттера. Плохо: OrderStatusChanged(old="draft", new="confirmed") — потребитель вынужден писать if new == ..., то есть знать чужую модель состояний. Хорошо: отдельные OrderConfirmed, OrderCancelled, OrderShipped. Типов больше, зато каждый потребитель подписан ровно на нужное и не ломается при добавлении статуса.


4. Когда публиковать: до или после коммита

Здесь ломается больше всего продовых систем.

Аксиомы, которые стоит запомнить:

  • Собирать события можно и нужно до коммита — они уже лежат в агрегате; доставлять наружу — только после успешного коммита, иначе склад зарезервирует товар под заказ, который откатился.
  • Обработчики, которым нужна та же транзакция (пересчёт денормализованной таблицы внутри контекста), могут работать внутри неё — но тогда их падение откатывает бизнес-операцию. Это осознанный выбор, а не случайность.
  • Обработчик не должен молча менять другой агрегат в той же транзакции. Правило Вернона: одна транзакция — один агрегат, остальное через события.

5. Двойная запись и транзакционный outbox

Наивное «после коммита опубликуем в Kafka» содержит фундаментальную ошибку — dual write: две независимые системы, атомарности между ними нет.

Проблема двойной записи и её решение через outbox

Распределённые транзакции (XA/2PC) формально решают задачу, но на практике почти не применяются: блокируют ресурсы, плохо переживают сетевые разделы, не поддерживаются большинством брокеров. Промышленный ответ — Transactional Outbox (microservices.io).

-- Outbox живёт в ТОЙ ЖЕ базе, что и агрегаты: только так возможна одна транзакция.
CREATE TABLE outbox (
    id             BIGSERIAL PRIMARY KEY,   -- монотонный порядок публикации
    event_id       UUID        NOT NULL UNIQUE,
    aggregate_type TEXT        NOT NULL,
    aggregate_id   TEXT        NOT NULL,    -- станет ключом партиционирования
    event_type     TEXT        NOT NULL,    -- 'sales.order.confirmed'
    event_version  INT         NOT NULL DEFAULT 1,
    payload        JSONB       NOT NULL,
    headers        JSONB       NOT NULL DEFAULT '{}',  -- trace_id, correlation_id
    occurred_at    TIMESTAMPTZ NOT NULL,
    published_at   TIMESTAMPTZ,             -- NULL = ещё не отправлено
    attempts       INT         NOT NULL DEFAULT 0
);

-- Частичный индекс: релей сканирует только «хвост», а не всю таблицу.
CREATE INDEX outbox_unpublished_idx ON outbox (id) WHERE published_at IS NULL;
class SqlUnitOfWork:
    """Транзакционная граница: агрегаты и outbox фиксируются вместе.
    __init__ хранит conn, dispatcher и список отслеживаемых агрегатов,
    __enter__ открывает BEGIN — опущены для краткости."""

    def __exit__(self, exc_type, exc, tb) -> None:      # commit / rollback
        if exc_type is not None:
            self._conn.execute("ROLLBACK")
            return
        pending = [e for agg in self._tracked for e in agg.pull_events()]
        for e in pending:
            public = translate_to_integration(e)   # внутреннее → публичный контракт
            if public is None:                     # не всякое событие уезжает наружу
                continue
            self._conn.execute(
                """INSERT INTO outbox (event_id, aggregate_type, aggregate_id,
                                       event_type, event_version, payload, occurred_at)
                   VALUES (%s, %s, %s, %s, %s, %s, %s)""",
                (str(e.event_id), public.aggregate_type, public.aggregate_id,
                 public.type, public.version, json.dumps(public.payload), e.occurred_at))
        self._conn.execute("COMMIT")
        for event in pending:                      # только теперь — локальные обработчики
            self._dispatcher.dispatch(event)

Релей бывает двух видов:

Способ Как работает Плюсы Минусы
Polling publisher SELECT … WHERE published_at IS NULL ORDER BY id LIMIT n FOR UPDATE SKIP LOCKED в цикле Просто, никакой доп. инфраструктуры Задержка = период опроса; нагрузка на БД
CDC / log tailing Debezium читает WAL/binlog и льёт строки outbox в Kafka Задержка десятки мс, нет нагрузки запросами Нужен Kafka Connect, сложнее эксплуатация
def relay_batch(conn, producer, batch: int = 200) -> int:
    """Один цикл релея. SKIP LOCKED позволяет запускать N инстансов параллельно."""
    rows = conn.execute(
        """SELECT id, event_id, aggregate_id, event_type, payload
             FROM outbox WHERE published_at IS NULL
            ORDER BY id LIMIT %s FOR UPDATE SKIP LOCKED""", (batch,)).fetchall()
    for row in rows:
        producer.send(topic=row["event_type"],
                      key=row["aggregate_id"],   # порядок в пределах одного агрегата
                      value=row["payload"],
                      headers={"event_id": row["event_id"]})
        # Падение здесь => строка не помечена => событие уедет второй раз.
        # Это и есть at-least-once; поэтому потребитель обязан дедуплицировать.
        conn.execute("UPDATE outbox SET published_at = now() WHERE id = %s", (row["id"],))
    conn.commit()
    return len(rows)

Сложность. Запись: O(1) дополнительных вставок на событие в той же транзакции — заметно дешевле 2PC. Релей: O(k) на цикл, где k — размер батча; благодаря частичному индексу стоимость не зависит от общего числа строк. Память на дедупликацию у потребителя — O(w), где w — число event_id в окне хранения; окно должно быть заведомо больше максимальной задержки повторной доставки (типично 24–72 часа). И обязательно ретеншен: outbox без уборки становится самой большой таблицей в базе — партиционирование по дате плюс DROP PARTITION дешевле, чем DELETE.


6. Идемпотентность потребителя и inbox

At-least-once означает: каждый потребитель рано или поздно получит дубль. Не «может быть» — получит. Значит, обработка обязана быть идемпотентной. Три уровня защиты, от лучшего к худшему:

  1. Естественная идемпотентность — операция сама повторяема: SET status = 'confirmed' (а не counter = counter + 1), INSERT … ON CONFLICT DO NOTHING по бизнес-ключу. Ничего дополнительного не нужно.
  2. Inbox — таблица обработанных сообщений; универсально, работает для любой логики.
  3. Проверка «уже сделано» по состоянию агрегата — если состояние однозначно кодирует факт обработки.
def handle_order_confirmed(conn, message) -> None:
    """Потребитель в контексте «Склад»: резервирует товар под подтверждённый заказ."""
    with conn.transaction():                    # inbox и эффект — одна транзакция!
        try:
            conn.execute(
                "INSERT INTO inbox (event_id, consumer, processed_at) VALUES (%s, %s, now())",
                (message.headers["event_id"], "warehouse.reservation"))
        except UniqueViolation:
            return                              # дубль: молча выходим, offset коммитится

        payload = translate_from_published_language(message.value)   # ACL на входе
        reserve_stock(conn, order_id=payload.order_id, lines=payload.lines)

Таблица тривиальна: PRIMARY KEY (consumer, event_id) — у каждого обработчика своя дедупликация, плюс индекс по processed_at для чистки старых записей. Критично другое: запись в inbox и бизнес-эффект должны быть в одной локальной транзакции. Пометите сообщение обработанным до эффекта и упадёте — событие потеряно навсегда.

Порядок. Kafka гарантирует его только внутри партиции, поэтому ключ партиционирования — aggregate_id: все события одного заказа придут по порядку. Между разными агрегатами порядка нет и быть не должно — не проектируйте логику, которая на него опирается. Если порядок всё же нарушится (ребаланс, ретрай, перепубликация), спасает версия агрегата в событии: проекция обновляется условно — UPDATE projection SET version = :v WHERE id = :id AND version < :v, и пришедшее с опозданием v5 при уже применённом v7 просто не даст изменённых строк, то есть будет отброшено.


7. Стили событийной интеграции

Мартин Фаулер разделил три вещи, которые все называют «event-driven» («What do you mean by Event-Driven?»).

Три стиля событийной интеграции

Компромисс, который чаще всего работает: тонкое событие плюс необходимый минимум состояния. Кладите поля, которые нужны большинству потребителей и семантически принадлежат факту (сумма заказа принадлежит факту подтверждения), но не выгружайте весь агрегат.

Отдельно про event sourcing: это способ хранить состояние агрегата как поток событий, ортогональный интеграции. Публиковать интеграционные события без event sourcing можно (обычно так и делают) и наоборот. «Мы используем доменные события» ≠ «у нас event sourcing». Разбор — у Фаулера в «Event Sourcing».


8. Контракты: published language и версионирование

Интеграционное событие — публичный API, причём сложнее REST: нельзя «попросить всех обновиться завтра». Published Language — согласованный формат обмена (JSON Schema, Avro, Protobuf), принадлежащий не одному сервису, а связке; живёт в отдельном репозитории или schema registry.

# sales/order-confirmed-v2.schema.yaml — published language, а не внутренняя модель
$schema: "https://json-schema.org/draft/2020-12/schema"
title: sales.order.confirmed
type: object
required: [event_id, occurred_at, order_id, customer_id, total]
properties:
  event_id:    { type: string, format: uuid }
  occurred_at: { type: string, format: date-time }
  order_id:    { type: string }
  customer_id: { type: string }
  total:                                      # деньги — строкой, без float
    type: object
    properties:
      amount:   { type: string }
      currency: { type: string, pattern: "^[A-Z]{3}$" }
  delivery_slot: { type: [string, "null"] }   # добавлено в v2, опциональное
additionalProperties: true                    # незнакомое потребитель игнорирует

Правила эволюции, проверенные болью:

  • В рамках мажорной версии — только совместимые изменения. Добавлять опциональные поля можно; удалять, переименовывать, сужать тип, менять смысл — нельзя.
  • Tolerant reader: потребитель читает нужные поля и игнорирует остальные (Fowler).
  • Несовместимое изменение = новый тип/топик …confirmed.v2. Обе версии публикуются параллельно, пока метрики не покажут ноль потребителей v1.
  • Деньги — строкой или целым в минимальных единицах. float в JSON — гарантированный баг в биллинге.
  • Апкастинг на входе потребителя: старые версии приводятся к текущей одной функцией, чтобы бизнес-логика осталась в единственном экземпляре, без if version == … по всему коду.

9. Саги и процесс-менеджеры

События хороши для «реагируй на факт». Но процесс, растянутый на несколько контекстов («оформление заказа» = резерв + оплата + отгрузка), требует ещё и компенсаций: не прошла оплата — снять резерв. Транзакции на всю цепочку нет, есть сага (microservices.io).

Хореография — каждый сервис слушает события и реагирует: нет центрального узла, минимальная связность, но процесс нигде не описан целиком и его приходится собирать из логов пяти сервисов. Хорошо для 2–3 шагов. Оркестрация (process manager) — отдельный компонент хранит состояние саги и рассылает команды: процесс виден в одном месте, легко добавить таймауты и ретраи; цена — новый сервис и риск получить «бога процессов». Хорошо для 4+ шагов с компенсациями.

class OrderFulfillmentSaga:
    """Состояние персистентно, реакции идемпотентны. Сага НЕ содержит
    бизнес-правил контекстов — только порядок шагов и компенсации."""

    def on_order_confirmed(self, e) -> list["Command"]:
        self.state = "AWAITING_STOCK"
        self.timeout_at = e.occurred_at + timedelta(minutes=15)
        return [ReserveStock(order_id=e.order_id, lines=e.lines)]

    def on_payment_failed(self, e) -> list["Command"]:
        if self.state in ("FAILED", "COMPLETED"):
            return []                       # поздний дубль игнорируем
        self.state = "COMPENSATING"
        return [ReleaseStock(order_id=e.order_id, reason="payment_failed")]

Компенсация — не rollback. Отменённый резерв, возврат средств, письмо «извините» — полноценные доменные факты, которые эксперт понимает и которые остаются в истории. Проектируйте их вместе с бизнесом, а не придумывайте в коде.


10. Согласованность в конечном счёте — продуктовое решение

Технический вопрос «а не увидит ли пользователь устаревшие данные?» на деле бизнес-вопрос: какая задержка допустима и что делать, если правило нарушится. Вернон предлагает простой приём — спросить эксперта: «если склад узнает о заказе через 3 секунды, что случится?»

  • «Ничего» → смело асинхронно.
  • «Может уйти в минус на пару штук, дозакажем» → асинхронно плюс процедура разрешения. Овербукинг — нормальная бизнес-практика, авиакомпании живут на ней десятилетиями.
  • «Никогда, это регуляторное требование» → значит, эти сущности лежат в одном агрегате и одной транзакции, и граница контекста проведена неверно.

Приёмы для UI: статус «обрабатывается», read-your-writes через локальную проекцию контекста-источника, correlation_id клиенту, чтобы он мог опросить результат.


11. Наблюдаемость и тестирование

Событийная система без наблюдаемости — чёрный ящик. Минимум, окупающийся с первого инцидента: correlation_id (общий на всю цепочку от HTTP-запроса до последнего потребителя) и causation_id (id события, породившего текущее) — вместе они дают дерево причинности; метрики размера неопубликованного хвоста outbox, consumer lag и задержки от occurred_at до обработки; DLQ с алертом, иначе одно «ядовитое» сообщение останавливает партицию.

Тесты ложатся на три уровня.

def test_confirm_records_event():
    """1. Юнит: агрегат фиксирует факт. Без БД и брокера — миллисекунды."""
    order = Order("ord-1", "cus-1", Decimal("100.00"))
    order.confirm()
    assert isinstance(order.pull_events()[0], OrderConfirmed)
    assert order.pull_events() == []          # второй вызов пуст: события забраны


def test_rollback_publishes_nothing(conn, dispatcher_spy):
    """2. Интеграция: падение транзакции не выпускает событий наружу."""
    with pytest.raises(ValueError):
        with SqlUnitOfWork(conn, dispatcher_spy) as uow:
            order = Order("ord-2", "cus-1", Decimal("0"))
            uow.track(order)
            order.confirm()                   # бросит: сумма не положительна
    assert dispatcher_spy.dispatched == [] and count_rows(conn, "outbox") == 0


def test_handler_is_idempotent(conn, message):
    """3. Контракт потребителя: повторная доставка не меняет результат."""
    handle_order_confirmed(conn, message)
    handle_order_confirmed(conn, message)     # тот же event_id
    assert reserved_quantity(conn, "ord-3") == 1

Четвёртый уровень — contract testing (Pact, проверка совместимости схем в CI): производитель физически не может задеплоить схему, ломающую зарегистрированных потребителей.


12. Типичные ошибки

  1. Публикация внутренних событий наружу — внутренняя модель утекает в контракт, рефакторинг становится невозможен. Лечение: слой трансляции (OHS/published language).
  2. Публикация до коммита — потребители реагируют на факты, которых не случилось: фантомные резервы и письма «ваш заказ отправлен» про отменённый заказ.
  3. Отсутствие идемпотентности. «У нас же exactly-once в Kafka» — нет: она работает только внутри Kafka-транзакций между топиками, а не для побочных эффектов в вашей БД.
  4. CRUD-события. EntityUpdated заставляет потребителя реконструировать намерение; событие должно называть бизнес-факт.
  5. Событие как замаскированная команда (OrderConfirmedSoSendEmail). Признак: производитель обижается, когда обработчик не сработал.
  6. Цепочка событий вместо явного процесса — никто не знает, в каком состоянии процесс и как его починить после сбоя; для 4+ шагов нужна сага.
  7. Синхронное изменение нескольких агрегатов в одном обработчике — возвращает распределённый монолит.
  8. Нет ретеншена у outbox/inbox — через год сотни ГБ и деградация записи.
  9. Расчёт на глобальный порядок (его нет) и события вместо запросов: если данные нужны «прямо сейчас и точно свежие» — это синхронный запрос, а не событие.

13. Как это выглядит в проде

  • Хранение и доставка: outbox в основной OLTP-базе с партиционированием по суткам и ретеншеном 7–30 дней; Debezium → Kafka для высоконагруженных потоков, polling-релей там, где задержка в секунду допустима (это большинство систем).
  • Схемы: Confluent Schema Registry или репозиторий JSON Schema плюс проверка совместимости в CI (про compatibility types).
  • Именование топиков: <контекст>.<агрегат>.<событие>.v<N>, например sales.order.confirmed.v2 — версия в имени упрощает параллельную работу двух контрактов.
  • Ретраи: экспоненциальный backoff с джиттером, лимит попыток, затем DLQ и алерт дежурному.
  • Каркасы: в .NET — MediatR для внутрипроцессной диспетчеризации и MassTransit (встроены outbox и saga state machine); в Java — Spring Modulith с транзакционным журналом событий; в Go и Python обычно пишут тонкий слой руками — его там немного.
  • Модульный монолит как промежуточный шаг: те же доменные события внутри одного процесса и одной БД. Вы получаете развязку модулей без операционной сложности брокера, а когда модуль реально понадобится выделить — контракт событий уже есть. Лучший старт для большинства команд.

Мини-итог

  • Доменное событие — неизменяемый факт в прошедшем времени, записываемый агрегатом вместе с переходом состояния и названный словом из единого языка.
  • Разделяйте внутренние события (приватные, свободно меняются) и интеграционные (публичный контракт, версионируется, транслируется через published language).
  • Публикуйте только после коммита, решайте двойную запись через transactional outbox; из at-least-once автоматически следует обязательная идемпотентность потребителя (inbox).
  • Порядок гарантирован только внутри партиции: ключ = aggregate_id, плюс версия агрегата для отбрасывания устаревших событий.
  • Многошаговые межконтекстные процессы оформляйте сагой с явными компенсациями, а согласованность в конечном счёте обсуждайте с бизнесом, а не решайте втихую в коде.

Источники


Что дальше

Мы разобрали, как события живут в коде и в инфраструктуре. Но откуда берётся сам список событий домена? Его не выдумывают у доски в одиночку — его добывают вместе с экспертами.

Event Storming и практики моделирования домена — как за один воркшоп получить карту событий, команд, агрегатов и границ контекстов и почему именно события оказываются лучшей отправной точкой для моделирования.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов