CQRS: разделить запись и чтение, проекция как свёртка
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
CQRS: разделить запись и чтение, проекция как свёртка
Четвёртая ката второго блока, финал релиза. Write-сторона готова: команды через Decider пишут факты в стрим. Сегодня строим read-сторону. Один и тот же стрим событий, но свёрнутый под вопрос UI, а не под состояние агрегата.
Сцена · агрегат это плохая модель для чтения
После урока 29 состояние брони это свёртка её стрима. Это идеально для записи: команда грузит одну бронь, проверяет инвариант, дописывает факт. Но попробуй ответить на вопрос UI: “кто сейчас в отеле?”
С event-sourced агрегатом это больно. Нет таблицы “все брони”. Чтобы собрать список заселённых, надо взять все стримы, свернуть каждый в состояние, отфильтровать CheckedIn. Дорого, неудобно, и форма данных (FSM одной брони) не совпадает с формой ответа (плоский список гостей).
Корень в том, что запись и чтение хотят разных моделей. Запись хочет нормированный агрегат с инвариантами. Чтение хочет денормализованную таблицу под конкретный экран. Натягивать одну модель на оба, значит, проигрывать на обоих.
CQRS признаёт это разделение и делает его явным. Write-сторона это Decider и стрим (уроки 28-29). Read-сторона это проекции: функции, которые сворачивают тот же стрим в read-models под запросы.
Карта урока · что заберёшь через 3 часа
- Поймёшь, зачем физически разделять модель записи и чтения.
- Напишешь проекцию
in_house_guestsкак свёртку глобального лога. - Увидишь, что проекция это тот же
fold, чтоevolve, но над другим состоянием и под другой вопрос. - Поймёшь eventual consistency: почему read-модель законно отстаёт от write-стороны.
Концепт · проекция это свёртка под вопрос
Проекция берёт поток событий и сворачивает его в read-model. Структурно это list.fold(events, empty, apply), ровно как replay у Decider. Разница в том, что мы накапливаем и под какой вопрос.
evolve(урок 28) сворачивает стрим одной брони в её состояние (Reserved,CheckedIn). Вопрос: “в каком состоянии эта бронь?”in_house_guestsсворачивает глобальный лог всех броней в список заселённых. Вопрос: “кто сейчас в отеле?”
Одни и те же события, две разные свёртки, два разных ответа. В этом сила: добавить новый экран значит написать новую проекцию над тем же логом, ничего не меняя в записи.
Один нюанс делает проекцию in_house_guests интереснее тривиального фильтра. Событие GuestCheckedIn несёт только id, без комнаты и гостя. Чтобы строка read-model была полной, проекция держит вспомогательный индекс pending забронированных, но ещё не заехавших, и при check-in достаёт детали оттуда:
// examples/ddd-hotel/gleam/src/projections/occupancy.gleam
pub type OccupancyRow {
OccupancyRow(reservation_id: ReservationId, room: RoomNumber, guest: CustomerId)
}
pub type Occupancy {
Occupancy(
in_house: Dict(ReservationId, OccupancyRow),
pending: Dict(ReservationId, OccupancyRow),
)
}
Языковые механики Gleam · apply и project
Сердце проекции это apply(state, item): одно событие двигает read-model. Это зеркало evolve, но накапливает таблицу, а не FSM:
pub fn apply(state: Occupancy, item: StreamItem) -> Occupancy {
let id = item.stream_id
case item.event {
ReservationPlaced(_, guest, room, _, _) ->
Occupancy(..state, pending: dict.insert(
state.pending, id,
OccupancyRow(reservation_id: id, room: room, guest: guest),
))
GuestCheckedIn(_, _) ->
case dict.get(state.pending, id) {
Error(Nil) -> state
Ok(row) ->
Occupancy(
in_house: dict.insert(state.in_house, id, row),
pending: dict.delete(state.pending, id),
)
}
GuestCheckedOut(_, _) ->
Occupancy(..state, in_house: dict.delete(state.in_house, id))
ReservationCancelled(_, _) ->
Occupancy(
in_house: dict.delete(state.in_house, id),
pending: dict.delete(state.pending, id),
)
}
}
pub fn project(log: List(StreamItem)) -> Occupancy {
list.fold(log, empty(), apply)
}
project это вся CQRS read-side в одну строку: list.fold(log, empty, apply). Бронь попадает в pending при ReservationPlaced, переезжает в in_house при GuestCheckedIn, уходит при выезде или отмене. Read-model всегда отражает “кто заехал и ещё не выехал”.
Задача · сигнатуры
Собери projections/occupancy.gleam:
pub type OccupancyRow
pub type Occupancy
pub fn empty() -> Occupancy
pub fn apply(state: Occupancy, item: StreamItem) -> Occupancy
pub fn project(log: List(StreamItem)) -> Occupancy
pub fn in_house_guests(state: Occupancy) -> List(OccupancyRow)
pub fn count(state: Occupancy) -> Int
Контракт:
- В
in_houseпопадают только заселённые гости (GuestCheckedIn), не просто забронировавшие. GuestCheckedOutиReservationCancelledубирают из read-model.projectчистая: один и тот же лог всегда даёт одну и ту же read-model. Это её главное свойство, его и тестируем.
Решай прямо здесь: тесты прогонятся в песочнице, а подсказки и разбор ниже открывай, только если застрял.
Подсказки
- Источник событий это
store.read_all()из урока 29: глобальный лог всех стримов по порядку. Проекция читает весь лог, а не один стрим, потому что отвечает на вопрос про весь отель. pendingнужен только потому, чтоGuestCheckedInбеден на данные. Если бы событие несло комнату и гостя, индекс был бы не нужен. Это напоминание: события проектируй так, чтобы read-сторона могла их использовать.Occupancy(..state, in_house: ...)это record update: меняем одно поле, остальные копируем. Иммутабельно.countэтоdict.size(state.in_house). Не считайpending, там ещё не заехавшие.
Разбор · проекция как чистый fold
Тест показывает суть CQRS: гость B забронировал, но не заехал, и его нет в read-model, хотя его события в стриме есть:
// examples/ddd-hotel/gleam/test/projections_test.gleam
pub fn only_checked_in_guests_are_in_house_test() {
let assert Ok(#(store, _stop)) = in_memory.start()
let d = decider.reservation_decider(now)
// Гость A: забронировал и заехал.
let #(id_a, place_a) = place_command("0001", room: "101")
let assert Ok(_) = event_sourcing.handle_command(store, d, id_a, place_a)
let assert Ok(_) =
event_sourcing.handle_command(store, d, id_a, CheckInGuest(id_a))
// Гость B: только забронировал.
let #(id_b, place_b) = place_command("0002", room: "102")
let assert Ok(_) = event_sourcing.handle_command(store, d, id_b, place_b)
let view = occupancy.project(store.read_all())
occupancy.count(view) |> should.equal(1)
}
Write-сторона записала события обоих гостей. Read-сторона показывает одного. Это не баг, это разделение: стрим хранит всё, проекция отвечает на конкретный вопрос (“кто заселён”), и для неё гость B пока невидим.
Концепт · eventual consistency
Проекция собирается из лога, и собирается когда-то. В нашем коде мы зовём project(store.read_all()) синхронно, поэтому read-model всегда свежая. Но в проде проекция живёт отдельно: она подписана на стрим и обновляется по мере прихода событий, с задержкой.
Eventual consistency это сознательный размен: read-model отстаёт от записи на доли секунды, и за это мы получаем независимость сторон. Write-сторона пишет факты, не зная и не заботясь, сколько проекций их читают. Можно добавить десять read-models на тот же стрим, и запись не станет медленнее.
Это значит, что UI иногда показывает слегка устаревшую картину. Для “кто в отеле” это нормально: гость заехал секунду назад, список обновится через мгновение. Где нужна строгая согласованность (баланс счёта перед списанием), читают write-сторону напрямую через current_state, а не проекцию.
Критика · что не покрыто
- Проекция пересобирается с нуля.
project(read_all())каждый раз сворачивает весь лог. В проде проекция инкрементальна: держит состояние и применяет только новые события по курсору (ДЗ). - Проекция не персистится. Read-model живёт в памяти, после рестарта пересобирается из лога. Для больших логов её материализуют в таблицу и хранят позицию курсора.
- Нет live-подписки. В эффект-треке проекция это daemon на
SubscriptionRef, UI читает изменения без опроса. В Gleam это actor, держащий read-model и подписанный на события store (за пределами этой каты). - Одна проекция. Реальный отель хочет
daily_occupancy,arrivals_today,unsettled_folios. Каждая это новая свёртка над тем же логом (ДЗ). - read_all в памяти. Наследие урока 29: для больших логов нужен курсор, а не “прочитать всё”.
Takeaway
Одна фраза:
CQRS разделяет модель записи и чтения: write это Decider плюс стрим, read это проекции. Проекция это
list.fold(log, empty, apply), тот же fold, чтоevolve, но под вопрос UI, и она eventually consistent с записью.
Конец второго блока на чтение
Релиз R3 закрыт. К урокам 22-26 (Repository, контексты, HTTP, SQLite, конфиг) добавились четыре кирпича event-sourced архитектуры: Event Modeling как способ договориться, Decider как единица домена, event store как источник истины, CQRS как разделение чтения и записи. Дальше (R4) идут process manager и saga на границе контекстов, aggregateless ES и финальный разбор на Uno.
ДЗ
Дальше
Следующая ката открывает релиз R4 · 31. Process Manager и Saga. Выходим за границу одного контекста: событие GuestCheckedOut во Front Desk запускает команду IssueFolio в Billing, и если оплата падает, шлём компенсацию. Process-manager это снова Decider, но над собственным маленьким состоянием.
Параллельно полезно перечитать:
- 29. Event Sourcing, где появился
read_all. Сегодня он наконец пригодился: проекция читает весь лог. - 23. Bounded Contexts, про границы. Проекция тоже может жить в своём контексте и слушать чужие события через явный контракт.