DarkRiDDeR15 мин

Событийная интеграция: какой envelope зафиксировать до первого consumer

АрхитектураСобытияПрактика

Симптом появляется не в момент отправки события, а после первого независимого consumer. Заказ уже сменил статус, но обработчик не может сказать, кто создал сообщение, какую схему он читает и можно ли безопасно повторить вход. Иногда событие приходит второй раз. Иногда producer добавляет поле, а consumer начинает трактовать его как обязательное. Цена — не один красный лог. Команда теряет связь между доменным изменением, конкретной доставкой и результатом обработчика; затем повтор превращается в случайный запуск.

В июне 2021 года я бы начал не с выбора broker, а с одного маленького договора. Producer обязан отдать envelope с происхождением и типом, payload обязан назвать свою schemaVersion, consumer обязан заранее записать, какие версии и поля он читает, а результат обязан сохранить версию самого consumer. Ниже это проверяется только в памяти Node. Все идентификаторы, даты, source, payload, результаты и правила ниже учебные. Fixture хранит состояние только в Map процесса Node: это не broker, не schema registry, не production event, не реализация CloudEvents или Kafka, не transport, не база и не измерение throughput. Поэтому пример помогает обсудить границу, но не доказывает настройку реальной инфраструктуры.

Сначала отделяем событие от доставки

Событие описывает факт, который producer решил передать: в учебном случае изменился статус заказа. Доставка описывает попытку передать этот факт конкретному consumer. Это не одно и то же. Одно событие может быть доставлено дважды после неясного результата подтверждения. Один consumer может прочитать событие позже другого. Если в коде есть только JSON без идентификатора, source и версии, эти случаи уже нельзя отличить по факту, остаётся угадывать по времени лога.

Envelope не должен становиться свалкой всей предметной модели. Его задача — дать устойчивые координаты: какой договор envelope применён, какой event id наблюдаем, кто его сформировал, к какому предмету он относится, когда producer его создал и по какой схеме лежит payload. Payload содержит данные конкретного типа. Result consumer содержит уже другое: какую версию своего договора применил consumer и записал ли он эффект. Это три разных объекта с разными владельцами.

Минимальный учебный contract: каждый слой отвечает на свой вопрос
СлойПоле или правилоВладелецЧто проверить до действия
EnvelopecontractVersion, id, source, typeproducer контрактавсе обязательные поля непустые и type ожидаем consumer-ом
Связь с предметомsubjectproducer и владелец доменаsubject указывает на объект, для которого имеет смысл разбор
PayloadschemaVersion и dataproducer payloadconsumer явно принимает именно эту версию и stable fields
Consumerсписок versions и интерпретация полейвладелец конкретного consumerнет молчаливого предположения о новом поле или значении
ResultresultVersion и ключ consumerId:source:idconsumer и его хранилище результатаповтор того же source:id не создаёт второй учебный result без нового договора

CloudEvents полезен здесь как дисциплина метаданных: исторический snapshot апреля 2021 года отдельно описывает context события, event data и protocol binding. Но нельзя взять несколько похожих имён и назвать любой object CloudEvent. У формата есть свои обязательные атрибуты и bindings. В этой статье contractVersion и schemaVersion принадлежат нашему учебному договору. Fixture не проверяет bindings и SDK.

Пишем envelope так, чтобы его можно было проверить

Поле id должно быть неизменным для одной логической записи. Оно не обязано одновременно быть ключом любого бизнес-эффекта: например, письмо может иметь отдельный idempotency key. Но без id нельзя доказать, что controlled replay относится к тому же входу. source отделяет producer от consumer и не должен строиться из случайного hostname. В учебном примере разрешён только закрытый prefix training://orders/; это делает ошибочный внешний source наблюдаемым до работы consumer.

// Собственный учебный envelope. Это не CloudEvent.
const event = {
  contractVersion: '2021-06',
  id: 'evt-order-104-paid-02',
  source: 'training://orders/order-104',
  type: 'order.status.changed',
  subject: 'order-104',
  occurredAt: '2021-06-14T10:31:00.000Z',
  schemaVersion: 2,
  data: {
    orderId: 'order-104',
    status: 'paid',
    paymentReference: 'training-pay-77',
  },
};

// id нужен для наблюдения и controlled replay.
// schemaVersion принадлежит payload contract, не transport.

Время occurredAt не назначает порядок обработки. Это время, которое сообщил producer, а не позиция в topic и не время, когда consumer выполнил effect. Если домен требует порядок, нужен отдельный sequence contract и его owner. Если порядок не нужен, не стоит создавать ложную гарантию из timestamp. В этой партии это ограничение остаётся явным: fixture не сортирует delivery, не моделирует clock skew и не выбирает partition.

const checked = validateTrainingEnvelope(event);

if (!checked.ok) {
  return { state: 'manual-review', reason: checked.reason };
}

// До consumer сохраняем доказательство: source, id, type, subject, schemaVersion.
return { state: "ready-for-declared-consumer", envelopeKey: checked.envelopeKey };

Проверка envelope должна завершиться до любой интерпретации payload. Иначе consumer сначала создаст effect, а потом обнаружит, что source или type были неожиданными. В результате validateTrainingEnvelope() возвращает причину вроде missing-or-empty-type или source-outside-training-orders-boundary. Это не универсальный validator. Реальный сервис может потребовать подпись, tenant, permissions, content type и ограничения размера. Важен принцип: причина отказа должна быть в result, а не скрываться за общим exception.

Вертикальная схема учебной топологии: producer формирует envelope и payload, delivery передаёт их consumer, consumer сначала валидирует contract, затем записывает versioned result в ledger; отдельная стрелка показывает duplicate или controlled replay с тем же source:id
Топология разделяет факт, доставку и результат. Broker на схеме — граница передачи, а не обещание конкретного продукта или гарантии.

SchemaVersion — не декоративное число

Версия payload нужна не для красивого суффикса в названии события. Она говорит consumer, по каким правилам он может прочитать data. В учебном schemaVersion: 1 содержит orderId и status. Версия 2 добавляет paymentReference. Старый consumer v1 читает только стабильную пару и осознанно игнорирует новое поле. Новый consumer v2 умеет отдать paymentReference: null, когда воспроизводит старый v1 event. Это ограниченная, проверяемая политика, а не слово «backward compatible» без границы.

Apache Avro разделяет writer schema и reader schema, а правила resolution зависят от конкретного формата и схем. Из этого полезно перенести не кодек, а вопрос: что именно writer записал и что reader имеет право ожидать. Наш JSON object не является Avro record. В нём нет writer schema, fingerprint, default из Avro и реального registry. Поэтому добавление поля здесь совместимо только потому, что два наших consumer contract явно так определены.

const v1ConsumerOnV2 = projectForConsumer(event, "orders-projection@1");

// v1 читает только стабильные поля.
// paymentReference он намеренно не интерпретирует.
v1ConsumerOnV2.projection;
// { orderId: 'order-104', status: 'paid' }

const v2ConsumerOnV1 = projectForConsumer(trainingEvents.orderStatusV1, "orders-projection@2");
// Новый consumer нормализует отсутствующее необязательное поле к null.
Матрица совместимости учебных contracts
ВходConsumer contractРазрешённый результатЧего не обещает правило
schema v1: orderId, statusorders-projection@1stable projection из двух полейчто v1 понимает будущие значения status
schema v2: v1 + paymentReferenceorders-projection@1то же stable projection, поле v2 намеренно игнорируетсячто любой added field всегда безопасен
schema v1 без нового поляorders-projection@2paymentReference: null в versioned resultчто null равен неизвестному business state
schema v2 с необязательным полемorders-projection@2projection с проверенным string или nullчто значение ссылки корректно в соседней системе
schema v3 с изменённым смыслом stateлюбой contract этой fixturecontract-update-requiredчто consumer может угадать семантику

Consumer contract виден до replay

У consumer должна быть короткая объявленная граница: accepted schema versions, набор читаемых полей, результат и ключ, по которому он видит повтор. Иначе внедрение становится опасной схемой «сначала обновим producer, потом посмотрим на ошибки». В примере orders-projection@1 и orders-projection@2 не обозначают версии broker. Это два независимых договора чтения одного type. Их результат хранит resultVersion, чтобы replay можно было связать с интерпретацией, которая действовала в момент обработки.

Не надо автоматически принимать каждую большую версию, если JSON проходит синтаксис. Версия 3 в fixture меняет привычный status на state. Мы не считаем settled синонимом paid без решения владельца домена. Consumer возвращает contract-update-required и не пишет effect. Такой отказ дешевле тихой подмены смысла: в логе остаётся event id, input schemaVersion, consumer id и причина, по которым можно подготовить миграцию или отдельный адаптер.

const decision = projectForConsumer(trainingEvents.orderStatusV3, "orders-projection@2");

decision;
// {
//   state: "contract-update-required",
//   inputSchemaVersion: 3,
//   effectAllowed: false,
// }

// Не угадываем, что state: "settled" эквивалентен status: "paid".

Маршрут перед подключением реального broker

  1. Выбрать один event type и назвать symptom: какой consumer сейчас не может безопасно понять или повторить вход.
  2. Отделить envelope, payload и consumer result. Для каждого записать владельца и только необходимые поля.
  3. Определить стабильные поля, которые v1 consumer действительно читает, и одно добавочное поле следующей версии. Не менять семантику под именем «добавили поле».
  4. Написать validator envelope до кода effect: id, source, type, subject, contractVersion, schemaVersion и обязательные stable fields.
  5. Зафиксировать consumer contract: accepted versions, projection, resultVersion и исход для неизвестной версии.
  6. Проверить controlled fixture с v1, v2, duplicate и replay. Убедиться, что replay сохраняет source:id, а не создаёт новый вход для обхода ledger.
  7. Только после этого выбрать конкретный broker, serializer, storage result и integration test. Отдельно записать их реальные guarantees и failure modes.

Что остаётся за границей этого шага

У этой модели нет настоящего topic, consumer group, offset, transaction, outbox, schema registry, authorisation, encryption, retries сети, retention или delivery SLA. Нет и real production event: значения order-104 и training-pay-77 специально вымышлены. Kafka 2.7 documentation показывает, почему retry нельзя автоматически отождествлять с единственной доставкой, но этот факт не превращает Map в Kafka client. Аналогично CloudEvents не даёт бизнес-совместимость просто наличием envelope.

Практический следующий шаг — взять один безобидный event type проекта и сделать такой же evidence packet: serialized envelope, declared payload schema, consumer version, sample result и controlled replay. Если хотя бы одно поле нельзя объяснить владельцем, не публикуйте его в общий contract. Сначала сузьте событие до проверяемого факта. Тогда новый consumer будет явным договором, а не догадкой по JSON.

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

  • 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.
  • Apache Kafka 2.7.0 KafkaProducer API — версионная документация линии 2.7 предупреждает, что retry может привести к duplicate. Fixture не запускает Kafka, не использует offset, partition, transaction или producer API.