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 — про целые модели и их хранилища.
PlaceOrder"] VAL["Валидация + инварианты"] AGG["Агрегат
нормализованный, с поведением"] WDB[("Хранилище записи")] end subgraph read["Сторона чтения — модели ЧТЕНИЯ"] Q["Запрос
GetOrdersPage"] RM1[("Плоская таблица
для списка")] RM2[("Поисковый индекс")] RM3[("Витрина аналитики")] end UI -->|"изменить"| CMD --> VAL --> AGG --> WDB UI -->|"прочитать"| Q Q --> RM1 Q --> RM2 Q --> RM3 WDB -.->|"проекции:
синхронно или асинхронно"| RM1 WDB -.-> RM2 WDB -.-> RM3
Ключевое в этой схеме — пунктирная стрелка. Именно она, а не сам факт разделения, определяет, во что вам обойдётся 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ₙ])
Отсюда сразу три следствия, которые надо запомнить:
evolveне имеет права отказать. Событие — это уже случившийся факт, прошедшее время.ItemAddedнельзя «не принять»: он уже в журнале. Вся валидация происходит вdecide, до записи.decideне имеет права ничего менять. Она возвращает список событий; их применение и запись — снаружи. Это делает домен тестируемым без единого мока.- Именование событий в прошедшем времени — не стилистика, а семантика.
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("недостижимо")
Разбор поведения при конфликте — важнее, чем кажется:
и заново принимаем решение H->>ES: load(cart-7) ES-->>H: state, version = 5 H->>ES: append(cart-7, expected=5, [ItemAdded SKU-2]) ES-->>H: OK, version = 6 ES-)P: события 5, 6 (по позиции в журнале) P->>RM: upsert строки корзины (идемпотентно) P->>ES: checkpoint = позиция события 6
Почему ретрай обязан перечитывать состояние, а не просто повторять 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"])
Три свойства, без которых проекция сломается в проде:
- Идемпотентность. Доставка событий — at-least-once (см. «Событийную архитектуру»). Если чекпоинт и применение не в одной транзакции — обязателен либо
ON CONFLICT DO NOTHING, либо хранениеlast_positionпрямо в строке read-модели с условиемWHERE last_position < $new. Инкрементальныеtotal = total + X— самая частая мина: при повторе они молча удваивают деньги. - Возможность полной пересборки. Проекция обязана строиться с нуля из журнала, за одну команду, без ручных шагов. Это ваша страховка от любого бага в проекции: нашли ошибку — правите код, сбрасываете чекпоинт в 0, перестраиваете. Если пересборка невозможна, вы потеряли главное преимущество event sourcing.
- Мониторинг лага. Метрика
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-системе просто нет).
как часть домена?"] -->|нет| B["Нужны быстрые
разнородные чтения?"] A -->|да| C["Достаточно ли
аудит-лога / temporal tables?"] B -->|нет| D["Обычный CRUD.
Ступень 1 CQRS — и всё"] B -->|да| E["CQRS ступень 2–3
БЕЗ event sourcing"] C -->|да| F["Аудит-таблица или
системное версионирование БД"] C -->|нет| G["Нужно ли принимать
решения по истории?
(лимиты, начисления, споры)"] G -->|нет| F G -->|да| H["Есть ли команда,
готовая к 3–6 мес. кривой обучения
и к версионированию событий?"] H -->|нет| I["Отложить.
Начать с ступени 2 CQRS
и богатого аудит-лога"] H -->|да| J["Event Sourcing
в ОДНОМ ограниченном контексте,
не во всей системе"] style D fill:#5aa46933,stroke:#5aa469 style E fill:#4a90d933,stroke:#4a90d9 style F fill:#5aa46933,stroke:#5aa469 style I fill:#c98a2b33,stroke:#c98a2b style J fill:#9b59b633,stroke:#9b59b6
Отдельно подчеркну ветку 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. Типичные ошибки
- «CQRS = обязательно две базы и шина». Нет. Ступень 1 бесплатна и полезна почти везде.
- Валидация в
evolve. Событие — свершившийся факт. Еслиevolveбросает исключение, вы не сможете восстановить агрегат из собственного журнала после изменения правил. - События-снимки:
OrderUpdated{...все поля...}. Потерян смысл, потеряна причина изменения, потеряна вся ценность. Событие должно отвечать на вопрос «что произошло в терминах бизнеса». - CRUD-события.
OrderRowInserted,StatusFieldChanged— это журнал репликации БД, а не доменные события. Имена берутся у бизнеса, а не у таблицы. - Проверка инварианта по read-модели. Она отстаёт по определению. Инварианты — только на стороне записи, внутри одного агрегата, в одной транзакции.
- Неидемпотентная проекция.
total = total + Xбез защиты от повтора — тихая порча данных, которую заметят через месяц. - Чекпоинт и применение в разных транзакциях. Гарантированный источник дублей или потерь при рестарте.
- Снимок как источник истины. Если систему нельзя запустить после
TRUNCATE snapshots, снимки перестали быть кешем. - Один гигантский агрегат. «Агрегат Пользователь» с 300 000 событий — сигнал, что границы выбраны неверно.
- Отсутствие плана версионирования с первого дня. Поле
event_versionв схеме стоит ноль усилий сейчас и спасает проект через год. - Публикация события до коммита транзакции. Классическая гонка: подписчик прочитал событие раньше, чем факт зафиксирован, или факт зафиксирован без публикации. Решается паттерном outbox — см. следующую статью.
- 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 по версии или синхронной проекцией для критичных экранов.
- Версионирование событий (слабая схема + апкастеры) закладывается в первый день, иначе через год оно закладывает проект.
- Оба паттерна — решение уровня одного ограниченного контекста. Тотальное применение почти всегда ошибка.
Источники
- Greg Young. CQRS Documents — первоисточник.
- Greg Young. Versioning in an Event Sourced System — единственная книга целиком про эволюцию схемы событий.
- Martin Fowler. CQRS и Event Sourcing.
- Udi Dahan. Clarified CQRS — про то, чего CQRS НЕ означает.
- Microsoft. CQRS pattern, Event Sourcing pattern.
- Pat Helland. Immutability Changes Everything, ACM Queue — почему append-only меняет всю системную архитектуру.
- Chris Richardson. Event sourcing pattern, microservices.io.
- Oskar Dudycz. event-driven.io — практический блог с разбором граничных случаев.
- Vaughn Vernon. Implementing Domain-Driven Design — главы 8 (Domain Events) и 4 (Architecture).
Что дальше
Мы аккуратно обошли стороной главный вопрос: события зафиксированы в журнале, но как гарантировать, что они дойдут до других сервисов ровно один раз по эффекту, и что делать, когда бизнес-операция охватывает несколько агрегатов и несколько сервисов сразу? Об этом — следующая статья: