← назад к разделу

В обычном приложении одна таблица — и для записи, и для чтения. Это удобно пока данных мало. Но когда страница заказов собирает информацию из шести таблиц через пять JOIN'ов, и это происходит тысячи раз в секунду, база начинает задыхаться. Именно здесь CQRS предлагает отдельный «взгляд» на данные, заточенный только под чтение.

Что такое read-model

Представьте, что у вас есть склад (write-сторона) и витрина магазина (read-model). На складе данные хранятся в нормализованном виде — каждая сущность в своей таблице. На витрине всё уже разложено по полочкам так, как удобно покупателю: один предмет содержит всё нужное, ничего не надо искать в других местах.

Read-model — это денормализованное представление данных, где информация уже сложена в форму, удобную конечному потребителю: один SELECT без JOIN'ов возвращает готовый объект для UI или API.

Цена этого удобства: данные на «витрине» немного отстают от «склада». Обновление доезжает до витрины не в той же транзакции, а следом, обычно за 100 мс — 1 с; это и называют eventual consistency, согласованностью в конечном счёте. Но выигрыш — несравнимо меньше нагрузки на базу и в разы быстрее ответ.

Где хранить read-model

Выбор хранилища зависит от того, как данные будут читаться. Не «у нас уже есть Redis, туда и положим», а «какой паттерн чтения нам нужен».

PG-таблица — первый шаг

Денормализованная таблица в той же PostgreSQL — почти всегда правильное начало. Никакой новой инфраструктуры, привычные транзакции и индексы, обновление через outbox.

Например, вместо того чтобы каждый раз JOIN'ить order, order_item и customer, создаём одну таблицу со всем нужным:

CREATE TABLE order_summary (
    order_id        BIGINT PRIMARY KEY,
    customer_id     BIGINT NOT NULL,
    customer_name   TEXT NOT NULL,
    customer_email  TEXT NOT NULL,
    status          TEXT NOT NULL,
    item_count      INTEGER NOT NULL,
    total_amount    NUMERIC(19,4) NOT NULL,
    currency        TEXT NOT NULL,
    created_at      TIMESTAMPTZ NOT NULL,
    confirmed_at    TIMESTAMPTZ,
    shipped_at      TIMESTAMPTZ,
    updated_at      TIMESTAMPTZ NOT NULL,
    version         BIGINT NOT NULL DEFAULT 0
);

CREATE INDEX ix_os_customer    ON order_summary (customer_id, created_at DESC);
CREATE INDEX ix_os_status_date ON order_summary (status, created_at DESC);

Поле version нужно для идемпотентного обновления. В событии едет номер версии заказа, и потребитель пишет строку только если его версия свежее той, что уже лежит: ... WHERE order_id = ? AND version < ?. Если то же событие придёт дважды, второй UPDATE не изменит ни строки.

Откуда эта версия берётся. Это счётчик на корне агрегата на стороне записи: каждая успешная команда увеличивает его на единицу, и репозиторий пишет с условием «версия та же, что читали» — та самая оптимистичная блокировка. Событие уносит новое значение (aggregateVersion = 8), и проекция сравнивает его со своим.

Почему это вообще работает — важная оговорка: у агрегата один писатель. Все команды по одному заказу проходят через один обработчик и одну проверку версии, поэтому номера строго последовательны и сравнимы. Если бы запись шла двумя путями (обработчик команд и ещё какой-нибудь скрипт, правящий таблицу напрямую), версии перестали бы быть надёжными, и вся идемпотентность проекции развалилась бы. Отсюда практическое правило, которое звучит грубее, чем кажется: правка данных агрегата мимо обработчика команд запрещена не из вкуса, а потому что ломает проекции.

Что делать, если версии нет (событие пришло из чужого сервиса, где её не ведут): использовать отметку времени события и сравнивать по ней (WHERE updated_at < :occurredAt) — хуже, потому что часы на узлах расходятся, но работает; или хранить у себя идентификаторы обработанных событий и отбрасывать повтор по ним. Второе надёжнее и дороже: таблица обработанных с уборкой по возрасту.

Идемпотентность в других хранилищах

Проверка версии в условии UPDATE — приём для реляционной базы. Для двух других вариантов, предложенных выше, механизм другой, и это надо знать до выбора хранилища.

Кеш в памяти (Redis). Условного обновления «только если версия свежее» здесь нет: команды записи безусловны. Способы три. Версия в самом значении и запись через небольшой сценарий на стороне сервера, который читает текущее значение, сравнивает версию и пишет только при необходимости — атомарно, потому что сценарий выполняется целиком. Проверка и запись в транзакции с отслеживанием ключа — сложнее и с повторами при конфликте. И самый простой вариант, который часто достаточен: записывать безусловно, приняв, что «последний записавший победил» — годится, когда события одного объекта приходят по порядку (один раздел брокера), а повтор пишет то же значение. Худший случай при этом — на мгновение записалось старое значение, и следующее событие его исправит.

Поисковый движок (Elasticsearch и подобные). Здесь есть встроенный механизм, и он ровно про это: внешняя версия. Документ пишется с указанием «версия такая-то, применяй только если она больше текущей» (параметры внешнего версионирования), и движок сам отклоняет запись с устаревшей версией — с понятной ошибкой конфликта, которую потребитель трактует как «повтор, пропускаем». Это и есть правильный способ, и он же объясняет, зачем в событии нужна версия. Второй механизм — запись по номеру последовательности и первичного терма — годится для чтения-изменения-записи, а для проекций берут внешнюю версию.

Аналитическое хранилище (колоночное). Обновлений там обычно нет вовсе: события дописываются, а текущее состояние получается выборкой «последняя запись по объекту». Тогда идемпотентность не нужна — дубль просто добавит ещё одну строку с той же версией, а выборка возьмёт последнюю. Цена: запросы сложнее, зато вставка дешёвая.

Общее правило: идемпотентность проекции — свойство пары «хранилище + механизм», и выбирать её надо вместе с хранилищем, а не после. Хранилище, в котором условная запись невозможна и порядок не гарантирован, для проекции подходит хуже — независимо от того, насколько оно быстрое.

Materialized view — для тяжёлых агрегаций

Когда нужна сводка вроде «оборот по продуктам за последний месяц», и пересчитывать её каждый раз дорого:

CREATE MATERIALIZED VIEW product_revenue_daily AS
SELECT
    p.product_id,
    p.name,
    DATE(o.confirmed_at) AS day,
    SUM(oi.quantity * oi.unit_price) AS revenue,
    COUNT(DISTINCT o.id) AS order_count
FROM order_item oi
JOIN product p ON p.product_id = oi.product_id
JOIN "order" o ON o.id = oi.order_id
WHERE o.status IN ('CONFIRMED', 'SHIPPED', 'DELIVERED')
GROUP BY p.product_id, p.name, DATE(o.confirmed_at);

CREATE UNIQUE INDEX ux_prd_pk ON product_revenue_daily (product_id, day);

Обратите внимание на день: он берётся от o.confirmed_at, а не от даты создания позиции. Выручку считают по моменту, когда заказ подтверждён, — иначе один и тот же заказ попадёт в один день позициями, а в другой статусом, и при отмене задним числом цифры за прошлые дни поедут.

Обновляется по расписанию или по событию:

REFRESH MATERIALIZED VIEW CONCURRENTLY product_revenue_daily;

Про CONCURRENTLY важно знать две вещи. Первая: он работает только если у представления есть уникальный индекс — тот самый ux_prd_pk выше; без него команда откажется выполняться. Вторая: пересчёт всё равно полный. CONCURRENTLY лишь не блокирует читателей на время работы — он строит новую копию рядом и подменяет старую, а значит стоит примерно как весь запрос целиком и занимает вдвое больше места. На сотнях тысяч строк такой пересчёт укладывается в секунды, и раз в пять минут — нормальный ритм. На десятках миллионов он идёт минутами, и раз в пять минут его запускать уже нельзя: следующий пересчёт начнётся раньше, чем закончится предыдущий. Там переходят на обычную таблицу, которую потребитель событий правит построчно.

Запросы к представлению работают быстро — данные уже посчитаны.

Redis — для горячих ключей

Когда данные читаются на каждый запрос пользователя и нужна минимальная задержка. Например, проекция «пользователь → его текущий тарифный план»:

SET customer:42:plan {"plan":"PRO","expires_at":"2026-12-01"}

Потребитель событий обновляет ключ при каждом изменении подписки. Важная тонкость: здесь Redis — источник ответа, а не кеш-ускоритель. Это разные роли: кеш можно сбросить и перечитать из БД, read-model в Redis — это и есть основное хранилище данного представления.

Отсюда и отсутствие EX в команде: срок жизни ключа здесь не нужен. Поставьте EX 3600 — через час ключ исчезнет, и отвечать станет нечем, потому что перечитать неоткуда: следующее событие про эту подписку может прийти через полгода. Если хочется TTL, значит вы всё-таки пишете кеш, и тогда за ним обязан стоять запрос в базу на случай промаха.

ElasticSearch — для поиска

Полнотекстовый поиск по описаниям, фильтры по 20+ атрибутам с релевантностью — реляционные базы плохо с этим справляются. ElasticSearch хранит готовый документ:

{
  "product_id": 12345,
  "name": "Pixel 9 Pro 256GB",
  "category_path": ["Электроника", "Смартфоны", "Google"],
  "price": 99990,
  "in_stock": true,
  "rating": 4.7,
  "description": "..."
}

Потребители событий ProductCreated, ProductPriceChanged, StockUpdated обновляют документ. Запрос GET /products?q=pixel&min_rating=4.5 идёт напрямую в ES.

Схема read-model не зависит от write-стороны

Read-схема может — и часто должна — отличаться от write-схемы. Это её главное преимущество.

order_summary order order_item customer

Три таблицы записи складываются в одну строку чтения. Запрос за заказом становится выборкой одной строки без соединений — платим тем, что при смене имени клиента эту строку придётся обновлять отдельно, событием.

Поле order_summaryОткуда берётся
order_id, statusиз order
customer_name, customer_emailскопированы из customer
item_countпосчитан из order_item

Что это даёт:

  • Чтение без JOIN. SELECT * FROM order_summary WHERE order_id = ? вместо трёх JOIN'ов.
  • Индексы под конкретные запросы. Не нужно идти на компромисс между нуждами чтения и записи — у каждой стороны свои индексы.
  • Простой маппинг в коде. Read-DTO совпадает с read-схемой один в один.

Цена: если изменилось имя клиента, это изменение должно дойти до order_summary тоже — через событие и потребителя. Это нормальная плата за быстрое чтение.

Обновление через события

Read-model не обновляется в той же транзакции, что и запись. Обновляет её потребитель событий — уже после того, как запись зафиксирована.

Каким путём событие до него доедет, зависит от расстояния. Если проекция лежит в той же базе и в том же сервисе, хватает событий внутри процесса: агрегат зарегистрировал OrderConfirmed, после коммита обработчик обновил order_summary. Брокер тут не нужен. Как только проекция уезжает в другое хранилище или другой сервис, события нужно доставлять надёжно — и появляется связка outbox + Kafka. Дальше в статье разобран именно этот, дальний вариант.

агрегат + outbox одна транзакция после COMMIT outbox-relay тик 100-500 мс send + ack Kafka топик order.events poll потребитель @KafkaListener UPDATE, если версия новее order_summary version < :v

Путь одного изменения сверху вниз: атомарна только верхняя пара, дальше событие едет через брокер, поэтому строка витрины обновляется не в момент коммита.

1. command-handler меняет Order, сохраняет,
   регистрирует OrderConfirmed → outbox-таблица (атомарно)
2. outbox-relay публикует событие в Kafka
3. read-side потребитель ловит OrderConfirmed
4. UPDATE order_summary SET status = 'CONFIRMED', confirmed_at = ..., version = :v
   WHERE order_id = ? AND version < :v

Задержка обычно 100ms–1s в нормальном режиме. При проблемах с Kafka может быть больше — это ожидаемо и нормально. UI должен это понимать.

Из чего складывается это число — полезно знать, потому что каждое слагаемое управляется отдельно:

СлагаемоеПорядокЧем управляется
Ожидание в таблице исходящих0–500 мсинтервалом отправщика (обычно 100–500 мс)
Публикация в брокер1–10 мснастройками подтверждения и пакетирования у производителя
Ожидание в брокере1–50 мстем, успевает ли потребитель; при отставании — минуты
Обработка потребителем5–50 мссложностью преобразования и записью в витрину

Отсюда видно главное: основной вклад даёт интервал отправщика, и именно он — первая ручка, если задержка не устраивает. Уменьшили до 100 мс — типичная задержка стала 150–200 мс; уменьшать дальше смысла мало, потому что растёт нагрузка на базу от опроса таблицы. Альтернатива опросу — уведомление о новой строке (в PostgreSQL это LISTEN/NOTIFY или чтение журнала репликации), тогда задержка падает до единиц миллисекунд ценой более сложного отправщика.

И обратное: если наблюдаемая задержка не укладывается в эту сумму, дело не в подходе, а в отставании потребителя — то есть это поломка, а не свойство. Различить их можно только метрикой отставания, и это следующий раздел.

Почему не обновлять синхронно в той же транзакции:

  • Read-model теряет независимость. Изменение схемы order_summary начинает блокировать write-транзакции.
  • Если write-транзакция откатится — read-model уже могла измениться (особенно если read-store в другой базе).
  • Для разных баз синхронный одновременный апдейт потребует двухфазного подтверждения (2PC) — сложного протокола, который несёт больше проблем, чем решает.

Подробнее о механизме доставки событий — в Синхронизация read-model через события.

Как заметить, что проекция разошлась

Вся статья строится на том, что расхождение возможно, — и без способа его обнаружить это означает «пользователи расскажут». Три уровня наблюдения, и нужны все три, потому что они ловят разное.

1. Отставание: проекция позади, но исправна. Три числа на графике:

  • Возраст самой старой неопубликованной строки в таблице исходящих. Растёт — сломался отправщик; это самая ранняя тревога из всех, потому что она срабатывает до того, как события вообще попали в брокер.
  • Отставание потребителя (сколько сообщений он не дочитал). Растёт — потребитель не успевает или упал.
  • Задержка от события до записи в витрину: в событии есть время его возникновения, при записи в витрину считаем разницу и складываем в метрику. Это то самое число «100 мс — 1 с», измеренное по факту, а не предположенное.

Тревоги ставят на все три, и пороги берут от требований к свежести: если интерфейс обещает «сразу», порог — секунды; если «в течение минуты» — минуты.

2. Поломка: события не обрабатываются. Отдельно от отставания: потребитель жив, читает — и падает на каждом сообщении. Признаки — счётчик ошибок обработки и число сообщений, ушедших в очередь недоставленных. Последнее должно быть равно нулю и иметь тревогу на любое ненулевое значение: одно сообщение в недоставленных означает, что одна запись витрины навсегда неверна, пока кто-то не разберётся.

3. Расхождение: всё работает, а данные неверны. Самый неприятный случай: потребитель не отстаёт, ошибок нет, а витрина не соответствует записи — пропущенное событие, ошибка в преобразовании, ручная правка данных в обход обработчика. Отставание этого не покажет; нужна сверка.

Сверка делается задачей по расписанию (обычно ночью) и бывает трёх глубин — берите по цене:

-- Уровень 1: контрольные числа. Секунды на любом объёме, ловит грубые расхождения.
SELECT (SELECT count(*) FROM orders)         AS write_rows,
       (SELECT count(*) FROM order_summary)  AS read_rows,
       (SELECT count(*) FROM orders  WHERE created_at >= current_date - 1) AS write_today,
       (SELECT count(*) FROM order_summary WHERE created_at >= current_date - 1) AS read_today;

-- Уровень 2: сверка по версиям. Находит конкретные расхождения, если витрина в той же базе.
SELECT o.id, o.version AS write_version, s.version AS read_version
  FROM orders o
  LEFT JOIN order_summary s ON s.order_id = o.id
 WHERE s.order_id IS NULL            -- вообще нет в витрине
    OR s.version <> o.version        -- отстала или уехала вперёд
 LIMIT 100;

-- Уровень 3: сверка значений на выборке. Дорого, поэтому по образцу.
-- Взять 1000 случайных объектов, собрать проекцию заново и сравнить поля.

Что делать с найденным: точечно пересобрать расхождения (тот же код, что перестроение, но по списку идентификаторов) и разобраться в причине — потому что расхождение, исправленное без разбора, вернётся. Регулярное появление расхождений одного вида почти всегда означает пропущенное событие в одном из обработчиков команд.

Если витрина в другом хранилище, сверка сложнее: контрольные числа считаются с двух сторон и сравниваются в скрипте, а сверка по версиям делается выборкой идентификаторов и версий из обеих систем пачками. Это дороже, поэтому уровень 1 гоняют ежедневно, уровень 2 — еженедельно, уровень 3 — по подозрению.

И самое дешёвое из всего: счётчик «событий получено» против «событий применено» в самом потребителе. Разница означает, что часть событий отброшена — иногда законно (повторы), иногда нет. График этой разницы ловит ошибки в логике идемпотентности, которые иначе не видны вовсе.

Read-model всегда можно восстановить

Это важный принцип: для любой read-model должен существовать скрипт, который заново строит проекцию из write-стороны. Если потеряли данные в Redis или нужно перелить данные в новый ElasticSearch-индекс — скрипт проходит по всем агрегатам и собирает проекцию с нуля:

@Component
@RequiredArgsConstructor
public class OrderSummaryRebuilder {

    private final OrderRepository orderRepository;
    private final OrderSummaryRepository orderSummaryRepository;

    public void rebuildAll() {
        long lastId = 0;
        int batchSize = 1000;
        while (true) {
            List<Order> batch = orderRepository.findAllAfter(lastId, batchSize);
            if (batch.isEmpty()) break;
            List<OrderSummaryRow> rows = batch.stream().map(this::toSummaryRow).toList();
            orderSummaryRepository.upsertBatchIfNewer(rows);
            lastId = batch.getLast().id().value();
        }
    }
}

Метод называется upsertBatchIfNewer, а не просто upsertBatch, и это не придирка к имени. Перестроение почти всегда идёт при живом потребителе: пока скрипт читает первую пачку заказов, по одному из них уже приходит новое событие и обновляет строку. Простой upsert затрёт свежее состояние старым снимком — и проекция разойдётся ровно там, где её чинили. Поэтому строки пишут той же проверкой версии, что и потребитель: WHERE order_id = ? AND version < ?. Версия берётся из самого агрегата, так что запись из перестроения и запись из события сравниваются по одной шкале.

Когда это нужно:

  • Авария. Redis-кластер упал, ElasticSearch-индекс потерян, произошла миграция.
  • Новое хранилище. Решили добавить ElasticSearch — он пустой, нужно залить исторические данные.
  • Изменение схемы read-model. Добавили поле в order_summary — для старых записей оно пустое, скрипт дозаполнит.

Перестроение без простоя

Скрипт выше отвечает на вопрос «как собрать», и оставляет главный практический — как это сделать на работающей системе. Два способа, и выбор между ними не вкус.

Способ 1: чинить на месте. Скрипт идёт по объектам и обновляет строки витрины той же проверкой версии, что потребитель. Чтение всё это время работает и видит смесь: часть строк уже пересобрана, часть ещё старая. Годится, когда витрина в целом верна и правится частично (дозаполнить поле, исправить расхождения). Не годится, когда меняется смысл полей: смесь старого и нового формата означает, что интерфейс показывает неверные данные для части записей.

Способ 2: новая витрина и переключение — правильный способ для смены формата. Порядок:

  1. Создать вторую таблицу (или второй индекс/коллекцию в другом хранилище) с новой схемой.
  2. Запустить потребителя, который пишет в обе — старую и новую. Это ключевой шаг: с этого момента новая витрина получает все свежие события и больше не отстаёт.
  3. Прогнать перестроение по новой витрине (с проверкой версии, чтобы не затирать свежие события, — ровно как в коде выше).
  4. Сверить: контрольные числа и выборочная сверка значений между старой и новой.
  5. Переключить чтение на новую — по одному запросу, за флагом, с возможностью вернуться.
  6. Убрать запись в старую и удалить её — когда чтение прожило на новой неделю.

Что решается этим порядком: события, приходящие во время перестроения, не теряются (пункт 2 включён раньше перестроения) и не затираются (проверка версии в пункте 3). Это и есть ответ на самый частый вопрос.

Сколько это занимает. Перестроение — это чтение всей стороны записи и запись всей витрины, то есть время считается от объёма: на миллионе заказов пачками по тысяче с преобразованием и записью — порядка десятков минут; на сотне миллионов — часы или сутки, и тогда его гоняют в несколько потоков по диапазонам идентификаторов и с ограничением темпа, чтобы не задавить базу записи. Важно знать это число заранее: «перестроим за ночь» и «перестроим за трое суток» — разные решения. Поэтому процедуру прогоняют на копии прода до того, как она понадобится.

Что ещё нужно для спокойного перестроения: ограничение темпа (чтобы чтение стороны записи не мешало обычной работе), возобновляемость (запомнить, до какого идентификатора дошли, и продолжить после сбоя — как в коде выше через lastId), и метрика прогресса, иначе через два часа никто не знает, сколько осталось.

Схема витрины меняется по тем же правилам, что схема записи

«Добавили поле — скрипт дозаполнит» верно для поля, допускающего пустое значение, и неверно для обязательного: обязательную колонку нельзя добавить в непустую таблицу без значения по умолчанию, а работающий потребитель ничего про неё не знает. Порядок ровно тот же, что для любой миграции — сначала добавить, потом заполнить, потом начать требовать:

  1. Добавить колонку, допускающую пустое значение (или с безопасным значением по умолчанию). Потребитель старой версии продолжает работать, он про неё не знает.
  2. Выкатить потребителя, который её заполняет. С этого момента новые события пишут поле.
  3. Дозаполнить старые строки скриптом (это то самое «перестроение по месту»).
  4. Начать читать поле в запросах — только теперь, когда оно заполнено везде.
  5. Сделать обязательным — если нужно, и отдельной миграцией.

Удаление поля — обратный порядок и с той же осторожностью: перестать читать, выкатить, перестать писать, выкатить, удалить колонку. Переименование не делают вовсе: это добавление плюс удаление, и в промежутке потребитель одной из версий пишет не туда.

Отдельно про смену типа и смысла. Изменить «сумма в рублях» на «сумма в копейках» в существующей колонке нельзя: часть строк окажется в одних единицах, часть в других, и различить их будет невозможно. Правильно — новая колонка, двойная запись, дозаполнение, переключение чтения, удаление старой. То есть ровно тот же приём «расширить и сжать», что на стороне записи.

Сколько проекций заводить

Самый частый практический вопрос, на который статья показывает одну таблицу и не даёт правила. Правило такое: проекция — на класс запросов, а не на экран и не на сущность.

«Класс запросов» означает: один способ доступа, одна форма данных, один набор фильтров. Практически:

  • Список заказов покупателя (фильтр по покупателю, сортировка по дате) — одна проекция.
  • Панель поддержки (поиск по номеру, телефону, адресу, статусу, за любой период) — другая проекция: другие фильтры, другие индексы, часто другое хранилище.
  • Отчёт «оборот по продуктам за месяц» — третья: это агрегаты, а не объекты, и у них своя структура.
  • Карточка заказа — обычно не проекция: одна строка по идентификатору прекрасно читается со стороны записи запросом, и заводить под неё витрину незачем.

Признаки, что проекций слишком мало (одна на всё): в таблице появляются поля, нужные одному экрану из пяти; индексов на таблице больше пяти, и каждый под свой запрос; одно событие обновляет её по десяти разным правилам; запросы к ней всё равно содержат условия, для которых нет индекса.

Признаки, что проекций слишком много: две проекции обновляются одним и тем же кодом по одним и тем же событиям и различаются двумя полями; проекция обслуживает один запрос, который вызывается раз в день; число потребителей событий больше, чем число людей в команде.

Цена каждой проекции — то, что удерживает от размножения: свой потребитель, своя схема с миграциями, своё перестроение с известным временем, своя сверка, своё место на диске и свои тревоги. Порядок величины — день работы на завод и постоянное внимание; это не бесплатно, и потому проекцию заводят под названный класс запросов, а не «чтобы было удобно».

Ориентир по числу для сервиса среднего размера: одна-три проекции (список, поиск для поддержки, агрегаты для отчётов) плюс обычное чтение со стороны записи для всего остального. Десять проекций — почти всегда признак, что витрины заводят вместо индексов.

Если скрипта восстановления нет — read-model фактически стала единственным источником правды, а это уже неправильно.

Частые ошибки

Бизнес-логика в read-таблице. Read-model — это проекция, не место для бизнес-правил. Добавлять CHECK-ограничения с бизнес-инвариантами в read-таблицу опасно: если write-сторона создала данные по бизнес-причинам, которые не проходят проверку в read-таблице, потребитель событий упадёт с ошибкой. Инварианты живут в агрегатах, read-модель просто хранит то, что пришло.

Read-model как единственный источник правды. Если данные есть только в read-model и нет write-стороны, из которой их можно восстановить — это уже не CQRS, а просто две несогласованные системы. Источник правды всегда один — write-side агрегаты.

Обратный поток данных. Поток идёт в одну сторону: write → события → read. Read-model не должна вносить изменения в write-сторону. Если read-сторона обнаружила несогласованность — нужно залогировать или перезапустить rebuild для этой записи, но не писать напрямую в write-таблицы.

Синхронное обновление read-model в write-транзакции. Это связывает две стороны: изменение схемы read-таблицы начинает влиять на write-производительность. Правило простое: только через события, уже после коммита записи.

Коротко

  • Read-model — денормализованная проекция данных, оптимизированная под конкретные запросы чтения. Один SELECT, без JOIN'ов. Хранилище выбирают под паттерн чтения: PG-таблица для табличных запросов, materialized view для тяжёлых агрегаций, Redis для горячих ключей, ElasticSearch для полнотекстового поиска.
  • Схема read-model независима от write-схемы — это нормально и правильно. Обновляется только через события и только после коммита записи. В пределах одного сервиса и одной базы хватает событий внутри процесса, через границу сервиса или хранилища — outbox + Kafka + потребитель.
  • Задержка обновления 100ms–1s — это не баг, а ожидаемое поведение (eventual consistency). Для каждой read-model должен быть скрипт восстановления из write-стороны.
  • Read-model — проекция, не источник правды. Источник правды — write-side агрегаты.
  • Версия в событии — счётчик на корне агрегата, и он надёжен только потому, что у агрегата один писатель; правка данных мимо обработчика команд ломает идемпотентность проекций.
  • Идемпотентность зависит от хранилища: в реляционной базе условие по версии, в кеше сценарий на стороне сервера или «последний победил», в поисковом движке внешняя версия, в аналитическом — дописывание и выборка последней записи.
  • Расхождение ловят тремя уровнями: отставание (возраст неопубликованной строки, отставание потребителя, измеренная задержка), поломка (ошибки и очередь недоставленных), расхождение (ночная сверка контрольных чисел, версий и значений на выборке).
  • Перестроение без простоя: новая витрина, потребитель пишет в обе, перестроение с проверкой версии, сверка, переключение чтения за флагом, удаление старой; время перестроения знают заранее по прогону на копии.
  • Схему витрины меняют как любую: добавить допускающее пустоту поле, научить писать, дозаполнить, начать читать, и только потом требовать; смена смысла колонки — только через новую колонку.
  • Проекцию заводят на класс запросов, а не на экран: одна-три на сервис, карточка по идентификатору читается со стороны записи, а каждая проекция стоит потребителя, миграций, перестроения, сверки и тревог.

Что почитать дальше