Раздел 27 · Gleam на практике

CQRS: разделить запись и чтение, проекция как свёртка

senior~30 мин

открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти

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, про границы. Проекция тоже может жить в своём контексте и слушать чужие события через явный контракт.