Пользователь подтвердил заказ, открыл список — а там заказ всё ещё «новый». Через секунду — уже «подтверждённый». Бывает хуже: в списке заказ «подтверждён», а транзакция на самом деле откатилась. Так снаружи выглядит главная трудность CQRS: данные лежат в двух местах, и read-store с проекциями должен узнавать об изменениях в write-store — без потерь, без выдумок и в правильном порядке.
В этой статье — путь события от строки в базе до проекции: где оно теряется или приходит дважды, что делать с событием, которое останавливает всех, откуда берётся окно отставания и как заметить, что синхронизация встала.
Почему нельзя обновить проекцию той же транзакцией
Первый порыв — добавить в обработчик команды строчку, которая обновляет проекцию сразу после сохранения агрегата:
@Transactional
public OrderId handle(ConfirmOrderCommand cmd) {
Order order = orderRepository.findById(cmd.orderId());
order.confirm();
orderRepository.save(order);
orderSummaryRepository.updateStatus(order.id(), "CONFIRMED"); // так не надо
return order.id();
}
Пока проекция лежит в той же базе, это даже работает. Ломается ровно в тот момент, ради которого проекцию и заводили: read-store переезжает в другую базу, в Elasticsearch или в другой сервис — и общей транзакции больше нет. Одна запись пройдёт, другая упадёт, откатить первую некому. Цена есть и до переезда: ошибка в UPDATE order_summary откатывает подтверждение заказа, хотя заказ ни при чём.
Триггер на write-таблицу — честный вариант для простой денормализованной таблицы в той же PostgreSQL, и хаб раздела его называет. Но у него две цены. Первая — невидимость: программист читает Java-код и не видит, что UPDATE order пишет ещё в одну таблицу. Вторая — массовые операции: UPDATE на миллион строк вызовет триггер миллион раз, и каждый вызов — отдельная запись в read-таблицу; операция, которая шла секунды, идёт минуты. А после переезда проекции в другую базу триггер не работает вовсе.
Связь, которая переживает переезд, — событие.
Outbox: событие ложится в базу вместе с заказом
Отправить событие в Kafka прямо из обработчика нельзя по той же причине: база и брокер — две системы без общей транзакции. Отправили до коммита, транзакция откатилась — в проекции появится подтверждение, которого не было. Отправили после коммита, а брокер недоступен или под погиб между коммитом и отправкой — заказ подтверждён, проекция не узнает об этом никогда. Повторять send в цикле не спасает: под может умереть посреди цикла.
Outbox убирает брокер из транзакции. Событие записывается не в Kafka, а в таблицу outbox — в той же базе и той же транзакции, что и изменение агрегата. База гарантирует: либо есть и новый статус, и событие, либо нет ни того, ни другого.
Атомарность есть только внутри пунктирной рамки: заказ и событие записываются одной транзакцией. Дальше событие едет по правилу «хотя бы один раз», и родиться дубль может в двух местах — в relay и у потребителя. Отбивает его проверка версии на последнем шаге.
Гарантия атомарности заканчивается на коммите. Дальше событие едет по правилу «хотя бы один раз»: каждый участник может доставить его повторно, но не может потерять. Где рождаются дубли и кто их отбивает — на схеме, про каждое место ниже.
Relay: кто доносит событие до брокера
После коммита в дело вступает outbox-relay — фоновый бин, который каждым тиком забирает неопубликованные строки и отправляет их в Kafka:
@Slf4j
@Component
@RequiredArgsConstructor
public class OutboxRelay {
private final OutboxRepository outboxRepository;
private final KafkaTemplate<String, String> kafkaTemplate;
@Scheduled(fixedDelay = 200)
@Transactional
public void relay() {
List<OutboxRecord> pending = outboxRepository.lockPending(100);
for (OutboxRecord record : pending) {
try {
kafkaTemplate.send(record.topic(), record.key(), record.payload())
.get(5, TimeUnit.SECONDS);
} catch (Exception e) {
log.warn("событие {} не ушло, вернёмся к нему на следующем тике",
record.id(), e);
break;
}
outboxRepository.markPublished(record.id());
}
}
}
Две детали здесь несут всю надёжность, и обе легко потерять.
Ждём подтверждения брокера. kafkaTemplate.send асинхронен: возвращает CompletableFuture до того, как брокер что-либо подтвердил. Если поставить markPublished сразу за ним, при недоступной Kafka relay зафиксирует «опубликовано» для событий, которые никуда не ушли, — ровно та потеря, ради которой outbox и заводили. Поэтому на future вызывают get с таймаутом: не дождались — выходим из цикла, ничего не помечаем, событие уедет на следующем тике. Здесь же рождается первый дубль: брокер записал событие, а под погиб до markPublished — на следующем тике оно уйдёт второй раз. Это плата за «не потерять».
Берём строки с блокировкой. lockPending — это не SELECT ... LIMIT 100, а:
SELECT * FROM outbox
WHERE published_at IS NULL
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED
В проде сервис живёт в нескольких экземплярах, и relay крутится в каждом. Без FOR UPDATE все они выберут одни и те же сто строк и отправят каждое событие столько раз, сколько подов. SKIP LOCKED говорит: строки, которые уже взял сосед, молча пропусти и возьми следующие.
Порядок событий держится на ключе партиции
Симптом: проекция показывает заказ «новым», хотя он подтверждён, а в логе потребителя — «нет строки проекции для заказа 42». OrderConfirmed пришёл раньше OrderCreated. Обработчик по версии, разобранный ниже, держится на том, что события одного заказа приходят в порядке, в каком случились. Kafka этого сама по себе не обещает.
Топик разбит на партиции, и порядок гарантирован только внутри одной. В какую партицию попадёт сообщение, решает ключ: продюсер берёт хеш ключа по модулю числа партиций. Relay передаёт ключом aggregate_id из строки outbox — тот самый record.key() в kafkaTemplate.send. Все события заказа 42 ложатся в одну партицию, а её внутри группы читает ровно один потребитель, поэтому он видит их в порядке записи. Без ключа продюсер раскладывает сообщения по партициям пачками: OrderCreated в третью, OrderConfirmed в седьмую, и второе может обработаться раньше первого. Сам relay порядок не ломает: он отправляет по одной строке и ждёт ответа.
Число партиций менять больно. Хеш по модулю 12 и по модулю 16 — разные партиции для одного ключа: старые события заказа 42 лежат в третьей, новые пойдут в седьмую, и пока третья не дочитана, порядок между ними нарушен. Расширяют топик при нулевом отставании потребителей — или мирятся с тем, что несколько минут проверка версии будет откладывать «обогнавшие» события.
Дубли неизбежны: как их отбить
Симптом: счётчик «заказов у покупателя» в проекции иногда больше, чем заказов на самом деле. Потребитель обработал событие, а подтверждение смещения до брокера не дошло — под перезапустили, группу перебалансировали, — и брокер отдал то же событие снова. Это не сбой, а нормальный режим: Kafka обещает доставку хотя бы один раз, и второй дубль рождается именно здесь.
У Kafka есть транзакции и режим «ровно один раз», и соблазн включить его понятен. Не выйдет: эта гарантия работает только на пути Kafka → Kafka — прочитали из топика, записали в другой топик, сдвинули смещение, всё одной транзакцией брокера. Как только обработчик пишет наружу, в PostgreSQL или Elasticsearch, брокер об этой записи не знает и откатить её не может. Идемпотентность приёмника нужна всё равно.
Сначала посмотрите, что именно повтор ломает. save строки по ключу заказа второй раз положит то же самое — повтор безвреден. increment счётчика прибавит единицу второй раз — вот он и накручивает. Операции, которые пишут состояние целиком, идемпотентны сами; накопительные — счётчики, суммы, письма — нет, и защищать нужно их.
Оба способа пропускают повтор, но узнают его по-разному: версия сравнивает событие со строкой проекции и не требует ничего помнить, таблица обработанных запоминает сам факт обработки и потому переживает повтор даже там, где операция накопительная — счётчик, сумма, отправка письма.
Проверка версии агрегата
Если в каждом событии есть aggregateVersion, а в проекции — колонка version, UPDATE проходит только когда событие новее строки:
UPDATE order_summary
SET status = 'CONFIRMED', confirmed_at = ?, version = ?, updated_at = NOW()
WHERE order_id = ? AND version < ?
@KafkaListener(topics = "order.events", groupId = "order-summary-projector")
@Transactional
public void onOrderConfirmed(@Payload OrderConfirmed event) {
int updated = orderSummaryRepository.updateStatusIfNewer(
event.orderId(),
"CONFIRMED",
event.confirmedAt(),
event.aggregateVersion()
);
if (updated == 0) {
if (orderSummaryRepository.exists(event.orderId())) {
log.debug("дубль или устаревшее событие, пропускаем: {}", event.eventId());
} else {
throw new ProjectionRowMissingException(event.orderId());
}
}
}
Ноль обновлённых строк означает три разные вещи, и валить их в один debug нельзя: событие уже обрабатывали; событие старее того, что применено; строки в проекции вообще нет. Первые два — норма, их пропускают молча. Третий — беда: OrderCreated ещё не доехал, и если промолчать, статус потеряется навсегда. Поэтому такой случай отделяют и бросают исключение — Kafka повторит доставку позже, когда строка уже появится. Как не повторять вечно — в следующем разделе.
Таблица обработанных событий
Второй способ — запоминать сам факт обработки. Перед обновлением потребитель проверяет, не видел ли он этот event_id, и в той же транзакции, что и обновление проекции, оставляет отметку:
CREATE TABLE processed_event (
event_id UUID NOT NULL,
consumer TEXT NOT NULL,
processed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
PRIMARY KEY (consumer, event_id)
);
В коде это одна строка в начале обработчика: if (!processedEvents.markIfAbsent(event.eventId(), "order-summary-projector")) return; — внутри INSERT ... ON CONFLICT DO NOTHING и проверка, вставилась ли строка. Отметка и обновление проекции — одной транзакцией, иначе между ними щель: отметили, упали, проекцию не обновили, а повтор уже отбит как «обработанное».
Что брать. Проверку версии — когда проекция хранит строку на агрегат и событие несёт его версию: ничего не нужно помнить, таблиц не прибавляется. Подвох один: если события несут не всё состояние, а только изменение, обогнавшее событие оставит проекцию без данных обогнанного — потому ключ партиции обязателен. Таблицу обработанных — когда операция накопительная (счётчик, сумма, письмо) или строка проекции собирается из событий разных агрегатов и «версии строки» нет. Цена таблицы — лишняя запись на каждое событие и уборка, о которой ниже.
Ядовитое сообщение: одно событие останавливает всех
Симптом: отставание группы растёт, проекция не обновляется ни по одному заказу, а в логе одно и то же событие падает раз за разом. Потребитель бросил исключение — сломанный JSON, ошибка в обработчике, тот самый «нет строки проекции». Партицию читает один поток, и пока событие не обработано, он стоит: за одним заказом встают все заказы той же партиции.
Что делает Spring Kafka без настройки: DefaultErrorHandler повторяет доставку десять раз без паузы, потом пишет ошибку в лог, фиксирует смещение и идёт дальше. Для «строки ещё нет» десять попыток за миллисекунды бесполезны — OrderCreated не доедет. А пропуск после них — тихая потеря: проекция навсегда без этого статуса.
Правильный обработчик повторяет с растущей паузой, чтобы соседнее событие успело доехать, а после исчерпания попыток не пропускает, а откладывает событие в отдельный топик и поднимает тревогу:
@Bean
DefaultErrorHandler kafkaErrorHandler(KafkaTemplate<Object, Object> template) {
var backOff = new ExponentialBackOffWithMaxRetries(5);
backOff.setInitialInterval(500);
backOff.setMultiplier(2.0);
backOff.setMaxInterval(8_000);
return new DefaultErrorHandler(new DeadLetterPublishingRecoverer(template), backOff);
}
Пять попыток с паузами 0,5 → 1 → 2 → 4 → 8 секунд — около шестнадцати секунд, и сумма важна: пауза идёт в потоке потребителя, и если она превысит max.poll.interval.ms (по умолчанию пять минут), брокер сочтёт потребителя мёртвым, перебалансирует группу и отдаст событие заново — получится цикл. Сломанный JSON — отдельная история: он падает ещё в poll, до обработчика, и без ErrorHandlingDeserializer в настройках потребителя контейнер будет падать на нём бесконечно, а обработчик ошибок его не увидит. С ним ошибка разбора доходит до DefaultErrorHandler, тот считает её неповторяемой и с первой попытки отправляет в order.events.DLT.
Событие в DLT — задача для человека: причину чинят, событие переигрывают. Тонкость: пока событие заказа 42 лежит в DLT, следующие события того же заказа тоже упадут на «нет строки» и уедут туда же. Проще после починки не переигрывать их по одному, а перестроить строку заказа 42 из write-store — тем же кодом, что и полное перестроение.
Окно отставания: откуда берутся 100 мс и секунда
Между коммитом и появлением изменения в проекции проходит время, и оно складывается из конкретных задержек.
Сумма худших нормальных задержек — около восьмисот миллисекунд, обычная — сто: relay в среднем ждёт половину тика, а потребитель забирает пачку сразу, как только она есть. Запрос клиента, сделанный сразу после ответа 200, попадает в это окно и видит старую строку — не ошибка, а свойство.
Relay тикает раз в 200 мс, значит в среднем строка ждёт сто, в худшем — двести. Брокер подтверждает запись за единицы миллисекунд. Потребитель спит в poll не дольше fetch.max.wait.ms — 500 мс по умолчанию, — но просыпается, как только в партиции что-то появилось, так что обычно и тут миллисекунды. UPDATE одной строки по первичному ключу — ещё несколько. Итого в норме около ста миллисекунд, в худшем нормальном случае — под секунду. При остановленном потребителе — часы: верхней границы у окна нет.
Клиент должен знать об этом заранее, иначе получаются ошибки, которые не воспроизвести. Ресурс отдаёт проекцию — это написано в контракте:
/orders/{id}/summary:
get:
summary: Сводка заказа из проекции
description: |
Обновляется событиями с задержкой, обычно до секунды.
Сразу после изменения заказа может вернуть прежнее состояние или 404.
Когда пользователь должен увидеть свою запись сразу
Создал заказ → открыл список → заказа нет. Для пользователя это «не сохранилось», и он нажимает ещё раз. Четыре приёма, от дешёвого к дорогому.
Вернуть результат из команды. Обработчик уже держит агрегат — пусть команда вернёт его данные, и экран покажет их, не спрашивая проекцию. Клиент добавляет заказ в список сам и заменяет при следующей загрузке. Это закрывает большинство случаев и не стоит ничего.
Два ресурса с разными гарантиями. GET /orders/{id} читает write-store и всегда свеж; GET /orders/{id}/summary читает проекцию и может отставать. Клиент выбирает сам, контракт говорит, что где.
Ответить 202. Если операция по-настоящему асинхронная — заказ проходит проверку, платёж ждёт банк, — честно ответить «принято» с Location, где смотреть состояние, а не делать вид, что всё уже случилось.
Дождаться версии проекции. Команда возвращает version агрегата, запрос принимает её параметром и ждёт, пока строка проекции её не достигнет:
public OrderSummary handle(GetOrderSummaryQuery query) {
Instant deadline = clock.instant().plus(Duration.ofSeconds(2));
while (true) {
Optional<OrderSummary> summary = orderSummaryRepository.findById(query.orderId());
boolean fresh = summary.map(s -> s.version() >= query.minVersion()).orElse(false);
if (fresh || !clock.instant().isBefore(deadline)) {
return summary.orElseThrow(() -> new OrderSummaryNotFoundException(query.orderId()));
}
sleeper.sleep(Duration.ofMillis(50));
}
}
Два решения здесь приняты явно. Без @Transactional вокруг цикла: транзакция на две секунды держит соединение из пула, и сорок таких запросов исчерпают пул. И по таймауту возвращаем то, что есть, с version в теле: запись уже прошла, а ошибка на чтении заставит клиента повторить команду и создать второй заказ. Клиент сравнивает версию с той, что вернула команда, и показывает «обновляется». Ошибку вместо старых данных отдают только там, где показать старое хуже, чем ничего, — остаток на счёте. Приём самый дорогой, до сорока чтений на запрос при шаге 50 мс, и берут его там, где свежесть обязательна.
Два приёма, которые не работают. Ждать после коммита внутри обработчика команды — POST держится до двух секунд, а решение по таймауту всё равно не принято. Липкая сессия к одному поду — проекция общая для всех подов, и пока она не догнала, свежести не даст ни один.
Пустая проекция: перестроение из write-store
Добавили новую проекцию — сводку заказов в Elasticsearch — или потеряли read-store. Ждать событий за 30 дней из Kafka нельзя: топик столько не хранит. Проекцию заполняют из write-store напрямую, пачками:
@Component
@RequiredArgsConstructor
public class OrderSummaryBootstrap {
private final OrderSummaryRepository orderSummaryRepository;
private final OrderRepository orderRepository;
private final AdvisoryLock advisoryLock;
@Async
@EventListener(ApplicationReadyEvent.class)
public void bootstrapIfEmpty() {
if (orderSummaryRepository.isEmpty()) {
advisoryLock.tryRun(REBUILD_LOCK_KEY, this::rebuildAll);
}
}
void rebuildAll() {
long lastId = 0;
while (true) {
List<Order> batch = orderRepository.findAllAfter(lastId, 1000);
if (batch.isEmpty()) break;
orderSummaryRepository.upsertBatchIfNewer(batch.stream().map(this::toSummary).toList());
lastId = batch.getLast().id().value();
}
}
}
Три вещи в этом коде стоят не случайно.
Advisory lock. Экземпляров сервиса несколько, и пустую проекцию увидят все сразу. Без блокировки перестраивать начнут все разом, и база ложится ровно в момент, когда сервис поднимается. pg_try_advisory_lock пускает одного, остальные уходят молча.
Перестроение после старта, а не на старте. Из ApplicationRunner приложение не завершит запуск, пока не пройдут все пачки: проба готовности не поднимется, оркестратор решит, что под не взлетел, и прибьёт его на середине. Поэтому перестроение вешают на ApplicationReadyEvent и уводят в отдельный поток: сервис уже принимает трафик, проекция догоняет параллельно.
Запись с проверкой версии. Пока перестроение идёт, потребитель событий работает. Обычный upsert затрёт свежую строку старым снимком, поэтому строки пишутся с тем же условием version < ?, что и в потребителе. Тот же метод с одним заказом — точечное перестроение после разбора DLT или сверки.
Таблицы растут: уборка outbox и processed_event
Обе таблицы только принимают строки, и никто из показанного кода их не удаляет. outbox получает строку на каждое изменение агрегата: при пятидесяти изменениях в секунду — четыре с лишним миллиона строк в сутки, за месяц — сто тридцать миллионов, из которых relay нужны только строки с published_at IS NULL. processed_event — строка на каждое событие у каждого потребителя.
Две вещи защищают relay от растущей таблицы. Первая — частичный индекс: CREATE INDEX ix_outbox_pending ON outbox (id) WHERE published_at IS NULL. В нём лежат только неопубликованные строки — обычно единицы, — и выборка relay не зависит от того, сколько миллионов уже опубликовано; без него запрос идёт по таблице целиком и с каждым днём медленнее. Вторая — уборка: фоновая задача удаляет опубликованные строки старше трёх суток пачками по десять тысяч с паузой между пачками. Один DELETE на сто миллионов строк — часовая транзакция, раздутый WAL и блокировки. Удалять сразу после публикации тоже можно, но тогда нечем посмотреть, что и когда ушло, — а при разборе рассинхрона это первое, куда смотрят.
Удалённые строки место сразу не освобождают: PostgreSQL помечает их мёртвыми, autovacuum отдаёт место под новые строки, но файл таблицы не сжимается. Пока пишут в конец и удаляют с начала, это нормальный режим; если же уборку запустили спустя месяц на раздутой таблице, размер останется прежним до VACUUM FULL с полной блокировкой или pg_repack. Радикальнее — секционировать по дню: PARTITION BY RANGE (created_at), и уборка превращается в DROP TABLE секции — мгновенно и без мёртвых строк.
Для processed_event срок хранения — не «сколько не жалко», а дольше, чем событие может прийти повторно. Повтор приходит только из топика, а топик хранит сообщения ограниченный срок — retention.ms, по умолчанию семь дней. Отметка старше него не нужна: события, которое она отбивает, уже нет. Держат срок топика плюс запас — десять дней при семи.
Как заметить, что синхронизация сломалась
Все режимы отказа выше объединяет одно: снаружи ничего не падает. Запись проходит, чтение отвечает — старым. Заметить это можно по трём цифрам.
Возраст самой старой неопубликованной строки в outbox: SELECT now() - min(created_at) FROM outbox WHERE published_at IS NULL. Глубина очереди хуже: при пике записи она большая и при здоровом relay. Возраст же в норме не превышает тика — 200 мс; минута означает, что relay триста тиков подряд не смог отправить: брокер недоступен или relay не крутится ни в одном поде. Тревога от минуты; метрику отдают через Micrometer Gauge этим же запросом.
Отставание группы потребителей — сколько сообщений записано в партицию, но не подтверждено группой. Его отдаёт kafka-exporter, и смотреть надо не на число, а на направление: сто сообщений при ста в секунду — секунда отставания, норма; сто сообщений, которые не убывают пять минут, — потребитель стоит. Сообщение в order.events.DLT — тревога сразу: за ним уже копятся события того же заказа.
Сверка. Обе метрики молчат, когда relay пометил строку опубликованной, а событие не ушло, или потребитель тихо пропустил событие. Ловит только сверка с write-store: раз в несколько минут взять заказы, изменённые за последний час, и сравнить version со строкой проекции. Расхождение старше окна отставания — рассинхрон, и он же говорит, какие строки перестроить.
Событие — не слепок таблицы: правила совместимости
Соблазнительно сделать событием сам record из jOOQ или сущность JPA — поля уже есть. Тогда каждая миграция write-таблицы становится изменением контракта для всех потребителей, включая чужие сервисы: переименовали колонку — у них сломался разбор JSON.
// так не надо: событие повторяет строку таблицы
public record OrderConfirmed(OrderRecord row) {}
// так надо: свои поля, независимые от схемы
public record OrderConfirmed(
UUID eventId,
Long orderId,
OrderStatus status,
Instant confirmedAt,
long aggregateVersion
) {}
Правила совместимости простые. Добавлять поле можно — старый потребитель его не знает и должен игнорировать; Jackson в Spring Boot так настроен по умолчанию, FAIL_ON_UNKNOWN_PROPERTIES выключен. У старых событий нового поля нет, значит у нового потребителя оно необязательное или со значением по умолчанию. Удалять поле, переименовывать и менять тип нельзя: в топике лежат события за семь дней в старом формате, и потребитель, который ждёт новый, на них упадёт. Когда без ломающего изменения не обойтись, рядом заводят OrderConfirmedV2 и публикуют оба, пока последний потребитель не перейдёт. При Avro и Schema Registry те же правила проверяет сам реестр.
Коротко
- Начинайте с outbox: событие пишется той же транзакцией, что и агрегат, и никогда не уходит в брокер из обработчика напрямую. Остальное — следствия.
- Relay помечает строку только после ответа брокера и берёт строки под
FOR UPDATE SKIP LOCKED. Первый дубль рождается здесь, и это правильно. - Ключ сообщения —
aggregate_id, иначе порядка нет. Партиции расширяют при нулевом отставании. - Дубли отбивают проверкой версии, если проекция хранит строку на агрегат, и таблицей обработанных, если операция накопительная. «Ровно один раз» в Kafka запись в вашу базу не защищает.
- Ошибка обработки — повтор с паузой, потом DLT и тревога, а не пропуск. Сумма пауз меньше
max.poll.interval.ms. - Окно отставания в норме 100 мс — 1 с и без верхней границы. Пишите это в контракте; свежесть после записи дешевле всего дать ответом самой команды.
- Уборку
outboxиprocessed_eventи частичный индекс поpublished_at IS NULLзаводят вместе с таблицами. - Мониторьте возраст старейшей неопубликованной строки, отставание группы и DLT; раз в несколько минут сверяйте версии с write-store.
Что почитать дальше
- Read-model в CQRS — таблица
order_summaryс колонкойversion, которую здесь обновляет потребитель. - Command side в CQRS — где в обработчике команды появляется запись в outbox.
- Распределённые паттерны — те же outbox и идемпотентный приёмник между сервисами, плюс сага.
- Event Sourcing — когда поток событий становится не транспортом, а хранилищем.