Архитектурные паттерны CQRS и Event Sourcing
0%

CQRS и Event Sourcing

CQRS и Event Sourcing

Есть пара паттернов, которые в разговорах почти всегда произносят одним словом — «сикьюэрэс-и-ивентсорсинг», будто это единая технология с одним переключателем. Из-за этого возникают две симметричные беды. Одни команды берут event sourcing там, где нужно было всего лишь отдельная денормализованная таблица для отчётов, и через год тонут в миграциях схемы событий. Другие боятся сказать слово «CQRS», потому что «мы не готовы к eventual consistency», хотя CQRS сам по себе никакой eventual consistency не требует.

Поэтому первое и главное утверждение этой статьи:

CQRS и Event Sourcing — это два разных, ортогональных паттерна. CQRS можно применять без event sourcing (и это самый частый полезный случай). Event sourcing без CQRS технически возможен, но почти всегда неудобен — и именно эта односторонняя связь порождает иллюзию, что они одно и то же.

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

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


1. Откуда вообще берётся идея разделять чтение и запись

Начнём с наблюдения, которое не требует никаких паттернов. Возьмите любую нетривиальную предметную область — скажем, интернет-магазин — и посмотрите на две операции.

Операция записи «оформить заказ». Ей нужно: проверить, что корзина не пуста; что товар в наличии; что промокод ещё действует; что адрес доставки валиден для региона; что у клиента нет незакрытой задолженности. Это работа с инвариантами — правилами, которые обязаны выполняться всегда. Чтобы их проверить, нужна модель, где данные собраны в связный, нормализованный, консистентный граф объектов.

Операция чтения «страница списка заказов в админке». Ей нужно: 40 заказов на страницу, у каждого — имя клиента, сумма, статус, город доставки, название службы, флажок «есть спор». Никаких инвариантов. Никакого поведения. Нужен один быстрый плоский результат, который в нормализованной модели собирается джойном шести таблиц или, что хуже, N+1 обращениями через ORM.

Эти две потребности тянут модель данных в противоположные стороны. Нормализация хороша для записи (одно место истины — нет рассогласования). Денормализация хороша для чтения (данные лежат ровно в той форме, в которой их спрашивают). Пока модель одна, вы вечно ищете компромисс и получаете структуру, которая плоха для обеих задач: объекты домена, обвешанные полями «для UI», и запросы, которые тащат из БД в двадцать раз больше, чем показывают.

Добавьте сюда асимметрию нагрузки. В типичной e-commerce системе на одну запись приходятся десятки-сотни чтений. Масштабировать чтение и запись хочется независимо и по-разному: чтение — репликами и кешем, запись — шардированием по ключу агрегата.

CQRS (Command Query Responsibility Segregation) — это ответ ровно на это: разделить модель для изменения состояния и модель для чтения состояния, позволив каждой развиваться и оптимизироваться независимо. Термин ввёл Greg Young, отталкиваясь от более старого принципа CQS Бертрана Мейера (Command-Query Separation): метод либо изменяет состояние, либо возвращает данные, но не то и другое сразу. CQS — про методы; CQRS — про целые модели и их хранилища.

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

1.1. Лестница CQRS: четыре ступени, а не переключатель

Огромная часть споров о CQRS происходит от того, что спорщики стоят на разных ступенях лестницы.

Ступень Что разделено Консистентность Цена
0. CQS Методы внутри класса Строгая Нулевая, просто дисциплина
1. Разные модели, одно хранилище Классы/DTO: агрегаты для записи, read-DTO для чтения Строгая Низкая: чуть больше кода
2. Разные схемы, одна БД Нормализованные таблицы + денормализованные проекции, обновляемые в той же транзакции Строгая Средняя: проекции надо писать и тестировать
3. Разные хранилища, асинхронно PostgreSQL для записи, Elasticsearch/Redis/ClickHouse для чтения Итоговая (eventual) Высокая: лаг, идемпотентность, пересборка, мониторинг

Ступень 1 — то, что стоит делать почти всегда, и это буквально «не тащить сущности ORM в контроллер, а писать отдельные read-запросы». Ступень 3 — то, что называют CQRS в докладах, и то, ради чего написана половина этой статьи.

Практическое правило: поднимайтесь на ступень выше только тогда, когда предыдущая уже не справляется с измеренной, а не воображаемой проблемой. Udi Dahan в «Clarified CQRS» формулирует это резко: если вам не нужна другая база под чтение, не делайте другую базу под чтение.

1.2. Как выглядит ступень 1–2 в коде

Смотрите: тут нет ни шины, ни брокера, ни событий. Только два разных пути в одном сервисе.

# ---------- СТОРОНА ЗАПИСИ: богатая модель, защищающая инварианты ----------

from dataclasses import dataclass, field
from decimal import Decimal
from uuid import UUID

class DomainError(Exception):
    """Нарушение бизнес-правила — это не 500, а 409/422."""

@dataclass
class Order:
    id: UUID
    customer_id: UUID
    lines: list["OrderLine"] = field(default_factory=list)
    status: str = "draft"

    # Инвариант живёт ЗДЕСЬ, а не в контроллере и не в триггере БД.
    MAX_LINES = 100

    def add_line(self, sku: str, qty: int, price: Decimal) -> None:
        if self.status != "draft":
            raise DomainError("нельзя менять состав подтверждённого заказа")
        if qty <= 0:
            raise DomainError("количество должно быть положительным")
        if len(self.lines) >= self.MAX_LINES:
            raise DomainError(f"в заказе не может быть больше {self.MAX_LINES} позиций")
        self.lines.append(OrderLine(sku, qty, price))

    def place(self) -> None:
        if not self.lines:
            raise DomainError("нельзя оформить пустой заказ")
        self.status = "placed"

@dataclass
class OrderLine:
    sku: str
    qty: int
    price: Decimal


# ---------- СТОРОНА ЧТЕНИЯ: никаких объектов домена, только SQL и плоский DTO ----------

@dataclass(frozen=True)
class OrderListRow:
    order_id: UUID
    customer_name: str
    total: Decimal
    status: str
    city: str

async def list_orders(conn, *, limit: int, offset: int) -> list[OrderListRow]:
    # Один запрос, ровно те колонки, что нужны экрану. ORM здесь не участвует
    # намеренно: read-модель не обязана совпадать с моделью домена.
    rows = await conn.fetch(
        """
        SELECT o.id, c.name AS customer_name, o.total, o.status, a.city
        FROM   orders o
        JOIN   customers c ON c.id = o.customer_id
        JOIN   addresses  a ON a.id = o.shipping_address_id
        ORDER  BY o.placed_at DESC
        LIMIT $1 OFFSET $2
        """,
        limit, offset,
    )
    return [OrderListRow(**dict(r)) for r in rows]

Это уже CQRS. Никакой инфраструктуры, а половина боли (толстые ORM-выборки, домен, изуродованный требованиями UI) снята. Дальше начинается интересное.


2. Event Sourcing: состояние как производная от журнала фактов

Теперь второй паттерн, полностью независимый от первого.

2.1. Интуиция

Ваш банк не хранит поле balance и не обновляет его при каждой операции. Он хранит проводки, а баланс вычисляет. Это не техническое решение, а требование учёта: нужно уметь ответить не только «сколько сейчас», но и «почему столько», «сколько было на 3 марта» и «кто и когда это изменил». Такая книга называется append-only: строки в неё дописываются, но никогда не переписываются; ошибочная проводка исправляется сторнирующей проводкой, а не ластиком.

Тот же приём вы используете каждый день: git не хранит текущее состояние файлов как истину — он хранит коммиты, а рабочая копия есть результат их применения. WAL в PostgreSQL, journal в ext4, redo log в Oracle — везде одна и та же конструкция: истина — это упорядоченный журнал изменений, а текущее состояние — материализованный кеш этого журнала.

Event Sourcing переносит эту конструкцию в прикладной домен: вместо строки в таблице orders, которую мы обновляем, мы храним последовательность событий этого заказа — OrderPlaced, ItemAdded, PaymentReceived, OrderShipped — а состояние получаем свёрткой.

Состояние как свёртка журнала событий

2.2. Строгое определение

Пусть E — множество типов событий, S — множество состояний агрегата. Event sourcing задаётся двумя чистыми функциями:

evolve : S × E → S           # как событие меняет состояние (никаких побочных эффектов)
decide : S × Command → [E]   # какие факты порождает команда в данном состоянии

и правилом восстановления:

state(n) = fold(evolve, S₀, [e₁, e₂, …, eₙ])

Отсюда сразу три следствия, которые надо запомнить:

  1. evolve не имеет права отказать. Событие — это уже случившийся факт, прошедшее время. ItemAdded нельзя «не принять»: он уже в журнале. Вся валидация происходит в decide, до записи.
  2. decide не имеет права ничего менять. Она возвращает список событий; их применение и запись — снаружи. Это делает домен тестируемым без единого мока.
  3. Именование событий в прошедшем времени — не стилистика, а семантика. CreateOrder — команда (может быть отвергнута), OrderCreated — событие (обсуждению не подлежит).

Мартин Фаулер описывает базовую конструкцию в Event Sourcing; каноническое изложение с прикладной стороны — CQRS Documents Грега Янга (PDF, ~50 страниц, читается за вечер и стоит того).

2.3. Жизненный цикл агрегата как автомат

Полезная проверка на здравость: если ваши события — это переходы конечного автомата, модель, скорее всего, хороша. Если события выглядят как OrderUpdated с полным снимком полей, вы просто переизобрели UPDATE и потеряли все преимущества.

2.4. Рабочая реализация

Ниже — минимальный, но честный event-sourced агрегат: чистое ядро плюс репозиторий с оптимистической блокировкой.

from __future__ import annotations
from dataclasses import dataclass, replace
from decimal import Decimal
from typing import Iterable
from uuid import UUID

# ---------- События: неизменяемые факты в прошедшем времени ----------

@dataclass(frozen=True)
class CartCreated:
    cart_id: UUID
    customer_id: UUID

@dataclass(frozen=True)
class ItemAdded:
    sku: str
    qty: int
    price: Decimal

@dataclass(frozen=True)
class ItemRemoved:
    sku: str
    qty: int

@dataclass(frozen=True)
class CartCheckedOut:
    total: Decimal

Event = CartCreated | ItemAdded | ItemRemoved | CartCheckedOut

# ---------- Состояние: результат свёртки, а не то, что хранится ----------

@dataclass(frozen=True)
class Cart:
    id: UUID | None = None
    items: dict[str, int] = None          # sku -> qty
    prices: dict[str, Decimal] = None
    checked_out: bool = False

EMPTY = Cart(items={}, prices={})

def evolve(state: Cart, event: Event) -> Cart:
    """Чистая функция. Не валидирует — факт уже произошёл."""
    match event:
        case CartCreated(cart_id=cid):
            return replace(state, id=cid, items={}, prices={})
        case ItemAdded(sku=sku, qty=q, price=p):
            items = {**state.items, sku: state.items.get(sku, 0) + q}
            return replace(state, items=items, prices={**state.prices, sku: p})
        case ItemRemoved(sku=sku, qty=q):
            left = state.items.get(sku, 0) - q
            items = {**state.items}
            if left > 0:
                items[sku] = left
            else:
                items.pop(sku, None)
            return replace(state, items=items)
        case CartCheckedOut():
            return replace(state, checked_out=True)
    return state

def fold(events: Iterable[Event], state: Cart = EMPTY) -> Cart:
    for e in events:
        state = evolve(state, e)
    return state

# ---------- Решения: вся валидация здесь ----------

class DomainError(Exception): ...

MAX_DISTINCT_SKUS = 50

def decide_add_item(state: Cart, sku: str, qty: int, price: Decimal) -> list[Event]:
    if state.checked_out:
        raise DomainError("корзина уже оформлена")
    if qty <= 0:
        raise DomainError("количество должно быть положительным")
    if sku not in state.items and len(state.items) >= MAX_DISTINCT_SKUS:
        raise DomainError(f"не больше {MAX_DISTINCT_SKUS} различных позиций")
    return [ItemAdded(sku, qty, price)]

def decide_checkout(state: Cart) -> list[Event]:
    if state.checked_out:
        return []                      # идемпотентность: повтор не порождает событий
    if not state.items:
        raise DomainError("нельзя оформить пустую корзину")
    total = sum(state.prices[s] * q for s, q in state.items.items())
    return [CartCheckedOut(total=Decimal(total))]

Обратите внимание на decide_checkout: повторная команда возвращает пустой список событий, а не ошибку. Это единственно верное поведение для команд, приходящих по сети с ретраями, — подробно про идемпотентность будет в статье «Saga, распределённые транзакции, outbox и идемпотентность».

2.5. Event store: схема и оптимистическая блокировка

Event store — это, по сути, таблица с двумя ограничениями: только вставки и уникальность (stream_id, version). Последнее и есть механизм конкурентного доступа.

-- Журнал событий. Единственный источник истины.
CREATE TABLE events (
    global_position BIGSERIAL PRIMARY KEY,   -- глобальный порядок для проекций
    stream_id       UUID        NOT NULL,    -- один поток = один агрегат
    version         INT         NOT NULL,    -- позиция внутри потока, с 1
    event_type      TEXT        NOT NULL,    -- 'ItemAdded'
    event_version   INT         NOT NULL DEFAULT 1,  -- версия СХЕМЫ события
    payload         JSONB       NOT NULL,
    metadata        JSONB       NOT NULL,    -- correlation_id, causation_id, user_id
    occurred_at     TIMESTAMPTZ NOT NULL DEFAULT now(),

    -- Сердце оптимистической блокировки: две конкурентные записи
    -- одной версии одного потока не могут обе выиграть.
    CONSTRAINT uq_stream_version UNIQUE (stream_id, version)
);

CREATE INDEX idx_events_stream ON events (stream_id, version);

-- Снимки — кеш, а не истина. Таблицу можно очистить целиком в любой момент.
CREATE TABLE snapshots (
    stream_id UUID PRIMARY KEY,
    version   INT   NOT NULL,
    state     JSONB NOT NULL
);

-- Позиции проекций. У каждой проекции своя, независимая.
CREATE TABLE projection_checkpoints (
    projection_name TEXT PRIMARY KEY,
    last_position   BIGINT NOT NULL DEFAULT 0,
    updated_at      TIMESTAMPTZ NOT NULL DEFAULT now()
);

Репозиторий поверх этой схемы:

import json

class ConcurrencyError(Exception):
    """Кто-то записал в поток, пока мы принимали решение."""

class EventStore:
    def __init__(self, conn):
        self.conn = conn

    async def load(self, stream_id) -> tuple[Cart, int]:
        """Восстановление: снимок + хвост событий. O(n − k) по событиям."""
        snap = await self.conn.fetchrow(
            "SELECT version, state FROM snapshots WHERE stream_id = $1", stream_id)
        state, version = (deserialize_state(snap["state"]), snap["version"]) if snap else (EMPTY, 0)

        rows = await self.conn.fetch(
            "SELECT version, event_type, event_version, payload FROM events "
            "WHERE stream_id = $1 AND version > $2 ORDER BY version",
            stream_id, version)
        for r in rows:
            # upcast поднимает старую схему события до актуальной — см. раздел 5
            state = evolve(state, upcast(r["event_type"], r["event_version"], json.loads(r["payload"])))
            version = r["version"]
        return state, version

    async def append(self, stream_id, expected_version: int, events: list[Event], meta: dict) -> int:
        """Атомарная дозапись. Конфликт версий ловим через UNIQUE-констрейнт."""
        if not events:
            return expected_version                  # нечего писать — не открываем транзакцию
        async with self.conn.transaction():
            v = expected_version
            for e in events:
                v += 1
                try:
                    await self.conn.execute(
                        "INSERT INTO events (stream_id, version, event_type, event_version, payload, metadata) "
                        "VALUES ($1,$2,$3,$4,$5,$6)",
                        stream_id, v, type(e).__name__, CURRENT_SCHEMA[type(e).__name__],
                        json.dumps(serialize(e)), json.dumps(meta))
                except UniqueViolationError:
                    raise ConcurrencyError(
                        f"поток {stream_id} изменён параллельно на версии {v}")
            return v

# ---------- Обработчик команды: load → decide → append, с ретраем ----------

async def handle_add_item(store: EventStore, cart_id, sku, qty, price, meta, attempts=3):
    for attempt in range(attempts):
        state, version = await store.load(cart_id)
        events = decide_add_item(state, sku, qty, price)   # чистая функция, легко тестируется
        try:
            return await store.append(cart_id, version, events, meta)
        except ConcurrencyError:
            if attempt == attempts - 1:
                raise
            # Повторяем ВЕСЬ цикл: решение принималось на устаревшем состоянии,
            # и оно могло стать невалидным (например, лимит позиций уже выбран).
    raise AssertionError("недостижимо")

Разбор поведения при конфликте — важнее, чем кажется:

Почему ретрай обязан перечитывать состояние, а не просто повторять INSERT с version+1. Между вашей загрузкой и записью кто-то мог оформить корзину. Слепой повтор впишет ItemAdded в оформленную корзину — то есть нарушит инвариант, ради защиты которого весь механизм и строился. Это самая частая ошибка в самописных event store.

2.6. Сложность и снимки

  • Восстановление агрегата: O(n) по времени, где n — число событий в потоке; O(1) дополнительной памяти при потоковой свёртке.
  • Со снимком на версии k: O(n − k).
  • Запись: O(m) вставок для m новых событий, плюс одна проверка уникальности.
  • Пересборка проекции: O(N) по всему журналу — это линейно, но N глобальное, и на 10⁹ событий пересборка может занимать часы. Планируйте это заранее.

Когда делать снимок? Практическое эмпирическое правило: каждые 50–200 событий потока, и только если профилирование показало, что загрузка стала узким местом. Снимок никогда не является источником истины: при изменении структуры состояния вы просто удаляете все снимки — и система продолжает работать, лишь медленнее, пока снимки не набегут заново. Это свойство надо явно проверять в тестах: «удалили таблицу snapshots — всё работает».

Главный запах: бесконечно растущий поток. Если у агрегата нет естественного финала (поток user-123 живёт годами и накопил 400 000 событий), это ошибка моделирования границ, а не повод для снимков. Обычно лечится сменой единицы потока: не account-42, а account-42-2026-07 (месячный период с переносом остатка), не device-7, а device-7-session-88. У Оскара Дудыча есть подробный разбор этого приёма — «closing the books» / temporal modelling.


3. Где два паттерна встречаются: проекции

Теперь соединим. Если состояние хранится как журнал, то запрос «покажи 40 последних заказов, отсортированных по сумме» становится физически невыполнимым по журналу: пришлось бы свернуть все потоки. Значит, для чтения нужна отдельная материализованная модель. Вот почему event sourcing практически всегда тянет за собой CQRS: не по идеологии, а по арифметике.

Проекция (projection, read model, view) — это долгоживущий процесс, который читает журнал в глобальном порядке и поддерживает денормализованную структуру под конкретный сценарий чтения.

class CartSummaryProjection:
    """Плоская витрина: одна строка на корзину, всё нужное экрану."""

    NAME = "cart_summary_v3"     # версия в имени — чтобы пересобирать рядом со старой
    HANDLES = {"CartCreated", "ItemAdded", "ItemRemoved", "CartCheckedOut"}

    async def run(self, conn, batch_size: int = 500):
        while True:
            pos = await conn.fetchval(
                "SELECT last_position FROM projection_checkpoints WHERE projection_name=$1",
                self.NAME)
            rows = await conn.fetch(
                "SELECT global_position, stream_id, event_type, event_version, payload "
                "FROM events WHERE global_position > $1 "
                "ORDER BY global_position LIMIT $2", pos, batch_size)
            if not rows:
                await asyncio.sleep(0.05)
                continue

            # КРИТИЧНО: применение батча и сдвиг чекпоинта — одна транзакция.
            # Иначе после падения между ними мы либо потеряем, либо задвоим эффект.
            async with conn.transaction():
                for r in rows:
                    if r["event_type"] in self.HANDLES:
                        await self.apply(conn, r)
                await conn.execute(
                    "UPDATE projection_checkpoints SET last_position=$1, updated_at=now() "
                    "WHERE projection_name=$2", rows[-1]["global_position"], self.NAME)

    async def apply(self, conn, r):
        p = json.loads(r["payload"])
        match r["event_type"]:
            case "CartCreated":
                await conn.execute(
                    "INSERT INTO cart_summary (cart_id, customer_id, items_count, total, status) "
                    "VALUES ($1,$2,0,0,'draft') ON CONFLICT (cart_id) DO NOTHING",
                    r["stream_id"], p["customer_id"])
            case "ItemAdded":
                await conn.execute(
                    "UPDATE cart_summary SET items_count = items_count + $2, "
                    "total = total + $3 WHERE cart_id = $1",
                    r["stream_id"], p["qty"], Decimal(p["price"]) * p["qty"])
            case "CartCheckedOut":
                await conn.execute(
                    "UPDATE cart_summary SET status='checked_out' WHERE cart_id=$1",
                    r["stream_id"])

Три свойства, без которых проекция сломается в проде:

  1. Идемпотентность. Доставка событий — at-least-once (см. «Событийную архитектуру»). Если чекпоинт и применение не в одной транзакции — обязателен либо ON CONFLICT DO NOTHING, либо хранение last_position прямо в строке read-модели с условием WHERE last_position < $new. Инкрементальные total = total + X — самая частая мина: при повторе они молча удваивают деньги.
  2. Возможность полной пересборки. Проекция обязана строиться с нуля из журнала, за одну команду, без ручных шагов. Это ваша страховка от любого бага в проекции: нашли ошибку — правите код, сбрасываете чекпоинт в 0, перестраиваете. Если пересборка невозможна, вы потеряли главное преимущество event sourcing.
  3. Мониторинг лага. Метрика max(global_position) − checkpoint по каждой проекции, алерт по порогу и по «чекпоинт не двигается N секунд». Растущий лаг = пользователи видят прошлое.

3.1. Blue/green пересборка проекций

Живой приём для пересборки без даунтайма: новая версия проекции пишется в новые таблицы (cart_summary_v4), догоняет журнал с нуля, и только когда лаг сошёлся к нулю — читатели переключаются переименованием/фиче-флагом. Старая версия остаётся ещё на день как откат. Именно поэтому версия входит в имя проекции и в имя таблицы.


4. Eventual consistency: единственная настоящая проблема пользователя

Пользователь нажал «добавить в корзину», страница перезагрузилась — и товара нет. Через 300 мс он появился, но доверие уже потеряно. Это не гипотетический сценарий, это режим отказа по умолчанию для асинхронного CQRS.

Окно рассинхронизации между моделью записи и моделью чтения

Плохое решение — sleep(500) на клиенте. Работающие решения, по возрастанию цены:

1. Клиент рисует результат сам. Команда выполнена успешно → значит, факт зафиксирован → UI применяет ожидаемое изменение локально, не дожидаясь проекции. Это то, что делают все современные фронтенд-фреймворки под названием optimistic UI. Самый дешёвый и самый частый вариант.

2. Read-your-own-writes через версию. Команда возвращает достигнутую версию потока (или глобальную позицию), клиент передаёт её в следующий запрос, а сторона чтения либо ждёт, пока проекция догонит, либо отвечает 409/Retry-After:

// Сторона записи возвращает позицию, до которой надо догнать
type CommandResult = { streamId: string; version: number; position: number };

// Сторона чтения умеет ждать — с жёстким потолком ожидания
async function getCart(id: string, minPosition?: number): Promise<CartView> {
  if (minPosition !== undefined) {
    const ok = await waitForProjection("cart_summary_v3", minPosition, { timeoutMs: 300 });
    // Честно сообщаем о рассинхронизации вместо тихой отдачи старых данных
    if (!ok) throw new StaleReadError("проекция отстаёт, повторите запрос");
  }
  return db.oneOrNone("SELECT * FROM cart_summary WHERE cart_id = $1", [id]);
}

3. Синхронная проекция для критичного экрана. Ничто не запрещает обновлять одну-две самые чувствительные проекции в той же транзакции, что и запись событий, а остальные — асинхронно. Гибрид совершенно легитимен и очень распространён.

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

Ключевой вопрос, который надо задавать бизнесу по каждому экрану: какое запаздывание здесь допустимо? Для витрины «товары дня» — минуты. Для остатка на счёте перед списанием — нисколько, и это признак того, что проверка должна жить на стороне записи, а не читаться из проекции. Проверять инвариант по read-модели — категорическая ошибка: read-модель по определению отстаёт, и вы получите двойное списание.


5. Версионирование событий — то, обо что разбиваются проекты

События живут вечно. Ваш код — нет. Через два года вам понадобится, чтобы в ItemAdded было поле warehouse_id, а в журнале лежит 40 миллионов событий без него. Мигрировать журнал UPDATE-ом нельзя: это разрушает саму гарантию неизменности (и ломает подписи/аудит, если они есть).

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

Слабая схема (weak schema). Сериализуйте в формат с необязательными полями (JSON, Protobuf, Avro со схемой) и никогда не делайте обязательным поле, которого не было. Добавление опционального поля — обратно совместимое изменение; удаление или смена типа — нет.

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

def upcast_item_added_v1_to_v2(p: dict) -> dict:
    # v2 добавила warehouse_id. Исторические события — со склада по умолчанию.
    return {**p, "warehouse_id": p.get("warehouse_id", "WH-MAIN")}

def upcast_item_added_v2_to_v3(p: dict) -> dict:
    # v3 разделила price на price_net + vat_rate.
    if "price_net" in p:
        return p
    net = Decimal(p["price"]) / Decimal("1.20")
    return {**p, "price_net": str(net.quantize(Decimal("0.01"))), "vat_rate": "0.20"}

UPCASTERS = {
    ("ItemAdded", 1): upcast_item_added_v1_to_v2,
    ("ItemAdded", 2): upcast_item_added_v2_to_v3,
}

def upcast(event_type: str, version: int, payload: dict):
    """Прогоняем по цепочке до актуальной версии схемы."""
    while (fn := UPCASTERS.get((event_type, version))) is not None:
        payload, version = fn(payload), version + 1
    return deserialize(event_type, payload)

Копирование-преобразование (copy-transform). Для несовместимых изменений, которые апкастером не выразить: новый журнал строится из старого прогоном трансформации, старый архивируется. Дорого, останавливает запись на время переключения — но иногда единственный выход. Исчерпывающий разбор всех стратегий — книга Грега Янга «Versioning in an Event Sourced System» (доступна бесплатно).

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


6. Продакшн-вопросы, которые всплывают на третий месяц

6.1. GDPR и «право на забвение» против неизменяемого журнала

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

  • Crypto-shredding (основной). Персональные данные в событиях шифруются ключом, отдельным для каждого субъекта; ключ лежит вне журнала. Запрос на удаление = уничтожение ключа. Журнал цел, структура и агрегаты не пострадали, данные необратимо нечитаемы. См. Crypto-shredding.
  • Вынос PII из событий. Событие несёт только customer_id; имя, адрес, телефон живут в отдельном изменяемом хранилище. Часто это проще всего и стоит рассматривать первым.
  • Точечное переписывание потока. Поток пересоздаётся без «плохих» событий. Работает, но ломает глобальные позиции и требует пересборки проекций — оставьте на крайний случай.

6.2. Kafka как event store — почему обычно нет

Соблазн понятен: Kafka уже стоит, она append-only и упорядочена. Но event store для агрегатов требует двух вещей, которых у Kafka нет из коробки:

  • чтение всего потока одного агрегата — Kafka читает партицию, а не ключ; чтобы свернуть один агрегат, придётся вычитать всю партицию или держать отдельный compacted-топик;
  • атомарная условная запись «добавь, если версия ровно N» — оптимистической блокировки на уровне ключа в Kafka нет.

Отсюда практический вывод: Kafka отлично работает как транспорт событий между сервисами и как вход для проекций, но плохо — как хранилище состояния агрегатов. Типичная зрелая конфигурация: PostgreSQL/EventStoreDB как event store внутри сервиса + публикация событий наружу в Kafka через outbox. Взгляд со стороны Confluent — Event sourcing with Apache Kafka; обратите внимание, сколько там оговорок.

6.3. Тестирование: given / when / then

Приятный побочный эффект чистых decide/evolve — тесты пишутся на языке бизнеса и не требуют ни базы, ни моков:

def test_нельзя_добавить_позицию_в_оформленную_корзину():
    # given: история фактов
    history = [
        CartCreated(cart_id=CID, customer_id=UID),
        ItemAdded("SKU-1", 2, Decimal("100")),
        CartCheckedOut(total=Decimal("200")),
    ]
    state = fold(history)

    # when / then: решение отвергается
    with pytest.raises(DomainError):
        decide_add_item(state, "SKU-2", 1, Decimal("50"))

def test_повторный_checkout_идемпотентен():
    state = fold([CartCreated(CID, UID), ItemAdded("SKU-1", 1, Decimal("10")),
                  CartCheckedOut(total=Decimal("10"))])
    assert decide_checkout(state) == []      # ни одного нового события

Такие тесты переживают любой рефакторинг внутренностей, потому что зафиксированы на контракте «факты → решение», а не на структуре классов.

6.4. Что ещё придётся построить

Список работ, которые в проекте с event sourcing никто не закладывает в оценку, а они занимают недели:

  • админ-инструмент просмотра потока событий (без него отладка невозможна);
  • перезапуск и пересборка проекций как самообслуживаемая операция;
  • дашборд лага всех проекций с алертами;
  • correlation_id / causation_id в метаданных каждого события — иначе трассировка причинно-следственных цепочек в распределённой системе безнадёжна;
  • политика архивации холодных потоков;
  • нагрузочный тест пересборки на объёме, ожидаемом через два года.

7. Когда НЕ надо

Честный разбор — половина ценности паттерна. Event sourcing и асинхронный CQRS не нужны, если:

  • домен по сути CRUD: справочники, настройки, контент. История изменений здесь либо не нужна, либо решается таблицей аудита за один день работы;
  • команда впервые видит эти паттерны и одновременно горит срок. Кривая обучения — месяцы, и ошибки моделирования событий фиксируются в журнале навсегда;
  • вам нужна только история — возьмите temporal tables в SQL Server, системное версионирование в MariaDB или простой аудит-лог;
  • вам нужны только быстрые отчёты — сделайте реплику для чтения или отдельную витрину, это ступень 3 CQRS без всякого event sourcing.

Наоборот, паттерн окупается, когда: история и причины изменений — часть предметной области (финансы, страхование, медицина, логистика, комплаенс); нужен аудит «кто, что, когда, почему»; требуется ретроспективный анализ «а что было бы, если» на реальных данных; или новые способы читать данные появляются постоянно, а исторические данные для них нужно построить задним числом (это, пожалуй, самое недооценённое преимущество: новая проекция строится по всей истории, которой в CRUD-системе просто нет).

Отдельно подчеркну ветку J: event sourcing — это решение уровня ограниченного контекста, а не системы. Совершенно нормальная зрелая архитектура: биллинг на event sourcing, каталог на обычном CRUD, поиск на Elasticsearch-проекции. Тотальный event sourcing «во всём приложении» — почти всегда признак того, что паттерн выбрали по восхищению, а не по задаче.


8. Как это выглядит в проде: инструменты

Инструмент Экосистема Чем интересен
EventStoreDB / Kurrent язык-агностична Специализированная БД под журналы: потоки, подписки, встроенные проекции, оптимистическая блокировка из коробки
Marten .NET + PostgreSQL Event store и document store поверх Postgres; сильная сторона — асинхронные демоны проекций и пересборка
Axon Framework JVM Полный стек: агрегаты, шина команд, саги, снимки, апкастеры
Akka Persistence JVM / Scala Event sourcing на акторах: persistent actor = агрегат в памяти, восстанавливается из журнала
Eventide Ruby Минималистичный event sourcing поверх Postgres
Собственная таблица в PostgreSQL любая Схема из раздела 2.5 — абсолютно рабочий и очень частый выбор для одного контекста; не бойтесь его

Отдельно стоит посмотреть открытые примеры: EventSourcing.NetCore Оскара Дудыча — вероятно, лучший публично доступный набор рабочих реализаций с разбором компромиссов.


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

  1. «CQRS = обязательно две базы и шина». Нет. Ступень 1 бесплатна и полезна почти везде.
  2. Валидация в evolve. Событие — свершившийся факт. Если evolve бросает исключение, вы не сможете восстановить агрегат из собственного журнала после изменения правил.
  3. События-снимки: OrderUpdated{...все поля...}. Потерян смысл, потеряна причина изменения, потеряна вся ценность. Событие должно отвечать на вопрос «что произошло в терминах бизнеса».
  4. CRUD-события. OrderRowInserted, StatusFieldChanged — это журнал репликации БД, а не доменные события. Имена берутся у бизнеса, а не у таблицы.
  5. Проверка инварианта по read-модели. Она отстаёт по определению. Инварианты — только на стороне записи, внутри одного агрегата, в одной транзакции.
  6. Неидемпотентная проекция. total = total + X без защиты от повтора — тихая порча данных, которую заметят через месяц.
  7. Чекпоинт и применение в разных транзакциях. Гарантированный источник дублей или потерь при рестарте.
  8. Снимок как источник истины. Если систему нельзя запустить после TRUNCATE snapshots, снимки перестали быть кешем.
  9. Один гигантский агрегат. «Агрегат Пользователь» с 300 000 событий — сигнал, что границы выбраны неверно.
  10. Отсутствие плана версионирования с первого дня. Поле event_version в схеме стоит ноль усилий сейчас и спасает проект через год.
  11. Публикация события до коммита транзакции. Классическая гонка: подписчик прочитал событие раньше, чем факт зафиксирован, или факт зафиксирован без публикации. Решается паттерном outbox — см. следующую статью.
  12. Event sourcing во всей системе разом. Начинайте с одного контекста, где ценность истории очевидна.

10. Мини-итог

  • CQRS — разделение модели записи (инварианты, нормализация, поведение) и модели чтения (денормализация под конкретный экран). Это лестница из четырёх ступеней, а не выключатель; начинать надо с самой дешёвой.
  • Event Sourcing — хранение состояния как неизменяемого журнала фактов; состояние есть fold(evolve, S₀, events). Снимки — кеш свёртки, не истина.
  • Ядро домена — две чистые функции: decide(state, command) → [events] и evolve(state, event) → state. Это даёт тестируемость без моков и полную обратимость решений.
  • Конкурентность решается оптимистической блокировкой по (stream_id, version), а ретрай обязан перечитывать состояние и заново принимать решение.
  • Проекции обязаны быть идемпотентными, полностью пересобираемыми и под мониторингом лага.
  • Eventual consistency — не техническая деталь, а контракт с пользователем; закрывается optimistic UI, read-your-own-writes по версии или синхронной проекцией для критичных экранов.
  • Версионирование событий (слабая схема + апкастеры) закладывается в первый день, иначе через год оно закладывает проект.
  • Оба паттерна — решение уровня одного ограниченного контекста. Тотальное применение почти всегда ошибка.

Источники


Что дальше

Мы аккуратно обошли стороной главный вопрос: события зафиксированы в журнале, но как гарантировать, что они дойдут до других сервисов ровно один раз по эффекту, и что делать, когда бизнес-операция охватывает несколько агрегатов и несколько сервисов сразу? Об этом — следующая статья:

Saga, распределённые транзакции, outbox и идемпотентность

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

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

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

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