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

Пользователь подтвердил заказ, открыл список — а там заказ всё ещё «новый». Через секунду — уже «подтверждённый». Бывает хуже: в списке заказ «подтверждён», а транзакция на самом деле откатилась. Так снаружи выглядит главная трудность 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 — в той же базе и той же транзакции, что и изменение агрегата. База гарантирует: либо есть и новый статус, и событие, либо нет ни того, ни другого.

одна транзакция UPDATE order status → CONFIRMED, version = 7 INSERT INTO outbox OrderConfirmed, version = 7 либо обе строки, либо ни одной: событие не потеряется и не появится без изменения в базе после COMMIT общей транзакции нет: дальше только «хотя бы один раз» outbox-relay тик каждые 200 мс дубль №1: отправил, упал до markPublished — пошлёт снова Kafka топик order.events дубль №2: перебалансировка группы — событие придёт повторно потребитель @KafkaListener UPDATE … WHERE version < 7: второй раз не пройдёт order_summary status = CONFIRMED, version = 7

Атомарность есть только внутри пунктирной рамки: заказ и событие записываются одной транзакцией. Дальше событие едет по правилу «хотя бы один раз», и родиться дубль может в двух местах — в 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 счётчика прибавит единицу второй раз — вот он и накручивает. Операции, которые пишут состояние целиком, идемпотентны сами; накопительные — счётчики, суммы, письма — нет, и защищать нужно их.

по версии агрегата первая доставка событие v7 строка v6 6 < 7 → UPDATE строка v7 повтор событие v7 строка v7 7 < 7 → 0 строк без изменений по таблице обработанных первая доставка событие e1 e1 нет INSERT e1, UPDATE строка новая повтор событие e1 e1 есть return без изменений

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

Проверка версии агрегата

Если в каждом событии есть 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 мс и секунда

Между коммитом и появлением изменения в проекции проходит время, и оно складывается из конкретных задержек.

0 мс COMMIT: заказ и outbox записаны, клиенту 200 OK +1 мс клиент: GET /orders/42/summary — строка старая до 200 мс relay: очередной тик fixedDelay = 200 взял строку +5–50 мс брокер подтвердил запись, published_at заполнен +1–500 мс потребитель получил пачку из топика (poll) +1–20 мс UPDATE order_summary: теперь строка свежая

Сумма худших нормальных задержек — около восьмисот миллисекунд, обычная — сто: 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 — когда поток событий становится не транспортом, а хранилищем.