From f5996ae296f3d362a88085385c97582568cfe115 Mon Sep 17 00:00:00 2001 From: "E.Gavrilov" Date: Thu, 3 Sep 2026 20:53:06 +0300 Subject: [PATCH] editorial: refine event replay article 235 --- editorial/agent-rewrites/235.json | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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. Это проверяемая граница, которая не даёт техническому повтору притвориться новым бизнес-фактом.

\n

Что именно повторяется

\n

Событие описывает уже произошедший факт. В учебном примере оно содержит source, id, type, subject, occurredAt, schemaVersion и data. Источник и id вместе образуют identity: training://orders/order-104:evt-order-104-paid-01. Новый запуск может иметь отдельную операционную причину, но эта причина не должна менять исходные поля.

\n

Consumer contract отвечает на другой вопрос: как этот consumer интерпретирует payload. Пусть orders-projection@1 читает orderId и status, а orders-projection@2 дополнительно читает paymentReference. Один event может законно дать два локальных projection result, если домен разрешает хранить обе версии. Но это не значит, что любой внешний effect можно повторить дважды. Для email, платежа или HTTP-вызова нужен отдельный idempotency contract на внешней границе.

\n
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 и три delivery

\n

Пусть пришёл учебный event order.status.changed со схемой v1. Первый consumer записывает результат для orders-projection@1. Сеть повторяет delivery. Ledger видит тот же ключ и подавляет второй result. Оператор запускает controlled replay. Ключ не меняется, поэтому replay тоже подавляется. Затем команда включает orders-projection@2. Это уже другой declared contract. Он может получить отдельный projection result, если такое решение принято явно и не смешано с внешним effect.

\n
Решение для одного учебного event
DeliveryIdentityПроверкаРезультат
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
\n

Нельзя удалять старую receipt перед историческим replay. Иначе новая запись скроет, что старый contract уже обработал event. Нельзя и автоматически считать любой новый contract безопасным: два projection допустимы не во всех доменах. Сначала называют effect и его владельца, потом выбирают ключ, receipt и способ восстановления.

\n
\"Дерево
Порядок проверки отделяет повтор того же result от новой явно объявленной интерпретации. Схема учебная и не описывает конкретный broker или production-систему.
\n

Схема проверяется до ledger

\n

Проверка identity не заменяет проверку schema. Если event v3 содержит поле state, которого consumer не знает, отсутствие ключа в ledger не даёт права выполнить effect. Сначала consumer проверяет envelope и accepted schema versions. Затем выбирает contract. Только после этого он читает ledger. Отказ должен сохранить причину: unknown type, unsupported schema или invalid field. Такой след отличает остановленный event от потерянной delivery.

\n

Добавление поля может быть совместимым для старого reader, если reader его игнорирует и обязательные поля сохраняют смысл. Новый reader может представить отсутствие поля как null или default, но это решение должно быть частью его contract. Переименование status в state нельзя выдавать за additive change. Учебный пример проверяет только эту пару контрактов; он не доказывает совместимость Avro, JSON Schema или вашей registry без отдельного теста.

\n

Симптом → причина → проверка → действие

\n
Диагностика перед replay
СимптомПричинаПроверкаДействие
Второй 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 записывает effectLedger читается до проверки contractПроверить порядок веток и отрицательный тест v3Запрещать effect до accepted schema
Внешний вызов повторился после сбояReceipt и effect не образуют атомарную границуСмоделировать crash между вызовамиОстановить auto-replay и согласовать внешний idempotency key
\n

Порядок controlled replay

\n
  1. Соберите evidence packet. Зафиксируйте исходные source, id, type, subject, schema version, consumer contract и причину replay.
  2. Проверьте envelope. Отклоните пустые поля, неожиданный source и неизвестный type до чтения payload и до любого effect.
  3. Проверьте схему и reader contract. Назовите принятые версии, правила отсутствующих полей и resultVersion. Не угадывайте новое поле по похожему имени.
  4. Постройте ledger key. Включите consumer identity, source и event id. Отдельно запишите, какой ключ защищает внешний business effect.
  5. Подавите duplicate. Если тот же contract уже имеет receipt, не удаляйте её и не создавайте новый id. Сохраните delivery evidence.
  6. Объявите новую интерпретацию. Если нужен новый projection, используйте новый contract или resultVersion и получите явное решение о допустимости второго результата.
  7. Проверьте аварийное окно. На интеграционном стенде повторите crash между effect и receipt, retries broker и восстановление consumer. Для внешнего сервиса подтвердите его фактический idempotency contract.
\n

Ограничения и критерий готовности

\n

Учебный пример хранит состояние в 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

Проверяемые источники

\n" + "title": "Replay события: как не записать второй result и не скрыть новую схему", + "excerpt": "Replay полезен только тогда, когда видно исходный event, consumer contract и уже записанный result. Разбираем evidence packet, duplicate, controlled replay, schema mismatch и маршрут без обещания exactly-once.", + "contentHtml": "

Симптом после восстановления 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.

\n

Replay сохраняет логическую идентичность

\n

Самая опасная «починка» — сделать новое 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 вместе с результатом.

\n
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: ожидаемый учебный результат
Delivery kindEvent identityLedger до шагаРешениеEffect
initial-deliveryтот же source:idнет ключа consumereffect-recordedодин учебный result записан
duplicate-deliveryтот же source:idключ уже естьduplicate-or-replay-suppressedновый result не пишется
controlled-replayтот же source:idключ уже естьduplicate-or-replay-suppressedreplay не обходил ledger
historical-replay в другом contractтот же source:idдругой ledger consumereffect-recorded с новой resultVersionразный reader result допустим
delivery schema v3новый event id, неизвестная schemaневажноcontract-update-requiredeffect запрещён
\n

Duplicate и новая интерпретация — разные развилки

\n

Duplicate по тому же contract не должен создавать второй result. Но новое правило consumer иногда действительно требует пересчитать projection. Тогда не стоит удалять старую ledger запись и притворяться, что история не существовала. В fixture orders-projection@2 имеет другой ledger key и resultVersion; historical replay v1 может создать свой result, потому что это другой declared reader contract. Такой выбор не делает результат «истиннее», он делает видимой новую интерпретацию.

\n

Перед этим шагом нужно определить, разрешён ли второй projection в домене. Для email, платёжа или изменения внешней системы одного consumerId:source:event.id почти наверняка недостаточно: потребуется effect key на внешней границе, транзакция или ручное решение. Для локальной read projection иногда допустимо хранить две версии и переключать reader после сверки. Статья не выбирает между этими архитектурами. Она требует, чтобы автор replay написал, какой effect будет создан и почему второй contract имеет право его создавать.

\n
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
\"Вертикальное
Диагностика не считает replay ошибкой сама по себе. Она отделяет повтор того же result от новой явно объявленной интерпретации.
\n

Схема не должна исчезать за duplicate

\n

Проверка ledger не заменяет проверку schema. Если v3 event пришёл с незнакомым state, consumer не должен сначала посмотреть id, не найти его и записать effect только потому, что это первый delivery. В учебном алгоритме порядок другой: validate envelope, выбрать consumer contract, проверить accepted schema version, затем читать ledger. Это сохраняет важный факт: новый event был получен, но effect запрещён из-за договора, а не потерян как «неизвестная ошибка».

\n

Также нельзя делать обратное: любой duplicate автоматически игнорировать до записи diagnostics. Result suppress содержит deliveryKind, ledgerKey, input schema и result version первого результата. Это не production audit trail, но минимальное evidence помогает отличить повтор сети от operator replay. Если реальная система не сохраняет эти факты, расследование начнётся с догадки: очередной запуск создал запись или просто повторил уже завершённую delivery.

\n
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

Fixture задаёт узкую, но полезную проверку

\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 всегда пишет ещё раз».

\n
const 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'
\n

Apache Kafka 2.7 producer API прямо рассматривает retries и возможность duplicate в некоторых режимах. Это хороший повод не обещать exactly-once одним словом. Конкретные guarantees зависят от producer, broker, consumer, storage и внешней границы. CloudEvents в историческом snapshot апреля 2021 года помогает разделить context и event data, но не задаёт business receipt. В статье оба источника ограничивают формулировку, а не дают готовую implementation.

\n
Диагноз перед replay: факт, проверка, действие
НаблюдениеЧто собратьБезопасное действиеЧего не делать
тот же source:id, same consumer contractledger key и первый resultVersionsuppress duplicate, сохранить evidence deliveryсоздавать новый id ради обхода проверки
тот же event, новый declared consumer contractстарый и новый resultVersion, тип effectотдельный controlled replay после явного решенияудалять старую receipt и терять историю
schemaVersion не объявлена consumerenvelope, input schema, причина отказаcontract update или manual route без effectугадывать поле по похожему имени
source или type неожиданныvalidation reason до payloadrejected envelope и проверка producer boundaryвыполнять partial effect для «похожего» входа
внешний effect уже мог быть сделанidempotency contract внешней системы и receiptостановить автоматический replay до доказательствасчитать Map заменой транзакции
\n

Маршрут controlled replay

\n
  1. Зафиксировать исходный event id, source, type, subject, payload schema и причину replay. Не заменять их новым JSON.
  2. Проверить envelope и consumer contract до ledger. Unknown type или schema должен дать объяснимый отказ без effect.
  3. Выбрать key, который соответствует именно данному consumer result; отдельно обсудить key доменного или внешнего effect.
  4. Если ledger уже содержит same consumer contract + source:id, suppress replay и сохранить delivery evidence.
  5. Если нужен новый расчёт, ввести новый declared consumer contract/resultVersion. Не перезаписывать старую интерпретацию без следа.
  6. Для historical event проверить old writer с new reader: default, null или manual route должны быть записаны явно.
  7. Перед реальным запуском выполнить integration test выбранного broker, storage и external API. Проверить их actual retries, crash window и rollback отдельно.
\n

Граница этого разбора

\n

Здесь нет реального 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 или ещё нет.

\n

После чтения стоит выбрать один невысокорисковый projection, а не внешний effect, и пройти тот же маршрут на интеграционном стенде. Хороший тест покажет initial delivery, duplicate, controlled replay и unknown schema с реальными версиями выбранных компонентов. Если для внешнего effect нет подтверждённого idempotency contract, автоматический replay не следует включать. В таком случае честный результат — сохранить evidence и передать решение владельцу операции.

\n

Проверяемые источники

\n" }