diff --git a/editorial/agent-rewrites/235.json b/editorial/agent-rewrites/235.json index a48ab32..b4b79e3 100644 --- a/editorial/agent-rewrites/235.json +++ b/editorial/agent-rewrites/235.json @@ -1,7 +1,7 @@ { "index": 235, "slug": "editorial-2021-06-field-event-driven", - "title": "Replay событий: как не создать второй effect и не потерять контракт схемы", - "excerpt": "Повторная доставка события не должна превращаться в новый факт. Разбираем identity, ledger, совместимость схемы и controlled replay на учебном примере, который явно отделяет duplicate от новой версии consumer.", - "contentHtml": "
После сбоя consumer команда часто видит один и тот же симптом: нужно проиграть события ещё раз. Оператор запускает replay, consumer снова получает запись, а в storage появляется второй result. Если replay получил новый event id, система уже не отличает повтор старого факта от нового факта. Если id сохранился, но consumer не ведёт ledger, он всё равно может повторить effect. Цена ошибки — не только дубль строки. Платёж, письмо или изменение внешней системы могут выполниться дважды. Затем команда теряет ответ на главный вопрос: какой contract создал каждый результат.
\nТезис простой: replay должен сохранять логическую identity исходного события, а consumer должен проверять схему и receipt до effect. Один и тот же source:id в том же consumer contract даёт duplicate и подавляется. Новый расчёт требует нового явно названного contract или resultVersion. Неизвестная схема останавливает обработку. Это не обещание exactly-once. Это проверяемая граница, которая не даёт техническому повтору притвориться новым бизнес-фактом.
Событие описывает уже произошедший факт. В учебном примере оно содержит source, id, type, subject, occurredAt, schemaVersion и data. Источник и id вместе образуют identity: training://orders/order-104:evt-order-104-paid-01. Новый запуск может иметь отдельную операционную причину, но эта причина не должна менять исходные поля.
Consumer contract отвечает на другой вопрос: как этот consumer интерпретирует payload. Пусть orders-projection@1 читает orderId и status, а orders-projection@2 дополнительно читает paymentReference. Один event может законно дать два локальных projection result, если домен разрешает хранить обе версии. Но это не значит, что любой внешний effect можно повторить дважды. Для email, платежа или HTTP-вызова нужен отдельный idempotency contract на внешней границе.
const ledgerKey = [consumerId, event.source, event.id].join(':');\n\nif (!acceptedSchemaVersions.includes(event.schemaVersion)) {\n return { state: 'contract-update-required', effectAllowed: false };\n}\n\nif (ledger.has(ledgerKey)) {\n return { state: 'duplicate-or-replay-suppressed', effectAllowed: false };\n}\n\nconst result = project(event, consumerId);\nledger.set(ledgerKey, { resultVersion, result });\nreturn { state: 'effect-recorded', effectAllowed: true };\nЭтот фрагмент показывает порядок, а не готовую библиотеку. В реальной системе запись receipt и effect должны иметь согласованный storage contract. Map процесса не защищает от падения между внешним вызовом и записью ledger. Если такая аварийная граница существует, автоматический replay нельзя объявлять безопасным без отдельного решения.
\nПусть пришёл учебный event order.status.changed со схемой v1. Первый consumer записывает результат для orders-projection@1. Сеть повторяет delivery. Ledger видит тот же ключ и подавляет второй result. Оператор запускает controlled replay. Ключ не меняется, поэтому replay тоже подавляется. Затем команда включает orders-projection@2. Это уже другой declared contract. Он может получить отдельный projection result, если такое решение принято явно и не смешано с внешним effect.
| Delivery | Identity | Проверка | Результат |
|---|---|---|---|
initial | тот же source:id | ключа нет | записать один result |
duplicate | тот же source:id | ключ есть у того же consumer | подавить новый effect |
controlled replay | тот же source:id | ключ есть у того же consumer | сохранить evidence, не писать второй result |
new contract | тот же event | другой consumerId и resultVersion | отдельный projection, если он разрешён |
schema v3 | новый или повторный event | версия не принята consumer | остановить effect |
Нельзя удалять старую receipt перед историческим replay. Иначе новая запись скроет, что старый contract уже обработал event. Нельзя и автоматически считать любой новый contract безопасным: два projection допустимы не во всех доменах. Сначала называют effect и его владельца, потом выбирают ключ, receipt и способ восстановления.
\nПроверка identity не заменяет проверку schema. Если event v3 содержит поле state, которого consumer не знает, отсутствие ключа в ledger не даёт права выполнить effect. Сначала consumer проверяет envelope и accepted schema versions. Затем выбирает contract. Только после этого он читает ledger. Отказ должен сохранить причину: unknown type, unsupported schema или invalid field. Такой след отличает остановленный event от потерянной delivery.
Добавление поля может быть совместимым для старого reader, если reader его игнорирует и обязательные поля сохраняют смысл. Новый reader может представить отсутствие поля как null или default, но это решение должно быть частью его contract. Переименование status в state нельзя выдавать за additive change. Учебный пример проверяет только эту пару контрактов; он не доказывает совместимость Avro, JSON Schema или вашей registry без отдельного теста.
| Симптом | Причина | Проверка | Действие |
|---|---|---|---|
| Второй result для того же event | Нет ledger или ключ построен без source | Сравнить consumer, source, id и receipt | Сделать identity явной и остановить повторный effect |
| Replay выглядит как новое событие | Оператор заменил исходный id | Сопоставить event log и операционную запись replay | Вернуть исходную identity, причину хранить отдельно |
| Старый event не читается новым consumer | Нет правила absence/default | Прогнать writer v1 через reader v2 | Добавить явный compatibility rule или manual route |
| Unknown schema записывает effect | Ledger читается до проверки contract | Проверить порядок веток и отрицательный тест v3 | Запрещать effect до accepted schema |
| Внешний вызов повторился после сбоя | Receipt и effect не образуют атомарную границу | Смоделировать crash между вызовами | Остановить auto-replay и согласовать внешний idempotency key |
source, id, type, subject, schema version, consumer contract и причину replay.Учебный пример хранит состояние в Map процесса Node. Он не запускает broker, schema registry, database, Kubernetes, HTTP или внешний платёж. Он не измеряет throughput и не доказывает exactly-once. Идентификаторы, даты, source, payload и результаты вымышлены. Реальная гарантия зависит от transaction boundary, retry policy, partitioning, storage и внешних API. Kafka отдельно предупреждает о возможности duplicate при retry; это ограничивает формулировку, но не даёт готовой архитектуры. CloudEvents описывает envelope и protocol binding, а не receipt бизнес-операции.
\nГотовность проверяется не фразой «replay прошёл». Для выбранного consumer должны воспроизводиться четыре результата: initial delivery создаёт один result; duplicate и controlled replay не создают второй; разрешённый новый contract создаёт отдельный versioned result; неизвестная schema останавливается без effect. Для внешней операции добавьте доказательство её idempotency key или ручной stop. Если хотя бы один результат нельзя объяснить по event identity, contract, ledger и receipt, replay ещё не готов.
\nСимптом после восстановления consumer звучит просто: «нужно проиграть события ещё раз». Затем один event приходит повторно, второй consumer уже обновлён, а в storage появляется непонятный result. Если replay создаёт новый id, невозможно отличить повтор старого факта от нового факта. Если он сохраняет id, но consumer не ведёт ledger, можно записать второй effect. Цена — не только дубль. Команда теряет доказательство, по какому contract был получен каждый результат, и не может безопасно решить, что делать со старой схемой.
\nПолезная точка старта — не ручной запуск всех 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.
Самая опасная «починка» — сделать новое event id, чтобы consumer не счёл запись duplicate. Так обходят проверку, но меняют вопрос. Новый id может означать новый факт, исправленную команду или технический replay; эти случаи нельзя сливать. Для controlled replay неизменным остаётся исходный id, source, type, subject и payload schema. Дополнительная причина replay может жить рядом с операционной записью, но не должна подменять исходный event. Тогда ledger способен ответить: этот contract уже обработал данный вход или нет.
\nКлюч 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();\nconst first = consumeTrainingEvent(ledger, event, "orders-projection@1", "initial-delivery");\nconst duplicate = consumeTrainingEvent(ledger, event, "orders-projection@1", "duplicate-delivery");\nconst replay = consumeTrainingEvent(ledger, event, "orders-projection@1", "controlled-replay");\n\nfirst.effectWritten; // true\nduplicate.effectWritten; // false\nreplay.effectWritten; // false\nledger.size; // 1\n| 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 по тому же 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(\n new Map(),\n trainingEvents.orderStatusV1,\n "orders-projection@2",\n "historical-replay",\n);\n\nresult.resultVersion; // "2021-06.orders-projection.2"\nresult.inputSchemaVersion; // 1\nresult.projection.paymentReference; // null\n\n// Result version описывает consumer contract, не версию broker-а.\nПроверка 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");\n\ndecision;\n// {\n// state: "contract-update-required",\n// inputSchemaVersion: 3,\n// effectAllowed: false,\n// }\n\n// Не угадываем, что state: "settled" эквивалентен status: "paid".\nФикстура использует два 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 всегда пишет ещё раз».
\nconst fixture = runEventIntegrationFixture();\nif (!Object.values(fixture.assertions).every(Boolean)) {\n throw new Error('event contract changed without an explicit decision');\n}\n\nconsole.log(fixture.deliveries.duplicateDelivery.state);\n// 'duplicate-or-replay-suppressed'\nApache Kafka 2.7 producer API прямо рассматривает retries и возможность duplicate в некоторых режимах. Это хороший повод не обещать exactly-once одним словом. Конкретные guarantees зависят от producer, broker, consumer, storage и внешней границы. CloudEvents в историческом snapshot апреля 2021 года помогает разделить context и event data, но не задаёт business receipt. В статье оба источника ограничивают формулировку, а не дают готовую implementation.
\n| Наблюдение | Что собрать | Безопасное действие | Чего не делать |
|---|---|---|---|
тот же 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 заменой транзакции |
Здесь нет реального 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 и передать решение владельцу операции.
\n