Всем привет, меня зовут Сергей Прощаев, и в этой статье расскажу про иллюзию, которая живёт почти в каждой событийной интеграции: если у всех событий заказа один ключ order_id, значит, они придут по порядку. Я Tech Lead и руководитель направления Java | Kotlin разработки в FinTech & E‑commerce и преподаю на курсах разработки и архитектуры в ОТУС.

Ниже — три инцидента за одну неделю в одном интернет‑магазине. Интеграция там построена на transactional outbox и Kafka, все основные паттерны применены, и на ревью схема выглядела корректной. Ваша задача — для каждого инцидента найти звено, на котором он возник. Подсказка: один инцидент может иметь несколько независимых причин.

Рис. 1. Один ключ на все события заказа — ещё не порядок
Рис. 1. Один ключ на все события заказа — ещё не порядок

Условие: схема и три инцидента

Сервис заказов на Kotlin и Spring Boot, база PostgreSQL. Доменное правило простое: собирать можно только оплаченный заказ, поэтому допустимая история заказа — сначала OrderPaid, потом OrderPacked. Схема такая.

  1. Запись событий. Изменение данных и строка в таблице outbox пишутся в одной транзакции. Оплату фиксирует обработчик колбэка платёжного шлюза (таблица payments), сборку — обработчик сообщений склада (таблица shipments). Так исторически сложилась интеграция: шлюз отдельно уведомляет склад, что заказ можно собирать, а наш сервис отдельно фиксирует оплату, и эти два сигнала приходят параллельно. Оба обработчика кладут событие в outbox с aggregate_id = order_id, а строку самого заказа в orders не трогают. Идентификатор строки outbox — bigserial.

  2. Relay. Три пода. Каждый раз в 200 мс выбирает пачку неотправленных событий запросом ORDER BY id LIMIT 100 FOR UPDATE SKIP LOCKED, отправляет их в Kafka асинхронно и после получения подтверждений от брокера помечает строки отправленными в той же транзакции.

  3. Kafka 4.2. Топик order-events, 12 партиций, ключ записи — order_id.

  4. Витрина заказов. Consumer group. Записи после poll() передаются в пул из 8 потоков. Поток, закончив свою запись, сообщает её offset основному потоку потребителя, и тот коммитит последний полученный offset + 1. Статус заказа витрина перезаписывает значением из последнего пришедшего события.

  5. Уведомления клиентам. Весной их перевели на share group: несколько экземпляров сервиса могут совместно читать одну партицию, и так проще масштабироваться, не изменяя топик.

Конфигурация продюсера в relay переехала из шаблона сервиса, которому лет семь:

acks=1
enable.idempotence=false
retries=10
max.in.flight.requests.per.connection=5
linger.ms=5

И требование из спецификации, которое согласовали все участники:

НФТ-07. События одного заказа доставляются потребителям в порядке их возникновения. Обеспечивается ключом партиционирования order_id.

Целиком схема показана на рисунке 2. Я специально нарисовал её так, как обычно рисуют в Confluence, без подсказок.

Рис. 2. Схема интеграции из условия задачи
Рис. 2. Схема интеграции из условия задачи

Каждый элемент схемы по отдельности выглядит корректно. Именно поэтому она и прошла ревью: ни на одном блоке не написано, какую гарантию он даёт, а какую нет.

Инциденты недели:

  • INC-1. Заказ 42: клиенту пришло «Ваш заказ собран», а через две секунды — «Оплата получена».

  • INC-2. Заказ 42: витрина показывает «Оплачен», хотя склад заказ уже собрал.

  • INC-3. Заказ 57: после планового перезапуска витрины заказ навсегда остался в статусе «Оплачен». Событие о сборке в топике есть, но витрина его так и не применила.

Дежурный уже снял содержимое партиции с событиями заказа 42:

$ kcat -C -b kafka:9092 -t order-events -p 7 -o 5119 -c 3 \
    -f 'offset=%o key=%k value=%s\n'
offset=5119 key=41 value={"type":"OrderCreated","outboxId":98}
offset=5120 key=42 value={"type":"OrderPacked","outboxId":101}
offset=5121 key=42 value={"type":"OrderPaid","outboxId":100}

Попробуйте решить

Для каждого инцидента определите, где он возник: до Kafka, у потребителя или и там, и там. Какие звенья схемы могли его вызвать?

Инцидент

Возник до Kafka?

Возник у потребителя?

Какие звенья?

INC-1

INC-2

INC-3

И вопрос, ради которого эту задачу дают аналитику: как переписать НФТ-07, чтобы ни один из трёх инцидентов не прошёл незамеченным?

Запишите ответ, прежде чем листать дальше.

Ответы, которые выглядят убедительно

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

«Оставим один под relay, и гонки не будет». Исчезнет только один источник конкуренции — параллельная работа экземпляров relay, и то не полностью. При стратегии RollingUpdate, которую Kubernetes использует для Deployment по умолчанию, старые и новые поды могут какое‑то время работать одновременно. А главное, даже внутри одного пода продюсер с выключенной идемпотентностью держит пять запросов в полёте: первый может упасть и уйти на повтор, а второй будет принят брокером раньше. Если же сделать отправку строго последовательной, при 15–20 мс на запись это примерно 50–67 событий в секунду на весь поток.

«Включим exactly‑once, и всё починится». Звучит убедительно, потому что exactly‑once воспринимается как самая сильная гарантия из возможных. Но у неё есть граница: она действует для записи и чтения внутри Kafka и не объединяет PostgreSQL и Kafka в одну транзакцию. А перестановку событий, которые отправили разные процессы, она не исправляет. Повтор и перестановка — разные проблемы с разными решениями.

Разбор по инцидентам

INC-1: причины есть и до Kafka, и после. Дамп — первый шаг, и он говорит, что событие переставилось ещё до брокера. Кандидатов три, и в этой схеме возможны все.

  • Источник. Склад и наш обработчик оплаты получают сигнал от шлюза параллельно, а общей блокировки на строке заказа у них нет. Транзакция сборки вполне может закоммититься раньше транзакции оплаты. Вдобавок id из bigserial выдаётся при INSERT, а строка становится видимой только после COMMIT: оплата могла получить id = 100, задержаться на записи чека и стать видимой позже сборки с id = 101. Номер строки говорит, кто раньше начал вставку, но не кто раньше завершил транзакцию.

  • Публикация. Три пода со SKIP LOCKED разбирают события одного заказа в разные пачки, и раньше в партиции окажется запись того пода, который первым получит подтверждение от брокера. Отдельно стоит помнить: блокировка FOR UPDATE SKIP LOCKED живёт только до конца транзакции. Либо relay держит её на время сетевой отправки в Kafka, и транзакции становятся долгими, либо отпускает раньше и отдаёт строки соседнему поду. Аренду записи она не заменяет.

  • Продюсер. В конфиге смешаны три разных свойства. acks отвечает за сохранность: при acks=all лидер ждёт подтверждения от реплик из текущего набора in‑sync, а запись принимается, только если в этом наборе не меньше min.insync.replicas реплик. Идемпотентность отвечает за повторы: для идемпотентного продюсера Kafka сохраняет порядок записей при max.in.flight.requests.per.connection не больше пяти, в том числе при повторной отправке. У нас она выключена. К тому же номера записей привязаны к идентичности конкретного продюсера. Разные поды — это разные продюсеры, и порядка между ними Kafka не обещает. transactional.id позволяет узнать продюсер после перезапуска и отсечь его старый экземпляр, но общей очерёдности для разных продюсеров не создаёт.

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

Но даже если починить всё до Kafka, INC-1 может повториться. Уведомления читают топик через share group. Эта модель стала production‑ready в Kafka 4.2, и документация рекомендует её там, где записи обрабатываются по одной, а не как упорядоченный поток. Если порядок для сервиса обязателен, а для уведомлений о смене статусов это так, share group ему не подходит. Один инцидент — две независимые причины. Поэтому расследование нельзя закрывать по первому найденному виновнику.

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

INC-3: это не перестановка, а потеря. Поток B быстро обработал запись со смещением 101 и сообщил об этом, основной поток закоммитил позицию 102. А поток A в это время ещё обрабатывал запись 100 — сборку заказа 57. Витрина перезапустилась и продолжила чтение с позиции 102. Запись 100 уже никогда не будет обработана.

Сама идея коммитить из основного потока правильная: KafkaConsumer не потокобезопасен, поэтому poll() и коммит должны оставаться в одном потоке. Ошибка в правиле «последний полученный offset + 1». Если завершены записи 101 и 102, а 100 ещё в работе, это правило закоммитит 103, хотя правильная позиция — 100. Kafka хранит для группы одну позицию на партицию, а не состояние каждой записи, поэтому коммитить можно только до конца непрерывного участка завершённых записей.

В нашем стеке продвижение позиции уже решено на уровне фреймворка. В Spring Kafka с версии 2.8 есть свойство контейнера asyncAcks: записи можно подтверждать в любом порядке, контейнер откладывает коммит до получения недостающих подтверждений и приостанавливает потребитель, пока не закоммичены все offset предыдущего poll(). Но asyncAcks решает только задачу прогресса чтения. Сериализацию обработки событий одного order_id он не обеспечивает, а после сбоя возможна повторная доставка уже обработанных записей, так что идемпотентную обработку он тоже не заменяет. Практическая сноска: в Spring Kafka 3.3.13 была регрессия asyncAcks, из‑за которой потребители навсегда оставались на паузе; исправление вышло в 3.3.14 и 4.0.4.

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

Для сериализации по ключу есть отдельные инструменты. В Confluent Parallel Consumer на момент написания статьи три режима: по ключу, по партиции и без порядка. В режиме по ключу записи одного ключа обрабатываются последовательно, а разные ключи — параллельно, и по умолчанию выбран именно он. Наш пул — это режим «без порядка», выбранный неявно, да ещё и с небезопасным коммитом.

Итог в таблице:

Инцидент

До Kafka

У потребителя

INC-1

Да: источник, публикация, продюсер

Да: share group

INC-2

Да: вход уже переставлен

Да: нет проверки номера события

INC-3

Нет

Да: коммит offset через незавершённую запись

Главная мысль: ключ даёт размещение, номер — позицию, правило — допустимость

order_id определяет, где событие хранится. Номер события определяет его позицию в истории заказа. Какая история вообще допустима, определяет бизнес‑правило. Ни ключ, ни номер сами по себе причинность не создают: если сборка закоммитится раньше оплаты, номер честно зафиксирует «сборка — 7, оплата — 8», и история будет упорядоченной, но недопустимой.

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

Звено

Что задаёт очерёдность

Что её ломает в нашей схеме

Что добавляем

Бизнес‑правило

Какие переходы заказа допустимы

Правило нигде не проверяется при записи

Проверку статуса в транзакции

Транзакция и outbox

Номер события в сериализованной истории

Нет общей блокировки, bigserial отражает порядок вставки

Блокировку строки заказа и aggregate_version

Публикация

Порядок чтения событий из базы

Три пода relay со SKIP LOCKED

CDC вместо самописного relay

Продюсер

Идентичность продюсера и номера его записей

Идемпотентность выключена, пять запросов в полёте

Идемпотентный продюсер

Партиция

Порядок поступления записей одного ключа

Изменение числа партиций

Миграцию через новый топик

Обработка у потребителя

Порядок побочных эффектов

Share group, пул без сериализации по ключу

Обработку по ключу, проверку номера, отсев по event_id

Прогресс чтения

Коммит offset

Коммит через незавершённую запись

Коммит только непрерывного участка

Чтобы в следующий раз быстро найти звено, на котором всё сломалось, удобно идти по дереву вопросов на рисунке 3. Под сериализацией публикации в нём понимается простое правило: событие с номером N+1 не публикуется впервые, пока не подтверждена публикация номера N.

Рис. 3. Дерево поиска звена, на котором нарушен порядок событий
Рис. 3. Дерево поиска звена, на котором нарушен порядок событий

Главное в этом дереве — первый вопрос. Дамп партиции делит расследование пополам: либо проблема выше Kafka, либо у потребителя. А второй вопрос показывает, почему без номера события расследование почти всегда упирается в тупик.

Похожий вывод можно найти и на другом стеке. Мне как‑то попался разбор системного архитектора Василия Миловидова, опубликованный на Хабре 19 августа 2026 года (ссылка — в конце статьи). Его команда публиковала события из PostgreSQL в Kafka через очередь задач River и обнаружила, что порядок событий одного агрегата теряется ещё до Kafka. Показательная деталь оттуда: служебный процесс River может повторно запустить задачу, которая выполняется дольше часа, параллельно с исходной, и тогда старая версия события окажется в партиции после новой. Стек другой, а суть та же: каждый компонент выполнил свой контракт, но контракта «очерёдность для заказа» не было ни у кого.

Что меняем

Шаг 1. Номер события и бизнес‑правило в источнике. Вводим aggregate_version — номер доменного события в потоке заказа, а не версию состояния строки. Номер растёт только вместе с записью события, а изменения заказа, которые событий не порождают, счётчик не трогают. Чтобы не усложнять пример, принимаем ограничение модели: одна транзакция порождает не более одного доменного события.

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

Блокировка сериализует изменения строки, а условие status = 'PAID' проверяет бизнес‑правило в рамках этой сериализации: сборка проходит, только если заказ уже оплачен.

CREATE TABLE outbox (
    id                bigint GENERATED BY DEFAULT AS IDENTITY PRIMARY KEY,
    event_id          uuid        NOT NULL UNIQUE,
    aggregate_type    text        NOT NULL DEFAULT 'order',
    aggregate_id      bigint      NOT NULL,
    aggregate_version bigint      NOT NULL,
    event_type        text        NOT NULL,
    payload           jsonb       NOT NULL,
    created_at        timestamptz NOT NULL DEFAULT now(),
    CONSTRAINT outbox_aggregate_version_uq
        UNIQUE (aggregate_id, aggregate_version)
);

-- Обработчик оплаты
BEGIN;
UPDATE orders
   SET aggregate_version = aggregate_version + 1,
       status            = 'PAID'
 WHERE id = 42
RETURNING aggregate_version;                    -- например, 7

INSERT INTO payments (order_id, amount, currency, status)
VALUES (42, 4990.00, 'RUB', 'CAPTURED');

INSERT INTO outbox (event_id, aggregate_id, aggregate_version, event_type, payload)
VALUES (gen_random_uuid(), 42, 7, 'OrderPaid',
        jsonb_build_object('orderId', 42, 'status', 'PAID',
                           'amount', 4990.00, 'currency', 'RUB'));
COMMIT;

-- Обработчик сборки. Если транзакция оплаты ещё не закоммичена,
-- UPDATE дождётся её и перепроверит условие.
BEGIN;
UPDATE orders
   SET aggregate_version = aggregate_version + 1,
       status            = 'PACKED'
 WHERE id = 42
   AND status = 'PAID'
RETURNING aggregate_version;   -- 8; если строк нет, заказ ещё не оплачен
-- ... INSERT в shipments и outbox с номером 8
COMMIT;

Если заказ ещё не оплачен, UPDATE просто не найдёт строку, и сам по себе он ничего не откладывает. Поэтому команда сборки не должна теряться: обработчик склада сохраняет её в собственную таблицу команд и повторяет после OrderPaid. Это отдельное требование к обработчику, и его стоит записать явно.

Шаг 2. CDC вместо самописного relay. Вариант, который я бы выбрал по умолчанию, если в компании уже есть Kafka Connect, — Debezium с Outbox Event Router. Логическое декодирование PostgreSQL передаёт успешно закоммиченные транзакции в порядке коммитов, и Debezium использует этот поток как источник публикации. Коннектор PostgreSQL использует одну задачу и последовательно обрабатывает свой поток изменений. Это упрощает сохранение порядка, но ограничивает горизонтальное масштабирование одного коннектора. Роутер использует aggregate_id как ключ записи и кладёт event_id в заголовок.

Конфигурация коннектора Debezium:

# Коннектор читает только таблицу outbox. Параметры подключения к БД опущены.
connector.class=io.debezium.connector.postgresql.PostgresConnector
plugin.name=pgoutput
slot.name=orders_outbox
table.include.list=public.outbox
heartbeat.interval.ms=10000

transforms=outbox
transforms.outbox.type=io.debezium.transforms.outbox.EventRouter
transforms.outbox.table.field.event.id=event_id
transforms.outbox.table.field.event.key=aggregate_id
transforms.outbox.table.field.event.payload=payload
# Поле маршрутизации указано явно, потому что имя по умолчанию (aggregatetype)
# не совпадает с нашей колонкой. В имени топика его значение пока не используется.
transforms.outbox.route.by.field=aggregate_type
transforms.outbox.route.topic.replacement=order-events
transforms.outbox.table.fields.additional.placement=aggregate_version:header:aggregateVersion,event_type:header:eventType

Конфигурация воркера Kafka Connect:

# Kafka Connect по умолчанию отключает идемпотентность своих продюсеров.
# Включается в конфигурации воркера, а не коннектора.
producer.enable.idempotence=true
producer.acks=all
producer.max.in.flight.requests.per.connection=5

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

$ kcat -C -b kafka:9092 -t order-events -p 7 -o 9340 -c 2 \
    -f 'offset=%o key=%k headers=%h\n'
offset=9340 key=42 headers=id=5d0c2a7e-8b1f-4c3e-9a55-1e2f3a4b5c6d,aggregateVersion=7,eventType=OrderPaid
offset=9341 key=42 headers=id=a81f6c2d-3e4b-4f5a-8c7d-9e0f1a2b3c4d,aggregateVersion=8,eventType=OrderPacked

При восстановлении после сбоя нужно учитывать возможность повторной доставки: в партиции может оказаться, например, 7, 8, 7, 8. Это повторы, а не перестановка, у вторых записей те же event_id. Потребители должны быть к ним устойчивы. Как сократить число повторов и как обойтись без Kafka Connect, — в разделе «Дополнительно».

Шаг 3. Потребители.

  • Все потребители распознают повторы по event_id из заголовка записи. Но распознать повтор — не то же самое, что исполнить побочный эффект ровно один раз. Отметка об обработке должна фиксироваться атомарно с локальным изменением состояния, как у витрины в одной транзакции с базой. Если же побочный эффект внешний, как push‑уведомление, между отправкой и сохранением отметки остаётся окно: процесс может упасть, и уведомление уйдёт повторно. Здесь нужен идемпотентный побочный эффект, например ключ идемпотентности у push‑провайдера, если он такое поддерживает.

  • Витрина применяет событие, только если номер больше уже применённого, и хранит полное состояние заказа.

  • Уведомления возвращаются в обычную consumer group. Назначение партиций остаётся за Kafka, а внутри каждого экземпляра обработка событий одного order_id сериализуется отдельно: сама consumer group такой гарантии не даёт.

  • Offset коммитится только до конца непрерывного участка обработанных записей, например через asyncAcks в Spring Kafka. При отзыве партиции собственный пул дожидается или отменяет незавершённые задачи.

  • Share group остаётся для независимых задач, например генерации чеков.

Шаг 4. Отравленное событие. Что делать, если событие с номером 9 не удаётся опубликовать никогда: остановить заказ, пропустить номер или пересоздать событие из текущего состояния? Здесь важен масштаб. Ошибка обработки записи по умолчанию может остановить задачу коннектора, а задача у Debezium PostgreSQL одна, поэтому может встать продвижение всего потока. Ограничение размера и проверку схемы события я бы делал ещё при записи в outbox.

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

НФТ-07: требование, а не механизм

Исходный НФТ-07 описывает механизм, причём неработающий, а не свойство системы. Я бы разделил его на две части: требование, которое переживёт смену базы или брокера, и реализационную гарантию текущей версии. И сразу ввёл бы словарь: порядок публикации — очерёдность событий в топике, порядок обработки — очерёдность, в которой потребитель их применяет, порядок побочных эффектов — очерёдность того, что видит внешний мир.

НФТ-07. Порядок событий заказа.

Требование:

  1. Допустимые переходы заказа задаются бизнес‑правилами (например, сборка только после оплаты) и проверяются при создании события. Команда, для которой правило пока не выполнено, не теряется, а повторяется после наступления нужного события.

  2. События одного заказа образуют последовательность номеров aggregate_version, начиная с 1 и без пропусков. Номер получают только доменные события; одна транзакция порождает не более одного события.

  3. Если изменение A в сериализованной истории заказа предшествует изменению B, событие A получает меньший номер.

  4. Каждое событие несёт уникальный event_id. Повторная доставка допустима, потребитель распознаёт повтор по event_id.

  5. Порядок публикации: для каждого order_id последовательность первых доставок уникальных event_id имеет строго возрастающие aggregate_version.

  6. Порядок побочных эффектов: сервис уведомлений не выполняет побочный эффект для номера N+1, пока не обработан номер N. Витрина применяет событие, только если его номер больше уже применённого. Повтор события не приводит к повторному побочному эффекту.

  7. Порядок обработки и прогресс чтения: обработка событий одного заказа сериализуется, а offset партиции фиксируется не выше непрерывной последовательности завершённых обработок.

  8. Политика для события, которое не удаётся опубликовать, согласована с владельцами всех потребителей. Каждый такой случай фиксируется метрикой и алертом.

Реализационная гарантия (может меняться без изменения требования): номер назначается в транзакции PostgreSQL под блокировкой строки заказа; Debezium переносит события в Kafka в порядке коммитов с ключом order_id; число партиций order-events меняется только через миграцию.

И критерии, по которым требование можно проверить:

# language: ru
Функция: Порядок событий заказа

  Сценарий: Сборка не создаёт событие раньше оплаты
    Дано заказ 42 не оплачен
    Когда обработчик склада получает команду на сборку заказа 42
    Тогда событие OrderPacked не создаётся
    И команда сохраняется и повторяется после фиксации оплаты

  Сценарий: Публикация сохраняет порядок номеров
    Дано для заказа 42 закоммичено событие с номером 7
    И после него в другой транзакции закоммичено событие с номером 8
    Когда сообщения этих событий появляются в топике order-events
    Тогда оба события находятся в одной партиции
    И первая доставка события с номером 7 имеет меньший offset, чем первая доставка события с номером 8

  Сценарий: Уведомление не опережает предыдущий переход
    Дано сервис уведомлений обработал событие заказа 42 с номером 6
    Когда он получает событие с номером 8 раньше события с номером 7
    Тогда уведомление по номеру 8 не отправляется
    И метрика gap_before_side_effect_total увеличивается на 1
    И если событие с номером 7 не пришло за 5 минут, дежурный получает алерт

  Сценарий: Витрина не откатывает состояние
    Дано витрина применила событие заказа 42 с номером 8
    Когда приходит событие с номером 7
    Тогда состояние заказа не меняется
    И метрика stale_event_dropped_total увеличивается на 1

  Сценарий: Перезапуск не теряет событие
    Дано витрина получила записи с offset 100 и 101
    И обработка записи 101 завершилась раньше, чем записи 100
    Когда витрина перезапускается до завершения обработки записи 100
    Тогда после перезапуска запись 100 обрабатывается

Дополнительно

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

А. Exactly‑once в Kafka Connect

Kafka Connect поддерживает exactly‑once для source‑коннекторов, и коннектор Debezium PostgreSQL умеет в этом участвовать. Режим требует распределённого Connect.

Конфигурация воркера Kafka Connect, на всех воркерах:

exactly.once.source.support=enabled

Конфигурация коннектора Debezium:

exactly.once.support=required
transaction.boundary=poll

Конфигурация потребителей топика order-events:

isolation.level=read_committed

Connect в этом режиме записывает события и свою позицию чтения в Kafka транзакциями, но транзакцию PostgreSQL с транзакцией Kafka не объединяет. Сам Debezium предупреждает, что корректность exactly‑once зависит от реализации транзакционного протокола Kafka и что известны пограничные случаи. Отсев повторов у потребителей я бы оставил.

Б. Слот репликации

Слот удерживает WAL, пока коннектор не подтвердит чтение. Если коннектор упал незаметно, диск заполняется, и база перестаёт принимать запись. Параметр max_slot_wal_keep_size ограничивает объём удерживаемого WAL. Но если нужные сегменты уже удалены, продолжить с прежней позиции нельзя, и источник придётся синхронизировать заново; конкретный сценарий зависит от режима снапшота и состояния слота. Объём удерживаемого WAL стоит вывести в мониторинг с алертом:

SELECT slot_name,
       active,
       pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained_wal
FROM pg_replication_slots;

В. Собственный relay вместо Debezium

Это упрощённая схема для случая, когда Kafka Connect нет. Публикацию нужно сериализовать по заказу, не держа транзакцию PostgreSQL открытой на время сетевого вызова. last_confirmed_version здесь — курсор публикации: номер, отправку которого брокер уже подтвердил. Relay публикует только N = last_confirmed_version + 1. Цикл такой: взять аренду, опубликовать N, дождаться подтверждения, сдвинуть курсор и повторять до конца очереди или до первой ошибки.

-- Курсор публикации по заказу; строка создаётся вместе с первым событием
CREATE TABLE order_publication (
    order_id               bigint PRIMARY KEY,
    last_confirmed_version bigint      NOT NULL DEFAULT 0,
    lease_owner            text,
    lease_until            timestamptz,
    fencing_token          bigint      NOT NULL DEFAULT 0
);

-- Короткая транзакция 1: взять аренду. Срок должен превышать нормальное время
-- публикации либо продлеваться во время работы; 30 секунд — только пример.
UPDATE order_publication
   SET lease_owner   = 'relay-2',
       lease_until   = now() + interval '30 seconds',
       fencing_token = fencing_token + 1
 WHERE order_id = 42
   AND (lease_until IS NULL OR lease_until < now())
RETURNING fencing_token, last_confirmed_version;   -- например, 17 и 6

-- Вне транзакции БД: опубликовать 7, дождаться подтверждения, затем 8 и так далее

-- Короткая транзакция 2: сдвинуть курсор, только если аренда всё ещё наша
UPDATE order_publication
   SET last_confirmed_version = 8,
       lease_until            = NULL
 WHERE order_id = 42
   AND fencing_token = 17;

Токен защищает от устаревшего владельца только сдвиг курсора в PostgreSQL. Уже начатую или возобновлённую отправку в Kafka он не отзывает: под, у которого истекла аренда, может проснуться и дописать в топик событие. Почему это не ломает порядок первых публикаций? Каждый владелец публикует номера строго подряд, начиная с подтверждённого курсора. Опоздавший владелец подчиняется тому же правилу: он продолжает строго со следующего номера после своего последнего подтверждённого. Значит, номер N+1 впервые появляется в топике только после того, как номер N там уже есть, кто бы из владельцев его ни отправил. Всё, что опоздавший под допишет с уже опубликованными номерами, — это повторы с теми же event_id. Например, 7, 8, 9, 8: последняя восьмёрка — повтор, а не перестановка.

Гарантия держится на трёх условиях: строгая последовательность номеров у каждого владельца, остановка на первой ошибке и сдвиг курсора только после подтверждения. Продюсер при этом — с acks=all, enable.idempotence=true и не больше чем пятью запросами в полёте. Если какое‑то условие нарушить, гарантия пропадает. Поэтому в production я бы всё‑таки предпочёл CDC.

Г. Контрольный потребитель

Чтобы видеть нарушения, а не узнавать о них из жалоб, полезен отдельный потребитель без побочных эффектов. Это инструмент мониторинга, а не механизм гарантии, и пример ниже намеренно упрощён.

Главное упрощение — хранилища seenEvents и highestSeenVersions обновляются отдельно. Если процесс упадёт между markSeen и put, при повторной доставке событие будет принято за повтор и в highestSeenVersions так и не попадёт: контрольный потребитель потеряет собственный сигнал. В реальной реализации оба состояния нужно обновлять атомарно, например в одной транзакции одного хранилища.

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

// Контрольный потребитель (упрощённо): сначала отсеивает повторы по event_id,
// затем сравнивает номер с наибольшим уже увиденным номером заказа.
// В реальной реализации markSeen и put должны выполняться атомарно.
class OrderSequenceChecker(
    private val seenEvents: ProcessedEventStore,        // долговечное хранилище event_id, TTL больше окна повтора
    private val highestSeenVersions: VersionStore,      // наибольший увиденный номер по заказу
    private val meters: MeterRegistry,
) {
    fun check(record: ConsumerRecord<String, String>) {
        val orderId = record.key()
        val eventId = record.headerString("id") ?: return               // event_id от Outbox Event Router
        val version = record.headerString("aggregateVersion")?.toLongOrNull() ?: return

        if (!seenEvents.markSeen(eventId)) {                              // false, если event_id уже встречался
            count("publication_duplicates_total")
            return
        }

        val highestSeen = highestSeenVersions.get(orderId)
        when {
            highestSeen == null -> highestSeenVersions.put(orderId, version)          // точка отсчёта
            version == highestSeen + 1 -> highestSeenVersions.put(orderId, version)
            version <= highestSeen -> count("publication_order_violations_total")   // впервые пришёл меньший номер
            else -> {
                count("publication_gap_signals_total")                  // сигнал, а не доказательство потери
                highestSeenVersions.put(orderId, version)
            }
        }
    }

    private fun count(name: String) = meters.counter(name).increment()
}

private fun ConsumerRecord<*, *>.headerString(name: String): String? =
    headers().lastHeader(name)?.value()?.toString(Charsets.UTF_8)

Ноль нарушений за неделю не доказывает гарантию. Он значит только, что за этот период нарушений не зафиксировано.

Д. Изменение числа партиций

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

  1. Остановить публикацию в старый топик.

  2. Убедиться, что outbox и CDC полностью опустошены.

  3. Зафиксировать последний offset старого топика отдельно для каждой партиции.

  4. Дождаться, пока каждая группа потребителей обработает каждую партицию до этого offset.

  5. Только после этого включить публикацию и чтение нового топика.

Вместо заключения: какой навык проверяла задача

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

Оценивать свой ответ стоит не по числу найденных звеньев, а по тому, удалось ли:

  • по дампу партиции разделить проблемы «до Kafka» и «у потребителя»;

  • заметить, что у INC-1 две независимые причины;

  • понять, что INC-3 — потеря события, а не перестановка;

  • отличить повтор от перестановки, а упорядоченную историю — от допустимой;

  • записать требование так, чтобы оно не зависело от конкретной базы и брокера.

Если какой‑то пункт ускользнул, это нормально. Такие ошибки не видны на схеме — их находят в проде.

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

Разбирать чужие инциденты по дампам — полезно. Но рано или поздно хочется, чтобы твоя схема не попала в такой разбор. В статье выше мы разобрали, где ломается порядок событий. Осталось разобраться с тем, что вокруг: как перестать тонуть в оформлении ТЗ и как договариваться с теми, кто тормозит внедрение. Об этом — два бесплатных вебинара:

  • 14 октября в 20:00. «ИИ‑агент как супероружие системного аналитика: составляем ТЗ за 10 минут». Записаться

  • 21 октября в 20:00. «Анатомия сопротивления: превращаем токсичных стейкхолдеров в союзников проекта». Записаться

Комментарии (1)


  1. nightweb
    02.10.2026 06:53

    Добавил бы к приёмочным тестам повтор уведомления об успешной оплате после перехода заказа в PACKED.
    В примере обработчик оплаты выставляет PAID без проверки текущего статуса. Если повтор не отсекается отдельно и приводит к созданию нового OrderPaid с новым event_id и большей aggregate_version, то проверки потребителя по ID и версии его не отфильтруют. Порядок доставки при этом может оставаться правильным.

    Вообще повтор подтверждения того же платежа не должен менять статус заказа, увеличивать версию или создавать новое доменное событие. Для этого можно использовать идентификатор операции у платёжного провайдера, который сохраняется при повторных уведомлениях об одной оплате, как ключ дедупликации с ограничением UNIQUE в БД. Регистрацию операции, изменение заказа и запись в outbox я бы выполнял в одной транзакции, а при обнаружении повтора завершал обработку без этих изменений.

    Поделитесь, какой ключ вы используете, чтобы повторное уведомление об одной оплате не создавало новое OrderPaid?