Async-проекции и CQRS на Rust: broadcast, tokio, watch
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
Async-проекции и CQRS на Rust: broadcast, tokio, watch
Сцена · писать и читать это разные задачи
Журнал событий отлично хранит правду, но плохо отвечает на вопросы. «Кто сейчас заселён?», «сколько занято на третьем этаже?», «выручка за день?». Чтобы ответить, пришлось бы реплеить все стримы и фильтровать, на каждый запрос, и это не масштабируется.
Отсюда идея: модель для записи и модель для чтения не обязаны совпадать. Запись это поток команд, которые порождают события и держат инварианты. Чтение это вопросы, на которые надо отвечать быстро. Разнести их, это CQRS: command-сторона (наш Decider и event store) и query-сторона (проекции). Это финальный урок трека: построим проекцию, которая держит read-model свежим сама.
Карта урока
- CQRS кратко. Почему read-model отдельная и чем за это платят.
- Проекция как чистый фолд.
(snapshot, event) -> snapshot, без асинхронности. - Идемпотентность. Почему фолд должен переживать повтор события.
- Live-проекция.
tokio::spawnповерхbroadcast, публикация вwatch. - watch как живой снимок. Rust-аналог подписки без опроса.
- Тест проекции. Гоним события через store и читаем live-view.
Раздел 1 · Проекция как чистый фолд
Сердце проекции, это чистая функция, и она устроена ровно как evolve из Decider: берёт текущий снимок read-модели и одно событие, возвращает новый снимок. Никакой асинхронности, никакого I/O, поэтому её можно проверять обычными unit-тестами.
Наша read-модель, это occupancy: кто сейчас в номерах. Тонкость в том, что комнату и гостя задаёт ReservationPlaced, а заселение это GuestCheckedIn позже. Значит, между ними данные надо где-то подержать. Заводим pending для забронированных, но ещё не заселённых:
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct OccupancyRow {
pub reservation_id: ReservationId,
pub room: RoomNumber,
pub guest: GuestEmail,
}
/// Текущие заселённые гости по id брони.
pub type Occupancy = HashMap<ReservationId, OccupancyRow>;
/// Снимок проекции: материализованный read-model плюс забронированные, но
/// ещё не заселённые брони (pending).
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct OccupancySnapshot {
pub occupancy: Occupancy,
pub pending: HashMap<ReservationId, OccupancyRow>,
}
Сам фолд, это match по событию: Placed кладёт строку в pending, CheckedIn переносит её в occupancy, CheckedOut и Cancelled убирают отовсюду:
/// Чистый фолд проекции: один шаг (snapshot, item) -> snapshot.
pub fn apply(mut snapshot: OccupancySnapshot, item: &AllEventsItem) -> OccupancySnapshot {
let id = &item.stream_id;
match &item.event {
ReservationEvent::ReservationPlaced { room, guest, .. } => {
snapshot.pending.insert(id.clone(), OccupancyRow {
reservation_id: id.clone(),
room: room.clone(),
guest: guest.clone(),
});
}
ReservationEvent::GuestCheckedIn { .. } => {
if let Some(row) = snapshot.pending.remove(id) {
snapshot.occupancy.insert(id.clone(), row);
}
}
ReservationEvent::GuestCheckedOut { .. }
| ReservationEvent::ReservationCancelled { .. } => {
snapshot.occupancy.remove(id);
snapshot.pending.remove(id);
}
}
snapshot
}
Read-модель это производная журнала, а не второй источник истины. Удалил проекцию, перестроил из событий, и получил ровно то же. Можно завести десять разных проекций над одним журналом, под десять разных вопросов.
Раздел 2 · Идемпотентность
Один момент, который кусает в проде. Доставка событий проекции почти никогда не бывает «ровно один раз»: при рестарте, переподписке или сбое одно и то же событие может прийти дважды. Поэтому фолд обязан быть идемпотентным: повтор не должен ломать модель.
Наш apply почти идемпотентен по построению, потому что работает через insert и remove по ключу: повторный Placed перезапишет ту же строку, повторный CheckedOut удалит уже отсутствующий ключ, и это no-op. Опасны были бы операции вроде счётчик += 1: их повтор задвоил бы счёт, и тогда понадобился бы last_seq, номер последнего обработанного события, мимо которого старое не пройдёт. Это материал домашки.
Раздел 3 · Live-проекция на tokio
Теперь оживим фолд. Проекция должна сама догонять журнал: подписаться на поток событий и применять apply по мере поступления. В прошлом уроке in-memory store на каждый append рассылал события в broadcast-канал. Подпишемся на него и запустим фоновую задачу через tokio::spawn:
/// Запустить live-проекцию occupancy. Берёт подписку на broadcast event store,
/// спавнит tokio-задачу, фолдит apply по событиям, публикует свежий Occupancy
/// в watch. Возвращает watch::Receiver: вызывающий всегда видит последний снимок.
pub fn run_occupancy_projection(
mut events: broadcast::Receiver<AllEventsItem>,
) -> watch::Receiver<Occupancy> {
let (tx, rx) = watch::channel(Occupancy::new());
tokio::spawn(async move {
let mut snapshot = OccupancySnapshot::default();
loop {
match events.recv().await {
Ok(item) => {
snapshot = apply(snapshot, &item);
// Если все читатели ушли, публиковать больше некому.
if tx.send(snapshot.occupancy.clone()).is_err() {
break;
}
}
// Отстали от буфера: продолжаем со следующего события.
Err(broadcast::error::RecvError::Lagged(_)) => continue,
// Источник закрыт: завершаемся.
Err(broadcast::error::RecvError::Closed) => break,
}
}
});
rx
}
Здесь спрятаны два важных решения. Первое: ветка Lagged это backpressure наоборот. broadcast хранит ограниченный буфер, и если проекция отстала сильнее, чем буфер, старые события для неё теряются, а канал говорит «ты отстал на N». Мы это явно обрабатываем (continue), а не делаем вид, что не бывает: в проде это сигнал, что пора перестраивать проекцию с нуля. Второе: задача завершается сама, когда watch::send возвращает ошибку, то есть когда ушли все читатели, и фоновый таск не висит зомби.
Раздел 4 · watch как живой снимок
Почему результат это watch::Receiver<Occupancy>, а не, скажем, mpsc? Потому что потребителю проекции не нужна вся история обновлений, ему нужен последний снимок. watch хранит ровно одно текущее значение: читатель берёт borrow() и видит актуальное, или ждёт changed().await следующего обновления, не опрашивая канал в цикле. Промежуточные снимки не копятся: подключился, сразу видишь последний. Это и есть живой read-model без опроса базы.
Цепочка целиком выглядит так: команда → decide → append в store → store рассылает событие в broadcast → проекция применяет apply → публикует снимок в watch → UI читает borrow(). Заметь односторонность: события текут от записи к чтению и никогда обратно. Это и есть та самая eventual consistency: между записью и появлением в проекции проходит короткий момент, и это сознательная цена CQRS.
Раздел 5 · Тест проекции
Проверим всё вместе: прогоним Placed плюс CheckedIn через in-memory store и убедимся, что live-view показывает гостя, а CheckedOut его убирает. Подписку берём до append-ов, иначе ранние события не дойдут:
#[tokio::test]
async fn occupancy_view_tracks_check_in_and_check_out() {
let store = InMemoryEventStore::new();
// Подписка до append-ов, иначе ранние события не дойдут до проекции.
let mut view = run_occupancy_projection(store.subscribe());
store.append(&id(), 0, vec![placed(), checked_in()]).await.unwrap();
// Ждём, пока проекция применит оба события и опубликует непустой снимок.
view.changed().await.unwrap();
while view.borrow().is_empty() {
view.changed().await.unwrap();
}
assert_eq!(view.borrow().get(&id()).unwrap().room.as_str(), "101");
store.append(&id(), 2, vec![checked_out()]).await.unwrap();
// Ждём, пока CheckOut очистит read-model.
view.changed().await.unwrap();
while !view.borrow().is_empty() {
view.changed().await.unwrap();
}
assert!(view.borrow().is_empty());
}
Тест асинхронный, но через changed().await он детерминирован: не спим на таймере, а ждём ровно следующего обновления снимка. Это правильный способ тестировать live-проекцию, без хрупких sleep.
Что ты собрал за трек
Останови взгляд на пройденном. Четыре урока назад был только Rust и идея «давай смоделируем отель». Сейчас у тебя на руках: словарь домена на newtype с инвариантами в типах, чистый Decider с decide, evolve и replay, event store на Postgres с оптимистичной блокировкой и live-проекции CQRS. Это не игрушка: ровно по этой схеме строят реальные event-sourced системы, здесь она целиком помещается в голове и в один читаемый пакет examples/ddd-hotel/rust. Три тезиса, которые стоит унести на уровне рефлекса:
- Невозможные состояния невыразимы. Тип это спецификация домена, а не контейнер для данных.
- Команды и события это два разных языка. Запрос в настоящем времени можно отвергнуть, факт в прошедшем нет.
- Decider универсальнее агрегата.
decideплюсevolveплюсinitial, и из этой тройки собирается всё: агрегат, процесс-менеджер, FSM, проекция.
Тот же домен живёт на четырёх языках курса (Gleam, Effect, Haskell, Rust), и если открыть их рядом, видно, что меняется синтаксис, а паттерн один. Это и была цель: DDD как мышление, не как фреймворк.
Чек-лист
Домашка
Финал трека
Трек Functional DDD на Rust закрыт, и вместе с ним раздел про Rust целиком. Тот же домен бронирования отеля ты теперь видел бы на любом из четырёх языков курса одинаково ясно, потому что выучил не синтаксис, а способ мышления: типы как спецификация, события как факты, Decider как ядро. Это переносится на любой стек и любой язык, и останется с тобой, когда конкретные библиотеки в коде давно сменятся.