Event Sourcing на практике: event store, проекции, снапшоты, replay

Концепцию Event Sourcing легко полюбить и трудно эксплуатировать: где хранить события (EventStoreDB vs Kafka vs Postgres), как строить и перестраивать проекции (read models), зачем снапшоты и как их делать, как реплеить историю, что такое upcasting при эволюции событий, оптимистичный конкаренси-контроль через expected version — и как жить с append-only, когда приходит GDPR-запрос на удаление

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

Эта статья — про эксплуатацию Event Sourcing: не «что это», а «как это держать в проде». Практический аддендум серии «Событийные паттерны» к концептуальной статье; read-сторону (проекции) продолжает CQRS на практике.

Архивный барабан event store выдаёт ленту штампованных событий-карточек в проекционные экраны read-моделей; закладка SNAPSHOT и петля REPLAY назад, приборные панели рядом

В статье

Где физически хранить события

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 проекции:

  1. Создаём новую read-таблицу balance_v2 рядом с живой balance_v1.
  2. Запускаем проектор v2 с global_pos = 0 — он проигрывает всю историю в balance_v2 (это может занять время).
  3. Когда v2 догнала хвост (её чекпоинт ≈ текущий максимум global_pos), переключаем чтение на balance_v2.
  4. Старую 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 везде») разобраны в базовой статье. Здесь — то, что ломается именно в эксплуатации:

  1. Неидемпотентное применение в проекции. Проектор падает между apply и сохранением чекпоинта, при рестарте применяет событие повторно — баланс задваивается. Лечится идемпотентным UPSERT и атомарностью «apply + checkpoint».
  2. Чекпоинт не в одной транзакции с read-моделью. Сохранили чекпоинт, но не успели применить событие (или наоборот) — проекция тихо разъезжается с историей. Оба изменения — в одной транзакции.
  3. Kafka как единственный event store. Захотели прочитать историю одного агрегата или получить оптимистичную блокировку — а partition ≠ stream, и обоего нет. Больно осознавать это после запуска.
  4. Снапшот пережил смену логики свёртки. Поменяли fold, старые снапшоты содержат состояние по старой логике — новые загрузки берут неверную базу. Снапшоты инвалидируются вместе с изменением кода свёртки.
  5. Проекция, которая вечно догоняет. Один проектор на весь поток не успевает за темпом записи. Дробят по типам агрегатов/партициям — но тогда следят за порядком внутри стрима (события одного агрегата — строго по версии).
  6. Реплей ломается на старых событиях. Новый код не умеет читать формат двухлетней давности, потому что upcasting не заложили. Перестройка проекции внезапно невозможна — а это главная страховка ES.

Checklist: эксплуатационная готовность

  1. Выбрано хранилище под реальные требования (нужны ли catch-up-подписки и оптимистичная блокировка из коробки, или хватит Postgres)?
  2. Есть уникальный инвариант (stream_id, version) и обработка конфликта версий на записи?
  3. Все проекторы идемпотентны, а чекпоинт сохраняется атомарно с read-моделью?
  4. Отработана перестройка проекции без даунтайма (blue-green) — хотя бы на одной?
  5. Заложено ли поле версии формата события и механизм upcasting?
  6. Решено, нужны ли снапшоты (длина стримов), и как они инвалидируются при смене свёртки?
  7. Если в событиях есть PII — выбрана стратегия забвения (crypto-shredding / вынос PII) до первого события?
  8. Есть наблюдаемость лага проекций (насколько чекпоинт отстаёт от хвоста)?

Если на большинство — «да», Event Sourcing у вас не только красив на слайде, но и эксплуатируем.

Демо и версии

  • Demo: digital-cookbook/architecture/event-sourcing — Postgres event store (схема events с global_pos и оптимистичной блокировкой) + асинхронный проектор баланса с чекпоинтом + снапшоты + перестройка проекции blue-green реплеем; Go и Java. Сценарий кислородного резервуара: запись событий, конфликт версий при параллельной записи, сборка новой проекции задним числом.
  • Версии Postgres/EventStoreDB/библиотек пиновать при написании (поведение подписок и проекций чувствительно к версии).

Документация и первоисточники

Обсуждение в Telegram

Присоединиться →

Комментарии