CQRS-проекции: Daemon, SubscriptionRef, live read-model
открытый урокЭтот раздел читается без входа. Войди, чтобы отмечать прогресс, вести заметки и решать задачи в редакторе. войти
CQRS-проекции: Daemon, SubscriptionRef, live read-model
Сцена · команды и запросы говорят на разных языках
У отеля две роли. Менеджер ресепшен меняет состояние: бронирует, отменяет, оформляет заезд. Менеджер ресторана читает состояние: “сколько гостей в отеле сейчас, кто из них на полу 1, кто на полу 2”. Их интересы не пересекаются: первому важна корректность транзакций, второму нужна свежая красивая таблица.
CRUD-подход пытается обслужить обе роли одной моделью. Получается компромисс: схема, которая нормализована для писателя, но плохо читается, или денормализована для читателя, но мучительна для писателя. CQRS говорит: разведём write-сторону и read-сторону по разным моделям. Писатель кладёт события в Event Store. Читатель смотрит на read-model, специально заточенную под его экран. Между ними стоит проекция, которая трансформирует поток событий в read-model.
В этом уроке мы строим живую проекцию “occupancy”, “кто сейчас в отеле”. Она читает поток событий из Event Store, обновляется при каждом GuestCheckedIn/GuestCheckedOut, и UI подписывается на её изменения через SubscriptionRef. Без polling, без своего websocket, без задержки.
Карта урока · что заберёшь домой
- CQRS как разделение языков. Команды vs запросы, write-модель vs read-модель.
- Проекция как чистая функция.
(state, event) -> state. То же, чтоevolve, но другой state. - SubscriptionRef. Ref с встроенной подпиской на изменения. Подписчики получают каждое новое значение.
- Daemon под Effect.forkScoped. Фоновый файбер, который держит read-model в актуальном состоянии.
- Идемпотентность. Почему проекция должна выдерживать повторное применение события.
- Last_seq tracking. Как пережить рестарт без потери позиции.
- Polling против LISTEN. Когда хватит первого, когда нужен второй.
Раздел 1 · CQRS, кратко
Сцена
CQRS (Command Query Responsibility Segregation) сформулирован Грегом Янгом в 2010 году. Идея простая: команды и запросы это разные операции, и им можно дать разные модели данных.
Command path Query path
───────────── ──────────
POST /reservations GET /occupancy
│ │
▼ ▼
CommandHandler ReadModel
│ ▲
▼ │
reservationDecider ┌─── Projection (Daemon)
│ │
▼ │
EventStore (Postgres) ────────┘ events stream
Левая ветка это уроки 14-16: команда -> decider -> event store. Правая ветка это сегодняшний урок: event store -> проекция -> read model -> UI.
Между ветками есть один общий ресурс: поток событий. Write-сторона его пишет, read-сторона его читает. Read-модели не обновляются напрямую командами; они материализуются исключительно из событий, и в этом залог их корректности.
Что взять с собой
- CQRS разделяет write и read модели, не репликацию БД.
- Источник истины это events; read-модели производны от него.
- Read-модель можно пересоздать с нуля, прокрутив все события.
Раздел 2 · Проекция как чистая функция
Сцена
Проекция, по форме, это та же функция, что evolve в Decider-е. Берёт текущее состояние read-модели и событие, возвращает новое состояние. Но state другой: не агрегат, а денормализованная view.
Шаг 1 · Что в read-model
Хотим показать таблицу “кто сейчас в отеле”. Каждая строка это { reservationId, room, guest }. Гость попадает в таблицу при GuestCheckedIn, выходит при GuestCheckedOut или ReservationCancelled (после check-in отмены не бывает, но защищаемся на всякий).
// examples/ddd-hotel/effect/src/projections/occupancy.ts
export type OccupancyRow = {
readonly reservationId: ReservationId;
readonly room: RoomNumber;
readonly guest: GuestEmail;
};
export type Occupancy = ReadonlyMap<ReservationId, OccupancyRow>;
Шаг 2 · Apply
import { Match } from 'effect';
type Pending = ReadonlyMap<ReservationId, OccupancyRow>;
const apply = (
state: Occupancy,
pending: Pending,
item: { readonly streamId: ReservationId; readonly event: ReservationEvent },
): { readonly state: Occupancy; readonly pending: Pending } =>
Match.value(item.event).pipe(
Match.tag('ReservationPlaced', (e) => {
const nextPending = new Map(pending);
nextPending.set(item.streamId, {
reservationId: item.streamId, room: e.room, guest: e.guest,
});
return { state, pending: nextPending };
}),
Match.tag('ReservationCancelled', () => ({
state: removeKey(state, item.streamId),
pending: removeKey(pending, item.streamId),
})),
Match.tag('GuestCheckedIn', () => {
const row = pending.get(item.streamId);
if (row === undefined) return { state, pending };
const next = new Map(state);
next.set(item.streamId, row);
return { state: next, pending: removeKey(pending, item.streamId) };
}),
Match.tag('GuestCheckedOut', () => ({
state: removeKey(state, item.streamId),
pending,
})),
Match.exhaustive,
);
Что тут происходит:
ReservationPlaced, кладём ещё-не-заехавшего гостя в pending. В основной state он попадёт послеGuestCheckedIn.ReservationCancelled, удаляем из обоих контейнеров (вдруг успели зачеркнуть).GuestCheckedIn, перекладываем из pending в state. Гость теперь “в отеле”.GuestCheckedOut, удаляем из state. Гость выехал.
pending это вспомогательная структура самой проекции, не часть публичного state. Это нормально: read-model для UI знает только Occupancy, внутренняя книга projection-а это её приватный мир.
Шаг 3 · Идемпотентность
Apply должна выдерживать повторное применение одного и того же события. Это критично, потому что:
- После рестарта daemon может начать читать с уже обработанного места.
- Сетевые retries могут продублировать NOTIFY.
- Repair-режим (пересоздать read-model с нуля) гоняет все события заново.
Проверь: повторное ReservationPlaced перезаписывает row в pending (один и тот же row). Повторное GuestCheckedIn ищет row в pending, не находит (уже перемещён в state), возвращает state без изменений. Все ветки идемпотентны.
Что взять с собой
- Проекция это
(state, event) -> state, то же по форме чтоevolve, другой state. - Внутренние структуры (
pending) не часть публичного API read-model. - Каждая ветка обязана быть идемпотентной: повторное применение не ломает state.
Покрути проекцию руками
Прежде чем читать про daemon, прогони apply сам. Жми события у броней и смотри, как строка едет pending -> в отеле -> прочь. Повтори Placed или CheckIn дважды подряд, увидишь, как загорается флаг идемпотентности: state не изменился.
Раздел 3 · SubscriptionRef
Сцена
Ref это атомарная переменная: get, set, update. Этого хватает для shared state, но не хватает для live UI. Когда state меняется, как сказать всем подписчикам “обновитесь”? Через timer и сравнение версий? Через ручной EventEmitter? Effect даёт прямой ответ: SubscriptionRef.
Идея словами
SubscriptionRef<A> это Ref<A> плюс встроенный Stream изменений. Каждый SubscriptionRef.set(ref, value) атомарно обновляет state и публикует новое значение в SubscriptionRef.changes(ref). Подписчик читает SubscriptionRef.changes(ref) как обычный Stream и получает все обновления.
Шаг 1 · Создание
import { SubscriptionRef } from 'effect';
const ref = yield* SubscriptionRef.make<Occupancy>(new Map());
SubscriptionRef.make(initial) это эффект, который создаёт ref. Работает внутри Effect.gen.
Шаг 2 · Чтение и запись
// чтение текущего значения
const current = yield* SubscriptionRef.get(ref);
// атомарная замена
yield* SubscriptionRef.set(ref, newOccupancy);
// функциональное обновление
yield* SubscriptionRef.update(ref, (current) => apply(current, ...));
API идентично Ref, плюс дополнительно SubscriptionRef.changes(ref) это Stream.
Шаг 3 · Подписка
import { Stream } from 'effect';
const consumer = Stream.runForEach(SubscriptionRef.changes(ref), (occupancy) =>
Effect.sync(() => console.log('новая read-model:', occupancy)),
);
SubscriptionRef.changes(ref) начинает с текущего значения, потом отдаёт каждое последующее set. Если поставить два set-а быстро подряд, подписчик гарантированно увидит последнее, но не обязательно промежуточные (стратегия “skip if behind”). Это правильное поведение для UI: важно показать актуальное, а не воспроизвести анимацию.
Что взять с собой
SubscriptionRefэто Ref плюс Stream изменений.- Подписчик читает
SubscriptionRef.changes(ref)без polling. - Если подписчик не успевает, он получает последнее значение, не все промежуточные.
Раздел 4 · Daemon под Effect.forkScoped
Сцена
Кто-то должен крутить цикл “прочитал событие -> применил к state -> обновил SubscriptionRef”. Это фоновый процесс, который должен жить, пока живёт владеющий Scope (например, всё время работы программы или время жизни одного агрегата).
Шаг 1 · forkScoped, не fork
В уроке 06 мы видели Effect.forkChild. Он привязывает файбер к родительскому файберу, и имя это прямо говорит: умирает родитель, умирает и ребёнок. Это не то, что нам нужно: родитель это конструктор Layer, он умирает сразу после return {...}. Daemon должен жить дольше.
Effect.forkScoped привязывает файбер к Scope из Context. Конструктор Layer.effect даёт нам такой Scope: файбер живёт, пока живёт Scope, и корректно прерывается при его закрытии. Отдельного Layer.scoped не существует, Layer.effect сам принимает эффект, которому нужен Scope, и убирает эту потребность из требований слоя.
Шаг 2 · Полная реализация
import { Effect, Ref, Stream, SubscriptionRef } from 'effect';
export const buildOccupancyProjection = Effect.gen(function* () {
const store = yield* InMemoryEventStore; // или PgEventStore
const ref = yield* SubscriptionRef.make<Occupancy>(new Map());
const pendingRef = yield* Ref.make<Pending>(new Map());
const stream = yield* store.subscribe; // Stream<AllEventsItem>
yield* Effect.forkScoped(
Stream.runForEach(stream, (item) =>
Effect.gen(function* () {
const pending = yield* Ref.get(pendingRef);
const state = yield* SubscriptionRef.get(ref);
const next = apply(state, pending, item);
yield* SubscriptionRef.set(ref, next.state);
yield* Ref.set(pendingRef, next.pending);
}),
),
);
return ref;
});
Шаги:
- Получаем store и создаём refs.
store.subscribeэто Effect, который возвращаетStream<AllEventsItem>. Важно: подписка к PubSub устанавливается синхронно в момент yield, поэтому все события после этого момента гарантированно дойдут. Внутри store это две строки:PubSub.subscribe(pubsub)отдаётSubscription, аStream.fromSubscription(subscription)превращает подписку в стрим. Подписка это отдельный тип, а не очередь, поэтому и конструктор стрима для неё свой.Effect.forkScopedзапускает фоновый файбер, который крутитStream.runForEach. Для каждого события мы читаем state, применяемapply, пишем обратно.- Возвращаем
SubscriptionRef. Caller подпишется черезSubscriptionRef.changes(ref).
Шаг 3 · Использование
const program = Effect.gen(function* () {
const occupancy = yield* buildOccupancyProjection;
// подписка
yield* Effect.forkScoped(
Stream.runForEach(SubscriptionRef.changes(occupancy), (current) =>
Effect.sync(() => console.log('гостей сейчас:', current.size)),
),
);
// ... тем временем кто-то пишет события в store
}).pipe(Effect.scoped, Effect.provide(InMemoryEventStoreLive));
Effect.scoped создаёт корневой Scope. Когда программа завершается, Scope закрывается, и оба наших forkScoped-файбера корректно прерываются.
Что взять с собой
Effect.forkScopedпривязывает файбер к Scope, не к родительскому файберу.- В
Layer.effectили вEffect.genплюсEffect.scopedэто правильный выбор для daemon-а. - Daemon крутит
Stream.runForEach, на каждом событии обновляет state.
Раздел 5 · Тест проекции
Сцена
Проекцию тестируем end-to-end: подкладываем InMemoryEventStore, append-аем события, читаем SubscriptionRef, проверяем содержимое. Никаких моков, реальная mini-инфраструктура.
Шаг 1 · Тестовый сценарий
// examples/ddd-hotel/effect/test/projections.test.ts
import { Effect, SubscriptionRef } from 'effect';
import { describe, expect, it } from 'vitest';
const placed: ReservationEvent = {
_tag: 'ReservationPlaced', reservationId: id, guest, room, range, occurredAt: at,
};
const checkedIn: ReservationEvent = { _tag: 'GuestCheckedIn', reservationId: id, occurredAt: at };
const checkedOut: ReservationEvent = { _tag: 'GuestCheckedOut', reservationId: id, occurredAt: at };
const waitTick = Effect.sleep('20 millis');
describe('occupancy projection', () => {
it('после CheckIn read-model содержит гостя, после CheckOut пустой', async () => {
const program = Effect.gen(function* () {
const ref = yield* buildOccupancyProjection;
const store = yield* InMemoryEventStore;
yield* store.append(id, 0, [placed, checkedIn]);
yield* waitTick;
const afterCheckIn = yield* SubscriptionRef.get(ref);
expect(afterCheckIn.size).toBe(1);
expect(afterCheckIn.get(id)?.room).toBe(room);
yield* store.append(id, 2, [checkedOut]);
yield* waitTick;
const afterCheckOut = yield* SubscriptionRef.get(ref);
expect(afterCheckOut.size).toBe(0);
}).pipe(Effect.scoped, Effect.provide(InMemoryEventStoreLive));
await Effect.runPromise(program);
});
});
Заметь waitTick = Effect.sleep('20 millis'). Это не “плохой sleep в тесте”, а необходимая пауза на eventually-consistent системе. Append-аем мы из одного файбера, проекция читает из другого. Чтобы daemon успел обработать событие, нам нужно отдать ему квант времени. В реальной системе UI и так увидит обновление с задержкой.
Альтернатива sleep: подписаться на SubscriptionRef.changes(ref) и ждать первого значения с нужным размером. Это надёжнее, но многословнее. Для иллюстрации sleep понятнее.
Шаг 2 · Запуск
cd examples/ddd-hotel/effect
pnpm test
Если всё собрано правильно (decider, in-memory store, проекция), тест проходит за пятьдесят миллисекунд. Можешь повторить append с теми же событиями, проверить идемпотентность: после повторного ReservationPlaced row в pending перезатёрся, после повторного GuestCheckedIn state не изменился (потому что pending уже пуст).
Что взять с собой
- Проекцию тестируем end-to-end с in-memory store.
Effect.sleepв тесте это пауза для daemon-а, не плохой паттерн.- Альтернатива: подписаться на
SubscriptionRef.changes(ref)и ждать заданного значения.
Раздел 6 · Идемпотентность и last_seq
Сцена
В реальной системе daemon может рестартнуться (CI deploy, OOM, плановая перезагрузка). Что произойдёт с проекцией?
Два сценария:
- Naive: проекция стартует с пустой read-model, и каждое чтение начинается с
id = 0. Это работает для маленьких систем, но на миллионе событий запуск превращается в часы. - Snapshot: проекция периодически сохраняет state в Postgres вместе с
last_seq. После рестарта читает snapshot, потом дочитывает события сid > last_seq.
Шаг 1 · Таблица projection_state
create table if not exists projection_state (
name text primary key,
last_seq bigint not null,
state jsonb not null,
updated_at timestamptz not null default now()
);
name это имя проекции (occupancy, daily_occupancy, и так далее). last_seq это id последнего обработанного события из events. state это сама материализованная read-model.
Шаг 2 · Цикл с восстановлением
Псевдокод (полную реализацию делаешь в ДЗ):
const buildOccupancyProjection = Effect.gen(function* () {
const sql = yield* PgClient.PgClient;
const ref = yield* SubscriptionRef.make<Occupancy>(/* загрузить из projection_state.state */);
let lastSeq = /* загрузить из projection_state.last_seq */;
// catchup: прочитать всё, что прибавилось с момента сохранения
const missed = yield* sql`select id, payload from events where id > ${lastSeq} order by id`;
for (const row of missed) {
/* apply, обновить state */
lastSeq = row.id;
}
yield* SubscriptionRef.set(ref, newState);
yield* saveState(name, lastSeq, newState);
// live: подписаться на свежее
yield* Effect.forkScoped(
/* polling каждые 100ms или LISTEN/NOTIFY */
Effect.gen(function* () {
const newEvents = yield* sql`select id, payload from events where id > ${lastSeq} order by id`;
for (const row of newEvents) {
/* apply, обновить ref, persist */
}
}).pipe(Effect.repeat({ schedule: Schedule.spaced('100 millis') })),
);
return ref;
});
Идемпотентность тут критична: при рестарте мы можем дважды применить событие на грани last_seq, и проекция не должна сломаться.
Что взять с собой
- Без persistence projection-state проекция начинает с нуля при каждом рестарте.
projection_state(name, last_seq, state)сохраняет место и снимок.- Идемпотентность apply страхует от дублирующего применения на границе catchup.
Раздел 7 · Polling vs LISTEN/NOTIFY
Сцена
Чтобы daemon знал, что в events появилась новая строка, есть два пути.
Polling
SELECT ... WHERE id > $last_seq каждые 100ms. Просто, переносимо, работает на любой БД. Минус: лишние запросы при отсутствии нагрузки, и задержка до 100ms на каждое событие.
LISTEN/NOTIFY
Postgres-специфичный механизм. Триггер на INSERT INTO events шлёт NOTIFY events_channel. Daemon делает LISTEN events_channel и получает push-уведомление мгновенно.
create or replace function notify_events() returns trigger as $$
begin
perform pg_notify('events_channel', new.stream_id || ':' || new.id);
return new;
end;
$$ language plpgsql;
create trigger events_notify after insert on events
for each row execute function notify_events();
Плюс: нулевая задержка. Минус: NOTIFY ограничен 8000 байт payload, требует постоянного соединения, не работает за PgBouncer в transaction-mode.
Гибрид
Хороший production-паттерн: LISTEN для уведомления “появилось новое”, catchup через SELECT. NOTIFY это сигнал (“проверь events”), не сами данные. Это снимает ограничения payload и работает корректно при пропусках уведомлений.
В нашем уроке мы оставляем in-memory PubSub для in-memory store и polling для Postgres. LISTEN/NOTIFY это ДЗ для тех, кому хочется глубже.
Что взять с собой
- Polling прост, переносим, имеет задержку.
- LISTEN/NOTIFY мгновенный, но Postgres-специфичный.
- Гибрид: NOTIFY это сигнал, SELECT это данные. Дополнительно страхует от пропусков.
Чек-лист
ДЗ
Финал блока
Этот урок закрывает четвёрку Effect.ts-трека по functional DDD. По карте:
- 14. DDD-типы: branded, opaque, smart constructors. Невозможные состояния невыразимы.
- 15. Decider pattern:
decide+evolve+initial. Бизнес-логика в трёх чистых функциях. - 16. Event Store на @effect/sql-pg: append-only, optimistic concurrency, command handler.
- 17. CQRS-проекции: SubscriptionRef, Daemon, live read-model.
У тебя в examples/ddd-hotel/effect/ сейчас лежит полный мини-домен Hotel Booking: пять состояний, четыре команды, четыре события, Event Store, проекция. Около 350 строк прикладного кода, и тесты, которые гоняются за полсекунды.
Что дальше: тот же домен в Gleam-треке, потом в Haskell и Rust. После четырёх треков в 16-arch появится один сводный урок, где мы сравним четыре языка на одном и том же Decider-е и поймём, где какой язык удобнее, где проигрывает, и что общего во всех вариантах ФП-DDD.
Дальше
Раздел продолжается блоком прод-обвязки, и он начинается с 18 · Observability: логи, спаны, метрики и экспорт наружу. Дальше HttpClient и API, конфигурация и секреты, Schema вглубь, приёмники потоков, стандартная библиотека и инструменты.
Полезно перечитать:
- 09-stream, про
StreamиSink. Проекции это classic stream-processing, а приёмники разбираются в 23 · Sink и продвинутый Stream. - 11-runtime-and-schedule, про
Schedule.spacedиSchedule.exponential. Polling-цикл этоEffect.repeat(spaced(100 millis)). - Greg Young, “8 Lines of Code”. Видео, в котором за десять минут показывается, как CQRS убирает бойлерплейт.