Симптом обычно выглядит безобидно: producer добавил поле paymentReference, JSON по-прежнему валиден, а старый consumer либо падает на строгой проверке, либо начинает использовать значение не по договору. Хуже другой случай: поле переименовали или поменяли его смысл, consumer продолжил работать и записал правдоподобный, но неверный результат. Цена тихой совместимости выше явного отказа: затем невозможно восстановить, какая версия входа породила конкретную запись.
Здесь важно развести две вещи. Формат может суметь распарсить bytes, а consumer может не иметь права интерпретировать бизнес-смысл. В учебной модели v2 добавляет одно необязательное поле к устойчивой паре orderId и status. Consumer v1 объявляет, что читает только устойчивую пару. Consumer v2 умеет сохранить новое поле, но при replay v1 нормализует его к null. Модель не использует Avro, Kafka, CloudEvents или schema registry. Все идентификаторы, даты, source, payload, результаты и правила ниже учебные. Fixture хранит состояние только в Map процесса Node: это не broker, не schema registry, не production event, не реализация CloudEvents или Kafka, не transport, не база и не измерение throughput.
Совместимость начинается с пары writer и reader
Фраза «схема совместима» бесполезна без двух участников. Нужно назвать writer schema, reader contract и направление проверки. Новый writer v2 может быть совместим со старым reader v1, если reader действительно игнорирует добавленное поле и его смысл не меняет старые поля. Старый writer v1 может быть совместим с новым reader v2, если новый reader знает, как честно обработать отсутствие нового поля. Но изменение status на state не становится совместимым только потому, что оба значения строки. Это уже смена интерпретации.
Apache Avro формулирует это через writer и reader schema resolution. Конкретные правила зависят от record, default, alias и serialization. Для автора прикладного contract полезна более простая привычка: на каждый change показать одну старую запись, один новый consumer и один новый event для старого consumer. Если эти два направления не проверены, словом compatible называют только надежду. В нашей fixture оба направления являются отдельными assertions.
| Writer event | Reader contract | Учебный verdict | Причина |
|---|---|---|---|
v1: orderId, status | v1 | готово | оба используют один набор stable fields |
v2: v1 + optional paymentReference | v1 | готово с игнорированием поля | v1 contract читает только описанную пару |
v1 без paymentReference | v2 | готово с null | v2 contract явно определяет default представления |
| v2 с пустой или неверной ссылкой | v2 | rejected envelope | валидность поля проверяется до projection |
v3: state вместо status | v1 или v2 | contract update required | смысл stable field больше не доказан |
Стабильное поле имеет не только имя
Стабильность — это имя, тип и договорённость о смысле. status: "paid" нельзя заменить на state: "settled" и сказать старому consumer, что он должен «как-нибудь понять». Даже если домен считает значения близкими, у consumer могут быть ветки, аудит, SQL-проекция или внешняя команда, которые используют старое значение как ключ. Поэтому schemaVersion должна расти при таком change, а consumer должен либо получить отдельный адаптер с тестом, либо остановить effect до решения.
Добавочное поле тоже не автоматически безопасно. Оно безопасно для конкретного reader, когда reader не использует unknown fields для валидации и новое поле не меняет meaning предыдущих. Например, paymentReference в нашем v2 — optional string, который v1 не читает. Но если producer вводит currency и одновременно начинает иначе понимать amount, это не «добавили currency». Это изменение смысла суммы, требующее отдельного type или миграции. Контракт удобнее держать маленьким, чем потом спасать широкую схему исключениями.
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.
Versioned result связывает replay с интерпретацией
Event schemaVersion недостаточно, когда один и тот же event воспроизводят разные consumer. В fixture результат содержит inputSchemaVersion и resultVersion. Первый отвечает на вопрос, с каким payload пришёл вход. Второй отвечает, какой consumer contract создал projection. Это особенно важно при исправлении consumer: нельзя сказать, что replay «пересчитал данные», если не видно, старый или новый код дал результат.
Здесь resultVersion не является номером deploy, Git commit или версией broker. Это стабильное имя интерпретации 2021-06.orders-projection.1 или 2021-06.orders-projection.2. В реальном проекте к нему могут добавиться build, schema fingerprint или migration id. Но не стоит приклеивать всё сразу: достаточно обеспечить один ответ на вопрос расследования — по какому договору consumer прочёл event и что он записал.
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-а.
| Факт | Зачем нужен | Недостаточный заменитель |
|---|---|---|
event.id и source | связать result с исходным логическим входом | только время обработки |
inputSchemaVersion | понять форму data, которую видел consumer | название topic или queue |
consumerId | отделить два самостоятельных read contract | общее имя сервиса |
resultVersion | отделить старую и новую интерпретацию при replay | случайный build timestamp |
решение effect-recorded или отказ | не спутать обработанный input с безопасно интерпретированным input | одна строка «consumer finished» |
Unknown version должна менять маршрут, а не парсинг
Некоторые команды делают consumer permissive: он принимает любое число schemaVersion и берёт знакомые поля, надеясь, что остальное неважно. Такой подход удобен до первого semantic change. В нашем примере v3 содержит state, а не status. Envelope по форме всё ещё даёт origin, id и type, но v1/v2 contracts не объявили это значение. Поэтому projectForConsumer() возвращает contract-update-required и запрещает effect.
Это не означает, что каждое новое поле останавливает весь поток. Значит другое: разные категории изменений имеют разные правила. Новый optional field, который reader не читает, может пройти по заранее описанной ветке. Новая обязательная семантика, удаление stable field, смена единицы или переименование требуют explicit consumer change. Если такой change нельзя выполнить быстро, лучше сохранить event и manual evidence, чем превратить неизвестное значение в default без владельца.
const decision = projectForConsumer(trainingEvents.orderStatusV3, "orders-projection@2");
decision;
// {
// state: "contract-update-required",
// inputSchemaVersion: 3,
// effectAllowed: false,
// }
// Не угадываем, что state: "settled" эквивалентен status: "paid".
Schema registry полезен, но не заменяет договор
Schema registry может хранить definitions, compatibility modes и историю. Но сам факт регистрации не доказывает, что конкретный consumer хранит result idempotently, что event type соответствует доменному действию или что versioned replay безопасен. И наоборот, маленький проект может начать с versioned fixture и JSON-schema-like checks без registry, пока договор виден в коде и review. В обоих случаях остаются одни и те же вопросы: кто публикует schema, что принимает reader, где записан migration и как остановить неизвестный вход.
CloudEvents в историческом snapshot апреля 2021 года описывает общую форму event metadata, а Kafka 2.7 documentation различает producer send и возможность duplicate при retry. Ни один источник не говорит, что добавление поля в любую JSON data автоматически совместимо со всеми business consumer. Вся совместимость в этой статье ограничена двумя contracts и одной функцией projection. Так и должно быть: общий стандарт помогает передать envelope, но не владеет семантикой заказа.
Маршрут изменения схемы
- Записать одну старую запись и один будущий event в отдельном fixture. Не начинать с массового изменения producer.
- Назвать stable fields, их тип и смысл. Если смысл меняется, считать это новым contract, даже когда JSON key похож.
- Проверить новый writer со старым reader: какие поля reader читает, какие игнорирует и почему это безопасно.
- Проверить старый writer с новым reader: какой explicit default или отдельный route получает отсутствующее поле.
- Добавить
inputSchemaVersionиresultVersionк result consumer, чтобы replay был объясним. - Для unknown schemaVersion вернуть contract update required без effect. Не пытаться перевести незнакомые данные по имени поля.
- После fixture выбрать реальный serialization format, registry policy и integration test; их правила записать отдельно от учебной модели.
Граница знания и следующий тест
Fixture не читает Avro bytes, не проверяет JSON Schema, не общается с registry и не запускает Kubernetes, broker или database. Она не доказывает backwards compatibility продукта и не измеряет lag consumer. Она проверяет только конкретный контракт: v2 добавляет поле, v1 его не читает, v2 consumer умеет представить absence как null, v3 не вызывает effect. Если ваш consumer использует enum, money, locale или permission, это должны быть отдельные assertions, а не перенос нашей пары полей.
Следующий практичный шаг — добавить к изменению схемы review-таблицу из этой статьи и один replay test на выбранном хранилище результата. В хорошем результате будет видно event id, writer schema, consumer contract и исход effect. Если такой след не получается собрать без догадок, schema evolution пока рано выпускать: сначала надо сделать наблюдаемой границу reader и writer.
Проверяемые источники
- Apache Avro 1.10.1 Specification — Schema Resolution — официальная спецификация различает writer и reader schema и описывает resolution. Здесь она служит ориентиром для явного consumer contract, а не форматом payload или schema registry.
- CloudEvents Core Specification snapshot, 16 апреля 2021 (v1.0.2-wip) — исторический working draft существовал к июню 2021 года и различает context attributes, event data и protocol binding. Учебный envelope ниже не является CloudEvent и не реализует binding.
- Apache Kafka 2.7.0 KafkaProducer API — версионная документация линии 2.7 предупреждает, что retry может привести к duplicate. Fixture не запускает Kafka, не использует offset, partition, transaction или producer API.