Event Sourcing: хранить факты, а не состояние
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Event Sourcing: хранить факты, а не состояние
Третья ката второго блока. Урок 28 доказал закон: состояние это свёртка событий. Сегодня делаем из этого вывод. Если состояние всегда восстановимо из истории, зачем хранить состояние? Будем хранить только факты, а состояние собирать по требованию.
Сцена · что теряет snapshot-репозиторий
Вспомни SQLite-репозиторий из урока 25. SQLite-репозиторий. Он хранил текущее состояние брони: одна строка, при каждом переходе перезаписывалась. Бронь заселилась, строка Reserved стала CheckedIn, прошлое исчезло.
Что мы потеряли:
- Историю. Когда забронировали? Когда заселили? Сколько раз меняли? Строка этого не помнит, она знает только “сейчас”.
- Аудит. “Кто и когда отменил эту бронь?” Ответа нет, факт отмены затёр предыдущий.
- Отладку во времени. Баг проявился после check-out. С snapshot мы видим только финал. С историей можно прокрутить плёнку и найти кадр, где всё сломалось.
Event Sourcing хранит вместо “сейчас” всю последовательность фактов. Состояние не лежит на диске, оно вычисляется: list.fold(events, initial_state, evolve). Закон Decider из урока 28 гарантирует, что этого достаточно.
Карта урока · что заберёшь через 3 часа
- Поймёшь, что event store это append-only лог фактов, а состояние это его свёртка.
- Соберёшь in-memory event store на OTP-actor с optimistic concurrency.
- Свяжешь Decider и store в command handler: load, replay, decide, append.
- Поймёшь optimistic concurrency через
expected_versionи зачем он нужен. - Узнаешь, что такое снапшот и когда он оправдан.
Концепт · event store как append-only лог
Event store это не таблица записей, а журнал. Две операции:
- load(stream) возвращает все события одного агрегата по порядку плюс версию (число событий). Стрим это история одной брони.
- append(stream, expected_version, events) дописывает события в конец, но только если версия совпадает с ожидаемой. Это и есть optimistic concurrency.
События никогда не меняются и не удаляются. Только дописываются. Это делает лог надёжным источником истины: прошлое неизменно по построению.
Контракт порта (та же запись с функциями, что в уроке 22):
// examples/ddd-hotel/gleam/src/event_store/types.gleam
pub type StreamSlice {
StreamSlice(version: Int, events: List(Event))
}
pub type ConcurrencyConflict {
ConcurrencyConflict(stream_id: String, expected: Int, actual: Int)
}
pub type EventStore {
EventStore(
load: fn(ReservationId) -> StreamSlice,
append: fn(ReservationId, Int, List(Event)) ->
Result(Nil, ConcurrencyConflict),
read_all: fn() -> List(StreamItem),
)
}
read_all отдаёт глобальный лог всех стримов по порядку. Он не нужен для команд, но нужен проекциям урока 30: read-model собирается из всех событий, а не одного стрима.
Концепт · optimistic concurrency через expected_version
Два запроса одновременно меняют одну бронь. Оба прочитали стрим версии 1, оба решили дописать событие. Без защиты второй затрёт решение первого, принятое на устаревшем состоянии.
Optimistic concurrency решает это без блокировок. Писатель при append передаёт версию, которую видел при load. Store сравнивает: если реальная версия другая, значит, кто-то записал между нашим чтением и записью, и append возвращает ConcurrencyConflict. Команду тогда перепроигрывают на свежем стриме.
Это “оптимистично”, потому что мы не блокируем заранее в надежде, что конфликта не будет, а ловим его по факту. На редких конфликтах это дешевле блокировок.
Языковые механики Gleam · actor с глобальным логом
In-memory store держит один список List(StreamItem) (новейшие в голове, для дешёвого prepend). Срез стрима это фильтр по id, версия это длина среза. Actor сериализует доступ, гонок нет (та же идея, что в уроке 22):
// examples/ddd-hotel/gleam/src/event_store/in_memory.gleam (фрагмент)
Append(reply, id, expected_version, events) -> {
let current = list.length(stream_events(log, id))
case current == expected_version {
False -> {
process.send(reply, Error(ConcurrencyConflict(
stream_id: reservation.reservation_id_to_string(id),
expected: expected_version,
actual: current,
)))
actor.continue(log)
}
True -> {
let new_items =
events
|> list.map(fn(event) { StreamItem(stream_id: id, event: event) })
|> list.reverse
process.send(reply, Ok(Nil))
actor.continue(list.append(new_items, log))
}
}
}
Проверка версии и запись происходят в одном сообщении actor-а, атомарно. Между “проверил версию” и “дописал” никто не вклинится, потому что actor обрабатывает сообщения по одному.
Задача · command handler
Свяжи Decider (урок 28) и store в одну функцию event_sourcing.handle_command:
pub type HandleError {
Rejected(DomainError)
Conflict(ConcurrencyConflict)
}
pub fn handle_command(
store: EventStore,
decider: ReservationDecider,
stream_id: ReservationId,
command: Command,
) -> Result(List(Event), HandleError)
Контракт, четыре шага:
- load срез стрима (события + версия).
- replay свернуть события в состояние через
evolve. - decide прогнать команду: либо новые события, либо
Rejected(domain_error). - append дописать с
expected_version = slice.version: успех илиConflict.
Две причины отказа разнесены в типе: Rejected это “домен сказал нет” (вина команды), Conflict это “кто-то опередил” (повод перепроиграть). Это разные исходы, и вызывающий обрабатывает их по-разному.
Решай прямо здесь: тесты прогонятся в песочнице, а подсказки и разбор ниже открывай, только если застрял.
Подсказки
expected_versionэто ровноslice.version, прочитанная на шаге 1. Мы говорим store: “дописываю, рассчитывая, что с момента моего чтения ничего не изменилось”.- На
Rejectedв store писать нечего: домен отверг команду, событий нет. Проверь тестом, что версия стрима осталась прежней. replayэтоlist.fold(slice.events, decider.initial_state, decider.evolve). Можно не звать отдельную функцию, форма короткая.- Имя поля стрима в типах называется
stream_id, неid. Один стрим это одна бронь, но завтра это может быть любой агрегат, и имя это подчёркивает.
Разбор · handle_command
// examples/ddd-hotel/gleam/src/event_sourcing.gleam
pub fn handle_command(
store: EventStore,
decider: ReservationDecider,
stream_id: ReservationId,
command: Command,
) -> Result(List(Event), HandleError) {
let slice = store.load(stream_id)
let state = list.fold(slice.events, decider.initial_state, decider.evolve)
case decider.decide(state, command) {
Error(domain_error) -> Error(Rejected(domain_error))
Ok(new_events) ->
case store.append(stream_id, slice.version, new_events) {
Ok(Nil) -> Ok(new_events)
Error(conflict) -> Error(Conflict(conflict))
}
}
}
Это вся write-сторона event sourcing в одной функции. Никакого мутабельного агрегата на диске: состояние собирается из истории на шаге 2, новые факты дописываются на шаге 4. Между ними чистый decide. Сравни с обычным репозиторием урока 22: там было load -> мутировать -> save, тут load -> replay -> decide -> append, и replay это то, что заменило “мутировать”.
Тест показывает полный цикл:
// examples/ddd-hotel/gleam/test/event_sourcing_test.gleam
pub fn place_then_check_in_builds_state_test() {
let assert Ok(#(store, _stop)) = in_memory.start()
let d = decider.reservation_decider(now)
let #(id, place) = place_command()
let assert Ok(_) = event_sourcing.handle_command(store, d, id, place)
let assert Ok(_) =
event_sourcing.handle_command(store, d, id, CheckInGuest(id))
event_sourcing.current_state(store, d, id)
|> reservation.tag_of
|> should.equal("CheckedIn")
}
current_state нигде не хранится: это load плюс replay. Состояние всегда вычисляется из фактов.
Подвигай ползунок и увидь, как состояние пересобирается свёрткой префикса стрима:
Концепт · снапшоты, когда история длинная
Замена состояния историей имеет цену: чтобы узнать “сейчас”, надо свернуть все события. У брони их единицы, fold мгновенен. Но у долгоживущего агрегата (счёт за год, игровая сессия) их тысячи, и replay с нуля на каждую команду становится дорогим.
Снапшот это кэш свёртки на версии N. Replay стартует не с initial_state, а с сохранённого снапшота, и доигрывает только события после версии N. Ключевое: снапшот это оптимизация, а не источник истины. Его можно удалить и пересобрать из событий в любой момент. Если снапшот разошёлся с историей, прав всегда лог фактов.
Мы не пишем снапшоты в коде урока (у брони они не нужны), но ДЗ предлагает их добавить: это ровно replay, стартующий не с нуля.
Критика · что не покрыто
- Store только in-memory. SQLite-версия это таблица
events(stream_id, version, payload)с уникальным индексом по(stream_id, version), который и обеспечивает optimistic concurrency на уровне БД (ДЗ). - Нет retry на конфликт.
handle_commandвозвращаетConflict, но не повторяет. Обёртка с retry (перечитать, перепроиграть, попробовать снова) это ДЗ. - Снапшоты в прозе. Реализованы как упражнение, не в основном коде. Для брони это переинженеринг.
- Один стрим на агрегат. Это классический ES. Альтернатива (один общий стрим, факт-функции без агрегатных границ) это aggregateless ES урока 32. Aggregateless ES.
- read_all держит всё в памяти. Для проекций урока 30 этого хватит, для прода нужна подписка на стрим с курсором, а не “прочитать весь лог”.
Takeaway
Одна фраза:
Event sourcing хранит факты, а не состояние: store это append-only лог, состояние это его свёртка через
evolve, а optimistic concurrency черезexpected_versionловит гонки без блокировок. Command handler этоload -> replay -> decide -> append.
ДЗ
Дальше
Следующая ката · 30. CQRS: read-models. Write-сторона готова: команды идут через Decider в стрим. Теперь read-сторона: проекции подписываются на тот же стрим и собирают денормализованные таблицы под запросы UI. read_all из сегодняшнего store наконец пригодится.
Параллельно полезно перечитать:
- 28. Decider, где доказан закон “состояние это свёртка”. Сегодня мы построили на нём хранилище.
- 22. Repository как порт, про OTP-actor и optimistic-семантику. Event store это тот же приём, но лог вместо последнего значения.