Симптом после восстановления consumer звучит просто: «нужно проиграть события ещё раз». Затем один event приходит повторно, второй consumer уже обновлён, а в storage появляется непонятный result. Если replay создаёт новый id, невозможно отличить повтор старого факта от нового факта. Если он сохраняет id, но consumer не ведёт ledger, можно записать второй effect. Цена — не только дубль. Команда теряет доказательство, по какому contract был получен каждый результат, и не может безопасно решить, что делать со старой схемой.
Полезная точка старта — не ручной запуск всех consumer, а evidence packet из пяти частей: envelope id и source, event type, input schemaVersion, consumer id/resultVersion, решение по ledger. В этой статье replay обозначает повторную delivery того же учебного source:id. Он не открывает реальный topic, не перемещает offset и не повторяет Kafka record. Все идентификаторы, даты, source, payload, результаты и правила ниже учебные. Fixture хранит состояние только в Map процесса Node: это не broker, не schema registry, не production event, не реализация CloudEvents или Kafka, не transport, не база и не измерение throughput. Поэтому result duplicate-or-replay-suppressed означает только, что Map уже видела тот же consumer contract и source:id.
Replay сохраняет логическую идентичность
Самая опасная «починка» — сделать новое event id, чтобы consumer не счёл запись duplicate. Так обходят проверку, но меняют вопрос. Новый id может означать новый факт, исправленную команду или технический replay; эти случаи нельзя сливать. Для controlled replay неизменным остаётся исходный id, source, type, subject и payload schema. Дополнительная причина replay может жить рядом с операционной записью, но не должна подменять исходный event. Тогда ledger способен ответить: этот contract уже обработал данный вход или нет.
Ключ ledger в fixture — consumerId:source:event.id. Source входит в identity: один и тот же id из другого producer не должен случайно подавить отдельный факт. Он подходит только для демонстрации одного projection result. В реальном домене этого может быть мало: внешний effect иногда нужно ключевать по business intent, а два разных consumer могут законно создать разные projections по одному event. Не надо переносить ключ как готовую идемпотентность. Сначала надо назвать effect и его owner, затем проверить, где хранится receipt вместе с результатом.
const ledger = new Map();
const first = consumeTrainingEvent(ledger, event, "orders-projection@1", "initial-delivery");
const duplicate = consumeTrainingEvent(ledger, event, "orders-projection@1", "duplicate-delivery");
const replay = consumeTrainingEvent(ledger, event, "orders-projection@1", "controlled-replay");
first.effectWritten; // true
duplicate.effectWritten; // false
replay.effectWritten; // false
ledger.size; // 1
| Delivery kind | Event identity | Ledger до шага | Решение | Effect |
|---|---|---|---|---|
initial-delivery | тот же source:id | нет ключа consumer | effect-recorded | один учебный result записан |
duplicate-delivery | тот же source:id | ключ уже есть | duplicate-or-replay-suppressed | новый result не пишется |
controlled-replay | тот же source:id | ключ уже есть | duplicate-or-replay-suppressed | replay не обходил ledger |
historical-replay в другом contract | тот же source:id | другой ledger consumer | effect-recorded с новой resultVersion | разный reader result допустим |
| delivery schema v3 | новый event id, неизвестная schema | неважно | contract-update-required | effect запрещён |
Duplicate и новая интерпретация — разные развилки
Duplicate по тому же contract не должен создавать второй result. Но новое правило consumer иногда действительно требует пересчитать projection. Тогда не стоит удалять старую ledger запись и притворяться, что история не существовала. В fixture orders-projection@2 имеет другой ledger key и resultVersion; historical replay v1 может создать свой result, потому что это другой declared reader contract. Такой выбор не делает результат «истиннее», он делает видимой новую интерпретацию.
Перед этим шагом нужно определить, разрешён ли второй projection в домене. Для email, платёжа или изменения внешней системы одного consumerId:source:event.id почти наверняка недостаточно: потребуется effect key на внешней границе, транзакция или ручное решение. Для локальной read projection иногда допустимо хранить две версии и переключать reader после сверки. Статья не выбирает между этими архитектурами. Она требует, чтобы автор replay написал, какой effect будет создан и почему второй contract имеет право его создавать.
const result = consumeTrainingEvent(
new Map(),
trainingEvents.orderStatusV1,
"orders-projection@2",
"historical-replay",
);
result.resultVersion; // "2021-06.orders-projection.2"
result.inputSchemaVersion; // 1
result.projection.paymentReference; // null
// Result version описывает consumer contract, не версию broker-а.
Схема не должна исчезать за duplicate
Проверка ledger не заменяет проверку schema. Если v3 event пришёл с незнакомым state, consumer не должен сначала посмотреть id, не найти его и записать effect только потому, что это первый delivery. В учебном алгоритме порядок другой: validate envelope, выбрать consumer contract, проверить accepted schema version, затем читать ledger. Это сохраняет важный факт: новый event был получен, но effect запрещён из-за договора, а не потерян как «неизвестная ошибка».
Также нельзя делать обратное: любой duplicate автоматически игнорировать до записи diagnostics. Result suppress содержит deliveryKind, ledgerKey, input schema и result version первого результата. Это не production audit trail, но минимальное evidence помогает отличить повтор сети от operator replay. Если реальная система не сохраняет эти факты, расследование начнётся с догадки: очередной запуск создал запись или просто повторил уже завершённую delivery.
const decision = projectForConsumer(trainingEvents.orderStatusV3, "orders-projection@2");
decision;
// {
// state: "contract-update-required",
// inputSchemaVersion: 3,
// effectAllowed: false,
// }
// Не угадываем, что state: "settled" эквивалентен status: "paid".
Fixture задаёт узкую, но полезную проверку
Фикстура использует два event v1/v2, один намеренно несовместимый v3, два consumer contract и две Map. Она проверяет envelope boundary, foreign source, additive v2 для v1 reader, normalisation old event в v2 reader, один записанный result, suppress duplicate, suppress controlled replay, отдельный resultVersion второго consumer и остановку v3. Assertions не измеряют время, не моделируют crash между внешним effect и receipt и не говорят ничего о exactly-once. Их задача скромнее: не дать редактуре незаметно поменять «тот же id подавляется» на «replay всегда пишет ещё раз».
const fixture = runEventIntegrationFixture();
if (!Object.values(fixture.assertions).every(Boolean)) {
throw new Error('event contract changed without an explicit decision');
}
console.log(fixture.deliveries.duplicateDelivery.state);
// 'duplicate-or-replay-suppressed'
Apache Kafka 2.7 producer API прямо рассматривает retries и возможность duplicate в некоторых режимах. Это хороший повод не обещать exactly-once одним словом. Конкретные guarantees зависят от producer, broker, consumer, storage и внешней границы. CloudEvents в историческом snapshot апреля 2021 года помогает разделить context и event data, но не задаёт business receipt. В статье оба источника ограничивают формулировку, а не дают готовую implementation.
| Наблюдение | Что собрать | Безопасное действие | Чего не делать |
|---|---|---|---|
тот же source:id, same consumer contract | ledger key и первый resultVersion | suppress duplicate, сохранить evidence delivery | создавать новый id ради обхода проверки |
| тот же event, новый declared consumer contract | старый и новый resultVersion, тип effect | отдельный controlled replay после явного решения | удалять старую receipt и терять историю |
| schemaVersion не объявлена consumer | envelope, input schema, причина отказа | contract update или manual route без effect | угадывать поле по похожему имени |
| source или type неожиданны | validation reason до payload | rejected envelope и проверка producer boundary | выполнять partial effect для «похожего» входа |
| внешний effect уже мог быть сделан | idempotency contract внешней системы и receipt | остановить автоматический replay до доказательства | считать Map заменой транзакции |
Маршрут controlled replay
- Зафиксировать исходный event id, source, type, subject, payload schema и причину replay. Не заменять их новым JSON.
- Проверить envelope и consumer contract до ledger. Unknown type или schema должен дать объяснимый отказ без effect.
- Выбрать key, который соответствует именно данному consumer result; отдельно обсудить key доменного или внешнего effect.
- Если ledger уже содержит same consumer contract + source:id, suppress replay и сохранить delivery evidence.
- Если нужен новый расчёт, ввести новый declared consumer contract/resultVersion. Не перезаписывать старую интерпретацию без следа.
- Для historical event проверить old writer с new reader: default, null или manual route должны быть записаны явно.
- Перед реальным запуском выполнить integration test выбранного broker, storage и external API. Проверить их actual retries, crash window и rollback отдельно.
Граница этого разбора
Здесь нет реального broker, schema registry, consumer offset, listener, HTTP, database, external payment, audit store, browser, CI или deployment. Нет данных пользователей, measured throughput, lag или историй production incident. Также нет claim, что consumerId:source:event.id гарантирует exactly-once. Он даёт один воспроизводимый answer внутри Map: текущий учебный consumer уже записал result для этого source:id или ещё нет.
После чтения стоит выбрать один невысокорисковый projection, а не внешний effect, и пройти тот же маршрут на интеграционном стенде. Хороший тест покажет initial delivery, duplicate, controlled replay и unknown schema с реальными версиями выбранных компонентов. Если для внешнего effect нет подтверждённого idempotency contract, автоматический replay не следует включать. В таком случае честный результат — сохранить evidence и передать решение владельцу операции.
Проверяемые источники
- Apache Kafka 2.7.0 KafkaProducer API — версионная документация линии 2.7 предупреждает, что retry может привести к duplicate. Fixture не запускает Kafka, не использует offset, partition, transaction или producer API.
- CloudEvents Core Specification snapshot, 16 апреля 2021 (v1.0.2-wip) — исторический working draft существовал к июню 2021 года и различает context attributes, event data и protocol binding. Учебный envelope ниже не является CloudEvent и не реализует binding.
- Apache Avro 1.10.1 Specification — Schema Resolution — официальная спецификация различает writer и reader schema и описывает resolution. Здесь она служит ориентиром для явного consumer contract, а не форматом payload или schema registry.