Event Sourcing красиво звучит — хранить не текущее состояние, а поток изменивших его событий — и концептуально мы это уже разобрали. Но между «мне нравится идея» и работающей системой лежит слой инженерных решений, о который спотыкаются: где физически хранить события, как из них собирать актуальное состояние и что делать, когда проекций стало десять и одну нужно перестроить, как не читать миллион событий на каждый запрос, и как удалить персональные данные из журнала, который по определению неизменяем.
Эта статья — про эксплуатацию Event Sourcing: не «что это», а «как это держать в проде». Практический аддендум серии «Событийные паттерны» к концептуальной статье; read-сторону (проекции) продолжает CQRS на практике.
В статье
- Где физически хранить события
- Запись: оптимистичный конкаренси-контроль
- Проекции: от событий к read-моделям
- Перестройка проекций без даунтайма
- Снапшоты: не читать всю историю агрегата
- Эволюция событий: upcasting
- Лунная база: журнал расхода кислорода на практике
- Append-only и право на забвение
- Эксплуатационные ошибки
- Checklist: эксплуатационная готовность
Где физически хранить события
Event store должен уметь три вещи: append (дописать событие в конец стрима), чтение по агрегату (прочитать все события одного стрима по порядку) и подписку на поток (чтобы проекции узнавали о новых событиях). Плюс оптимистичный конкаренси-контроль — иначе два параллельных изменения агрегата затрут друг друга. От того, какие из этих требований критичны, зависит выбор хранилища.
| Хранилище | Append | Чтение по агрегату | Подписка | Оптимистичная блокировка | Когда брать |
|---|---|---|---|---|---|
| PostgreSQL (таблица событий) | INSERT | WHERE stream_id=… |
LISTEN/NOTIFY или CDC | уникальный индекс (stream_id, version) |
дефолт для большинства; уже есть в стеке |
| EventStoreDB | нативно | нативно (стримы — первичная абстракция) | встроенные подписки/проекции | нативно (expected version) |
когда ES — ядро системы, нужны из коробки catch-up-подписки |
| Kafka | нативно | ⚠️ неудобно (partition ≠ stream) | нативно | ⚠️ нет на уровне стрима | как шина событий поверх store, не как сам store |
Kafka здесь — частая ловушка: лог выглядит как event store, но partition ≠ stream, точечно прочитать историю одного агрегата дорого, а оптимистичной блокировки на уровне стрима нет. Kafka хороша как транспорт событий наружу (в проекции, в другие сервисы), но источником правды агрегата её делают редко и осознанно.
Минимальная схема на Postgres:
CREATE TABLE events (
stream_id uuid NOT NULL, -- идентификатор агрегата
version bigint NOT NULL, -- порядковый номер в стриме, с 1
event_type text NOT NULL, -- 'OxygenConsumed'
schema_ver int NOT NULL DEFAULT 1,-- версия формата события (для upcasting)
data jsonb NOT NULL, -- полезная нагрузка
metadata jsonb NOT NULL, -- correlation_id, actor, causation_id
created_at timestamptz NOT NULL DEFAULT now(),
global_pos bigserial, -- сквозной порядок для проекций
PRIMARY KEY (stream_id, version) -- инвариант: версия уникальна в стриме
);
CREATE INDEX ON events (global_pos); -- проекции читают по глобальной позиции
global_pos (сквозной монотонный номер) — то, за что цепляются проекции: они хранят «дочитал до позиции N» и продолжают с неё.
Запись: оптимистичный конкаренси-контроль
Два запроса одновременно меняют один агрегат — например, две подсистемы списывают кислород. Без защиты оба прочитают версию 42, оба запишут версию 43, один инвариант потеряется. Решение — оптимистичная блокировка: команда несёт ожидаемую версию, а PRIMARY KEY (stream_id, version) физически не даст записать второе событие с той же версией.
загрузить агрегат → текущая версия = 42
проверить бизнес-правило
INSERT event (stream_id, version=43, …) -- ожидаем, что 43 ещё нет
├─ успех → команда применена
└─ конфликт (нарушение PK) → кто-то опередил:
перечитать агрегат и повторить (или вернуть ошибку клиенту)
Это то же оптимистичное управление, что и в обычных БД (см. аномалии и изоляцию), только версия живёт в стриме событий. Пессимистичные блокировки в ES почти не применяют — они убивают пропускную способность append-only-модели.
Проекции: от событий к read-моделям
Раз текущее состояние не хранится, его вычисляют — это проекция (read-модель). Два режима, и выбор между ними определяет консистентность:
- Синхронная (in-memory replay): на запрос читаем события стрима и сворачиваем в состояние. Строго консистентно, но дорого на длинных стримах — отсюда снапшоты (ниже).
- Асинхронная (materialized view): отдельный проектор читает поток по
global_posи обновляет отдельную таблицу/индекс/кэш. Быстрые чтения ценой eventual consistency — сразу после записи проекция ещё не догнала.
Проекций может быть несколько — по одной на каждый сценарий чтения (список, деталь, поиск, аналитика), каждая в своём хранилище. Это и есть стык с CQRS: write-модель (стрим событий) и read-модели (проекции) разделены. Подробно read-сторону — отставание, read-your-writes, несколько моделей — разбирает CQRS на практике.
Каркас асинхронного проектора держится на двух свойствах — чекпоинт позиции и идемпотентность применения:
loop:
checkpoint = load_checkpoint("balance_projection") -- напр. global_pos = 10456
for e in read_events(after=checkpoint, limit=500):
apply(e) -- UPSERT в read-таблицу, идемпотентно
checkpoint = e.global_pos
save_checkpoint("balance_projection", checkpoint) -- атомарно с apply
Ключевое: apply должен быть идемпотентным (событие может прийти повторно после падения между apply и save_checkpoint), а чекпоинт — сохраняться в той же транзакции, что и изменение read-модели. Иначе после сбоя проекция либо задвоит эффект, либо потеряет событие.
Перестройка проекций без даунтайма
Проекции перестраивают чаще, чем кажется: поменялась схема read-модели, нашли баг в проекторе, добавили новую проекцию задним числом. Поскольку события — источник правды, проекцию всегда можно собрать заново, проиграв историю с нуля. Вопрос — как это сделать, не останавливая чтения.
Рабочий приём — blue-green проекции:
- Создаём новую read-таблицу
balance_v2рядом с живойbalance_v1. - Запускаем проектор
v2сglobal_pos = 0— он проигрывает всю историю вbalance_v2(это может занять время). - Когда
v2догнала хвост (её чекпоинт ≈ текущий максимумglobal_pos), переключаем чтение наbalance_v2. - Старую
balance_v1удаляем.
Требования к этому — те же идемпотентность и чекпоинты: реплей должен давать тот же результат, что и онлайн-обработка. Именно поэтому события пишут неизменяемыми и самодостаточными: проектор через год должен понять старое событие без обращения к «текущему» коду. Здесь же всплывает upcasting (ниже) — реплей проходит через события всех исторических версий формата.
Снапшоты: не читать всю историю агрегата
Если агрегат накопил десятки тысяч событий, синхронный replay на каждую загрузку становится дорогим. Снапшот — сохранённое состояние агрегата на конкретной версии. Загрузка тогда: взять последний снапшот + проиграть только события после него.
snapshot(stream_id) = { version: 8000, state: {...}, created_at }
загрузка агрегата:
s = load_latest_snapshot(stream_id) -- version 8000
events = read_events(stream_id, after_version=8000) -- только хвост
state = fold(s.state, events)
Практические решения:
- Частота. По числу событий (каждые N) или по времени; снапшот делают асинхронно, он никогда не источник правды — только кэш свёртки.
- Версионирование снапшота. Снапшот хранит состояние в текущем формате кода. Поменялась логика свёртки — старые снапшоты надо инвалидировать и пересобрать из событий (события переживут смену кода, снапшот — нет).
- Не всем нужны. Если стримы короткие (сотни событий) — снапшоты только добавляют движущихся частей. Это оптимизация под конкретную боль, а не обязательный элемент.
Эволюция событий: upcasting
События неизменяемы, но код вокруг них живёт годами, и формат меняется: добавилось поле, переименовалось, разбилось надвое. Читать старые события новый код должен уметь — иначе история становится нечитаемой. Приём называется upcasting: при чтении событие старой версии на лету преобразуется в актуальную форму.
read_event(raw):
e = deserialize(raw)
while e.schema_ver < CURRENT_VERSION:
e = upcast[e.event_type][e.schema_ver](e) -- v1→v2→v3
return e
Отсюда поле schema_ver в таблице. Upcasting — это версионирование контрактов, применённое внутрь стрима; полностью тема раскрыта в контрактах событий и schema evolution. Главное правило — никогда не переписывать старые события в БД задним числом: преобразование живёт в коде чтения, а не в UPDATE по таблице (события неизменяемы, и этот инвариант ломать нельзя).
Лунная база: журнал расхода кислорода на практике
Вернёмся к учёту кислорода из концептуальной статьи, но теперь на уровне таблиц. Стрим одного резервуара — это строки в events:
| stream_id | version | event_type | data | global_pos |
|---|---|---|---|---|
| tank-A | 1 | OxygenDelivered |
{amount: 1000, source: "cargo-7"} |
5001 |
| tank-A | 2 | OxygenConsumed |
{amount: 120, module: "habitat-A"} |
5002 |
| tank-A | 3 | OxygenReserveAllocated |
{amount: 33, reason: "emergency"} |
5004 |
Из одного стрима собираются разные проекции под разные вопросы:
balance(текущий остаток по резервуарам) — свёртка+delivered −consumed −reserved; читает диспетчерская;consumption_by_module(расход по модулям во времени) — для аналитики и прогноза исчерпания;reserve_audit(когда и почему трогали аварийный резерв) — для комплаенса.
Теперь реальный эксплуатационный сценарий. Аналитики просят новую проекцию — прогноз исчерпания по скорости расхода. Событий за полгода — миллионы. Действуем по blue-green: запускаем новый проектор с global_pos = 0, он проигрывает всю историю в forecast, мы смотрим, как растёт его чекпоинт, и когда он догнал хвост — подключаем к API. Живые проекции при этом не останавливались, запись новых событий шла как обычно. Через месяц у balance замечаем, что синхронная загрузка резервуара с длинной историей тормозит диспетчерскую — добавляем снапшоты каждые 5000 событий, и загрузка снова читает лишь хвост.
Именно эта способность — задним числом собрать проекцию, которую при проектировании никто не предусмотрел, — и есть практическая выгода Event Sourcing, ради которой терпят его сложность.
Append-only и право на забвение
Первый вопрос комплаенса к ES: события неизменяемы, а GDPR/152-ФЗ требуют удалить персональные данные. «Просто вырезать» событие нельзя — сломаются проекции и инвариант версий. Рабочие подходы (подробнее — в крипте для разработчикаСкоро):
- crypto-shredding — персональные поля события шифруются индивидуальным ключом субъекта, ключ хранится в отдельном мутабельном хранилище; «удаление» = уничтожение ключа. Событие остаётся, версии целы, проекции работают, но данные превратились в нечитаемый шифротекст.
- вынос PII из события — в стриме лежит только ссылка (ID), сами персональные данные — в отдельной таблице, которую можно физически чистить. Event store при этом остаётся полностью append-only.
- компенсирующее событие с маскированием — годится для отображения («данные скрыты»), но не закрывает требование физического удаления, поэтому как единственная мера не подходит.
Стратегию забвения закладывают до первого события, а не после запроса регулятора: задним числом расшифровать, какие поля были персональными в миллионах старых событий, — отдельная дорогая история.
Эксплуатационные ошибки
Концептуальные ошибки ES (слишком крупные события, мутабельность, «ES везде») разобраны в базовой статье. Здесь — то, что ломается именно в эксплуатации:
- Неидемпотентное применение в проекции. Проектор падает между
applyи сохранением чекпоинта, при рестарте применяет событие повторно — баланс задваивается. Лечится идемпотентным UPSERT и атомарностью «apply + checkpoint». - Чекпоинт не в одной транзакции с read-моделью. Сохранили чекпоинт, но не успели применить событие (или наоборот) — проекция тихо разъезжается с историей. Оба изменения — в одной транзакции.
- Kafka как единственный event store. Захотели прочитать историю одного агрегата или получить оптимистичную блокировку — а
partition ≠ stream, и обоего нет. Больно осознавать это после запуска. - Снапшот пережил смену логики свёртки. Поменяли
fold, старые снапшоты содержат состояние по старой логике — новые загрузки берут неверную базу. Снапшоты инвалидируются вместе с изменением кода свёртки. - Проекция, которая вечно догоняет. Один проектор на весь поток не успевает за темпом записи. Дробят по типам агрегатов/партициям — но тогда следят за порядком внутри стрима (события одного агрегата — строго по версии).
- Реплей ломается на старых событиях. Новый код не умеет читать формат двухлетней давности, потому что upcasting не заложили. Перестройка проекции внезапно невозможна — а это главная страховка ES.
Checklist: эксплуатационная готовность
- Выбрано хранилище под реальные требования (нужны ли catch-up-подписки и оптимистичная блокировка из коробки, или хватит Postgres)?
- Есть уникальный инвариант
(stream_id, version)и обработка конфликта версий на записи? - Все проекторы идемпотентны, а чекпоинт сохраняется атомарно с read-моделью?
- Отработана перестройка проекции без даунтайма (blue-green) — хотя бы на одной?
- Заложено ли поле версии формата события и механизм upcasting?
- Решено, нужны ли снапшоты (длина стримов), и как они инвалидируются при смене свёртки?
- Если в событиях есть PII — выбрана стратегия забвения (crypto-shredding / вынос PII) до первого события?
- Есть наблюдаемость лага проекций (насколько чекпоинт отстаёт от хвоста)?
Если на большинство — «да», Event Sourcing у вас не только красив на слайде, но и эксплуатируем.
Демо и версии
- Demo:
digital-cookbook/architecture/event-sourcing— Postgres event store (схемаeventsсglobal_posи оптимистичной блокировкой) + асинхронный проектор баланса с чекпоинтом + снапшоты + перестройка проекции blue-green реплеем; Go и Java. Сценарий кислородного резервуара: запись событий, конфликт версий при параллельной записи, сборка новой проекции задним числом. - Версии Postgres/EventStoreDB/библиотек пиновать при написании (поведение подписок и проекций чувствительно к версии).
Документация и первоисточники
- Канон: Martin Fowler — Event Sourcing, Greg Young — CQRS Documents (PDF), Microsoft — Event Sourcing pattern.
- Продукт/event stores: EventStoreDB, Marten (Postgres, .NET), Axon Server.
- Смежное на сайте: Event Sourcing — концепция, CQRS, CQRS на практике, контракты событий (upcasting), CDC (Debezium), аномалии и изоляция, крипта для разработчика (crypto-shredding/GDPR)Скоро, матрица решений.
Комментарии