Schema Evolution

Schema Evolution

 

Schema Evolution (эволюция схемы) это практика управляемого изменения схемы данных во времени, при которой уже работающие producer и consumer продолжают обмениваться сообщениями без остановки и без синхронного редеплоя. Механизм держится на сопоставлении writer schema, схемы, с которой данные записаны, и reader schema, схемы, с которой их читают, а разрешённость расхождений между ними определяет режим совместимости. Практика подробно разобрана применительно к Apache Kafka в статье блога школы «AVRO и JSON в Apache Kafka: краткий ликбез по реестру схем».

 

Что такое Schema Evolution и почему это управляемый процесс

Данные редко живут по неизменной схеме. Добавляется новое поле под свежую фичу, меняется тип поля под более точные расчёты, а удаляется то, что больше не нужно producer. Проблема не в самом изменении, а в том, что схему меняет одна команда, а читают эти данные десятки чужих сервисов, которые узнают о переименовании поля только когда их код падает на чтении. По определению документации Confluent, Schema Evolution это именно практика безопасного изменения схем во времени с сохранением совместимости с уже работающими producer и consumer, а не разовая миграция данных.

Ключевое отличие от разовой миграции в том, что эволюция схемы не требует одновременного обновления всех участников обмена. Producer и consumer обновляются каждый в своём темпе, а формальные правила совместимости заранее говорят, какое изменение это позволяет, а какое потребует синхронного деплоя всех клиентов.

 

Архитектура и ключевые особенности

У практики эволюции схемы в событийных системах есть три участника. Producer пишет данные, сериализуя их по своей версии схемы. Consumer читает данные, десериализуя их уже по своей версии, которая может отличаться от версии producer. Schema Registry (реестр схем) хранит все версии схем, присваивает каждой уникальный идентификатор и порядковый номер версии и проверяет новую версию на совместимость с предыдущими прежде чем принять её.

 

Writer schema, reader schema и реестр схем

Сопоставление двух версий схемы во время чтения называется schema resolution (разрешение схемы). Это не абстрактная идея, а конкретный алгоритм. Writer schema и reader schema сравниваются поле за полем по имени, а не по позиции, лишние поля writer’а отбрасываются, а поля reader’а без соответствия у writer’а заполняются default-значением, если оно задано в схеме читателя.

Сам реестр схем при этом не участвует в сериализации данных напрямую, его задача другая. Перед регистрацией новой версии схемы он прогоняет её через проверку совместимости с одной или несколькими предыдущими версиями subject’а и только при успехе присваивает новой версии идентификатор.

 

Режимы совместимости backward, forward, full и transitive

Формально разрешённость изменения задаёт режим совместимости, который реестр схем проверяет перед регистрацией новой версии.

  • Backward. Новая схема способна прочитать данные, записанные предыдущей версией. Обновлять первым нужно consumer.
  • Forward. Старая схема способна прочитать данные, записанные новой версией. Обновлять первым нужно producer.
  • Full. Совместимость работает в обе стороны одновременно, порядок обновления producer и consumer значения не имеет.
  • Transitive-варианты. Каждый из трёх режимов проверяется не только против последней версии схемы, а против всех предыдущих версий сразу, что снижает риск накопленной несовместимости в длинной цепочке изменений.

Ниже показано, кто из участников обмена обновляется первым при каждом режиме, и единственное изменение, безопасное при любом режиме, это поле, добавленное со значением по умолчанию.

Порядок обновления producer и consumer при режимах совместимости backward, forward и full, единственное всегда безопасное изменение схемы

Apache Kafka для инженеров данных

Код курса
DEVKI
Ближайшая дата курса
12 октября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Принцип работы schema resolution под капотом

Практика эволюции схемы становится конкретной там, где формат сериализации описывает точные правила сопоставления. У Avro эти правила зафиксированы спецификацией, у Schema Registry поверх них надстроены режимы совместимости.

 

Как Avro сопоставляет поля схем через type promotion и default-значения

Схемы совпадают, если обе описывают один и тот же примитивный тип, обе являются массивами или картами с совпадающими типами элементов, обе являются records, enum или fixed с одинаковыми именами, либо тип writer’а входит в фиксированный список повышаемых типов. Спецификация Avro 1.12.0 разрешает следующие повышения, int в long, float или double, long в float или double, float в double, а также string и bytes взаимно. Любое другое изменение типа поля, например long в int или string в int, автоматически не резолвится.

Поля сопоставляются по имени, а не по позиции. Лишнее поле writer schema читатель просто игнорирует. Поле reader schema, которого нет у writer’а, заполняется default-значением, если оно задано, а если default не задан, резолюция завершается ошибкой чтения. Для enum действует похожее правило, символ, которого нет у читателя, заменяется default-значением enum, если оно есть, иначе резолюция тоже отказывает.

Разрешение схемы writer и reader в Avro schema resolution, type promotion и заполнение default-значений

 

Как реестр схем проверяет совместимость перед принятием новой версии

Единица версионирования в реестре называется subject (субъект), обычно имя топика с суффиксом value или key, например orders-value. Уровень совместимости настраивается для каждого subject’а отдельно и действует глобально, пока его не переключат явным запросом. Реестр сравнивает кандидата новой версии с предыдущими версиями subject’а и, если проверка провалилась, отклоняет регистрацию, откатывать правку или менять режим приходится отдельно.

 

Эволюция схемы за пределами Avro

Механизм schema resolution специфичен для Avro, но сама практика эволюции схемы шире одного формата сериализации и одной событийной системы. Protobuf, JSON Schema и табличные форматы вроде Apache Iceberg и Delta Lake решают ту же задачу, совместимость схемы во времени, но разными средствами.

Технология Как задаётся схема Механизм эволюции
Apache Avro Schema attached к данным, writer schema Schema resolution поле за полем во время чтения, type promotion, default-значения
Confluent Schema Registry Централизованный реестр версий по subject Проверка совместимости backward/forward/full перед регистрацией новой версии
Protobuf Номера полей вместо имён Совместимость по стабильным номерам полей, переиспользовать номер удалённого поля нельзя
JSON Schema Декларативная валидация структуры Правила совместимости зависят от конкретного реестра, единого стандарта резолюции нет
Apache Iceberg, Delta Lake Метаданные таблицы, а не сами файлы данных Добавление, переименование и изменение типа колонки как операция над метаданными, без пересчёта уже записанных файлов

Общее для всех пяти строк таблицы правило простое, чем более формально описан контракт между писателем и читателем, тем меньше в эволюции схемы ручной дисциплины и тем больше её берёт на себя инструмент.

Архитектура Данных

Код курса
ARMG
Ближайшая дата курса
28 сентября, 2026
Продолжительность
24 ак.часов
Стоимость обучения
76 800

 

Ограничения и подводные камни

Формальный режим совместимости не отменяет несколько практических ограничений, с которыми демонстрация столкнулась напрямую.

  • Kafka Streams поддерживает только BACKWARD. Прямое ограничение из документации Confluent, режимы FORWARD и FULL для Kafka Streams недоступны.
  • Breaking change не лечится сменой режима на лету. Если изменение не укладывается ни в один режим совместимости, выход только один, либо синхронно обновить всех клиентов, либо завести новый топик под новую схему.
  • Переименование поля без alias это не переименование. Для schema resolution старое и новое имя это два разных поля. Значение из данных писателя под старым именем просто отбрасывается как незнакомое читателю, а новое поле получает собственный default, если он задан, вместо значения, которое реально лежало в записи.
  • Сужение типа иногда проходит там, где не ожидаешь. Проверка совместимости смотрит не на риск изменения в бытовом смысле, а на конкретное направление проверки, поэтому сужение double в int может пройти под FORWARD и провалиться под BACKWARD и FULL одновременно.

Общий вывод из ограничений такой, режим совместимости формализует конкретный контракт между producer и consumer, но не избавляет от необходимости думать о порядке деплоя и о том, что произойдёт с уже записанными данными.

 

Сценарии использования и отличительные черты

Выбор режима совместимости зависит от того, какая сторона обмена обновляется первой в конкретной инфраструктуре.

  • Backward. Подходит, когда consumer умеет перечитать топик с начала и должен работать со старыми данными, а обновляют сначала его. Это режим Confluent Schema Registry по умолчанию.
  • Forward. Подходит, когда сначала выкатывается producer, а consumer пока читает старой схемой и должен понимать новые данные.
  • Full. Нужен, когда producer и consumer обновляются в непредсказуемом порядке и нельзя заранее гарантировать, кто из них окажется новее.
  • Transitive-варианты любого из трёх режимов. Нужны, когда совместимость должна сохраняться не только с последней версией схемы, а сразу со всеми предыдущими, типичный случай для долгоживущих топиков с длинной историей изменений.

Есть и обратный случай, когда механизм эволюции схемы не нужен вовсе. В закрытой системе с одним producer и одним consumer в общем репозитории, где деплой обоих компонентов синхронизирован одним пайплайном, реестр схем и режимы совместимости добавляют инфраструктуры больше, чем пользы, схему в такой системе проще менять напрямую вместе с кодом. Разобрать механизм Avro и режимы совместимости руками, на живом кластере, можно на курсе «Основы Apache Kafka».

 

Практика

Демонстрация состоит из двух независимых примеров. Первый показывает чистый механизм Avro schema resolution без Docker и без реестра, средствами библиотеки fastavro 1.12.2. Второй прогоняет матрицу режимов совместимости против настоящего Confluent Schema Registry 8.3.1, реестр и брокер Apache Kafka 4.3.0 подняты в Docker.

Schema Evolution from airflow import DAG from airflow.operators.bash import BashOperator from datetime import datetime with DAG( dag_id="spark_submit_demo", start_date=datetime(2025, 1, 1), schedule="@daily", catchup=False ) as dag: run = BashOperator( task_id="run_job", bash_command="spark-submit app.py" ) GitHub code example Schema Evolution

Первое демо, schema_resolution_avro.py, читает данные, записанные writer schema v1, через две разные reader schema, чтобы показать сразу три сценария, совместимое добавление поля, удаление обязательного поля без default и переименование поля без alias.

# fastavro 1.12.2, локально без Docker и реестра схем, прогнано на стенде 2026-09-04
#
# Механизм Avro schema resolution: писатель и читатель используют РАЗНЫЕ схемы одной и
# той же записи, и fastavro сопоставляет их поле за полем по имени, а не по позиции.
# Три сценария: совместимое добавление поля, удаление поля, переименование без alias.

import io

import fastavro
from fastavro.read import SchemaResolutionError

# writer schema v1: заказ с идентификатором, суммой и статусом
writer_schema_v1 = {
    "type": "record",
    "name": "Order",
    "fields": [
        {"name": "id", "type": "string"},
        {"name": "amount", "type": "double"},
        {"name": "status", "type": "string"},
    ],
}

# reader schema v2: добавлено поле currency с default, удалено поле status
reader_schema_v2 = {
    "type": "record",
    "name": "Order",
    "fields": [
        {"name": "id", "type": "string"},
        {"name": "amount", "type": "double"},
        {"name": "currency", "type": "string", "default": "RUB"},
    ],
}

# reader schema v3: то же самое поле currency, но переименовано в curr без alias —
# для резолюции это не переименование, а новое неизвестное поле читателя
reader_schema_v3_renamed_no_alias = {
    "type": "record",
    "name": "Order",
    "fields": [
        {"name": "id", "type": "string"},
        {"name": "amount", "type": "double"},
        {"name": "curr", "type": "string", "default": "RUB"},
    ],
}


def write_with_schema(record: dict, schema: dict) -> bytes:
    """Пишет одну запись в бинарный Avro писателем указанной схемы."""
    buf = io.BytesIO()
    fastavro.schemaless_writer(buf, schema, record)
    return buf.getvalue()


def read_with_schema(data: bytes, writer_schema: dict, reader_schema: dict) -> dict:
    """Читает байты, записанные writer_schema, через reader_schema — это и есть
    schema resolution: fastavro сопоставляет поля по имени между двумя схемами."""
    buf = io.BytesIO(data)
    return fastavro.schemaless_reader(buf, writer_schema, reader_schema)


def scenario_field_added_with_default():
    print("--- Сценарий 1: writer v1 -> reader v2, поле currency добавлено с default ---")
    record = {"id": "A-1001", "amount": 149.5, "status": "paid"}
    data = write_with_schema(record, writer_schema_v1)
    print(f"записано writer v1: {record}")
    print(f"байт на диске: {len(data)}")
    result = read_with_schema(data, writer_schema_v1, reader_schema_v2)
    print(f"прочитано reader v2: {result}")
    print("поле status отсутствует в reader v2 и молча отброшено")
    print("поле currency отсутствует у writer и заполнено default-значением")


def scenario_removed_field_breaks_old_reader():
    print()
    print("--- Сценарий 2: writer v2 (без status) -> reader v1 (status обязателен, без default) ---")
    record = {"id": "A-1002", "amount": 89.0, "currency": "USD"}
    data = write_with_schema(record, reader_schema_v2)
    print(f"записано writer v2: {record}")
    try:
        result = read_with_schema(data, reader_schema_v2, writer_schema_v1)
        print(f"прочитано reader v1: {result}")
    except SchemaResolutionError as exc:
        print(f"резолюция отказала: {exc}")
        print("это forward-совместимость (старый читатель против новых данных), а не")
        print("backward: поле status удалено в v2 без default в v1, и читателю, который")
        print("ещё не обновился на v2, взять значение status неоткуда")


def scenario_rename_without_alias_fails():
    print()
    print("--- Сценарий 3: переименование currency -> curr БЕЗ alias ---")
    record = {"id": "A-1003", "amount": 12.3, "currency": "EUR"}
    data = write_with_schema(record, reader_schema_v2)
    print(f"записано writer v2: {record}")
    try:
        result = read_with_schema(data, reader_schema_v2, reader_schema_v3_renamed_no_alias)
        print(f"прочитано reader v3: {result}")
        print("резолюция НЕ считает curr тем же полем, что currency: поле writer'а с")
        print("именем currency проигнорировано как незнакомое читателю, curr в результате")
        print("получил собственное default-значение, а не значение из данных")
        assert result["curr"] == "RUB", "curr должен взять default, а не значение currency"
        assert "currency" not in result, "currency не должно попасть в результат reader v3"
    except SchemaResolutionError as exc:
        print(f"резолюция отказала: {exc}")


if __name__ == "__main__":
    scenario_field_added_with_default()
    scenario_removed_field_breaks_old_reader()
    scenario_rename_without_alias_fails()

Вывод прогона показывает все три сценария на реальных данных.

--- Сценарий 1: writer v1 -> reader v2, поле currency добавлено с default ---
записано writer v1: {'id': 'A-1001', 'amount': 149.5, 'status': 'paid'}
байт на диске: 20
прочитано reader v2: {'id': 'A-1001', 'amount': 149.5, 'currency': 'RUB'}
поле status отсутствует в reader v2 и молча отброшено
поле currency отсутствует у writer и заполнено default-значением

--- Сценарий 2: writer v2 (без status) -> reader v1 (status обязателен, без default) ---
записано writer v2: {'id': 'A-1002', 'amount': 89.0, 'currency': 'USD'}
резолюция отказала: No default value for field status in Order
это forward-совместимость (старый читатель против новых данных), а не
backward: поле status удалено в v2 без default в v1, и читателю, который
ещё не обновился на v2, взять значение status неоткуда

--- Сценарий 3: переименование currency -> curr БЕЗ alias ---
записано writer v2: {'id': 'A-1003', 'amount': 12.3, 'currency': 'EUR'}
прочитано reader v3: {'id': 'A-1003', 'amount': 12.3, 'curr': 'RUB'}
резолюция НЕ считает curr тем же полем, что currency: поле writer'а с
именем currency проигнорировано как незнакомое читателю, curr в результате
получил собственное default-значение, а не значение из данных

Во втором сценарии резолюция реально отказывает с ошибкой No default value for field status in Order, это не сбой демонстрации, а прямая иллюстрация forward-совместимости, старый читатель, ещё не обновившийся на v2, не может взять значение status оттуда, где его больше нет. В третьем сценарии поле curr получает default-значение RUB, хотя в данных под именем currency лежало EUR, резолюция потеряла это значение молча, без единой ошибки, потому что переименование без alias для неё выглядит как одно поле исчезло, а другое появилось из ниоткуда.

Второе демо, compatibility_matrix.py, регистрирует схему заказа в Schema Registry и прогоняет шесть типовых изменений через эндпоинт совместимости под каждым из трёх режимов, переключая режим subject’а перед каждой проверкой отдельным запросом.

# confluentinc/cp-schema-registry 8.3.1, Docker образ, прогнано на стенде 2026-09-04
# REST API v1, requests 2.34.2
#
# Регистрирует базовую схему заказа, затем прогоняет набор типовых изменений через
# эндпоинт /compatibility под каждым из трёх режимов (BACKWARD, FORWARD, FULL) и печатает
# матрицу pass/fail. Реестр меняет режим глобально на время своей проверки — это ограничение
# самого REST API v1, отдельного per-request параметра режима нет.

import json

import requests

SR = "http://localhost:8081"
CT = {"Content-Type": "application/vnd.schemaregistry.v1+json"}
SUBJECT = "schema_evolution_order-value"

# базовая схема: заказ с идентификатором и суммой
schema_v1 = {
    "type": "record",
    "name": "Order",
    "fields": [
        {"name": "id", "type": "string"},
        {"name": "amount", "type": "double"},
    ],
}

# кандидаты на изменение схемы: имя, схема-кандидат
CHANGES = [
    (
        "поле добавлено с default",
        {
            "type": "record",
            "name": "Order",
            "fields": [
                {"name": "id", "type": "string"},
                {"name": "amount", "type": "double"},
                {"name": "currency", "type": "string", "default": "RUB"},
            ],
        },
    ),
    (
        "поле добавлено без default",
        {
            "type": "record",
            "name": "Order",
            "fields": [
                {"name": "id", "type": "string"},
                {"name": "amount", "type": "double"},
                {"name": "currency", "type": "string"},
            ],
        },
    ),
    (
        "поле удалено",
        {
            "type": "record",
            "name": "Order",
            "fields": [
                {"name": "id", "type": "string"},
            ],
        },
    ),
    (
        "поле переименовано без alias",
        {
            "type": "record",
            "name": "Order",
            "fields": [
                {"name": "id", "type": "string"},
                {"name": "total", "type": "double"},
            ],
        },
    ),
    (
        "тип сужен double -> int",
        {
            "type": "record",
            "name": "Order",
            "fields": [
                {"name": "id", "type": "string"},
                {"name": "amount", "type": "int"},
            ],
        },
    ),
    (
        "тип расширен double -> string",
        {
            "type": "record",
            "name": "Order",
            "fields": [
                {"name": "id", "type": "string"},
                {"name": "amount", "type": "string"},
            ],
        },
    ),
]

MODES = ["BACKWARD", "FORWARD", "FULL"]


def reset_subject():
    """Удаляет subject от прошлого прогона: мягкое удаление, затем окончательное."""
    requests.delete(f"{SR}/subjects/{SUBJECT}")
    requests.delete(f"{SR}/subjects/{SUBJECT}?permanent=true")


def register_v1():
    resp = requests.post(
        f"{SR}/subjects/{SUBJECT}/versions",
        headers=CT,
        data=json.dumps({"schemaType": "AVRO", "schema": json.dumps(schema_v1)}),
    )
    resp.raise_for_status()
    print(f"зарегистрирована схема v1: {resp.json()}")


def set_subject_mode(mode: str):
    resp = requests.put(
        f"{SR}/config/{SUBJECT}",
        headers=CT,
        data=json.dumps({"compatibility": mode}),
    )
    resp.raise_for_status()


def check_compatibility(candidate_schema: dict) -> bool:
    resp = requests.post(
        f"{SR}/compatibility/subjects/{SUBJECT}/versions/latest",
        headers=CT,
        data=json.dumps({"schemaType": "AVRO", "schema": json.dumps(candidate_schema)}),
    )
    resp.raise_for_status()
    return resp.json()["is_compatible"]


def build_matrix():
    reset_subject()
    register_v1()

    matrix = {}
    for mode in MODES:
        set_subject_mode(mode)
        print(f"\nрежим subject выставлен: {mode}")
        matrix[mode] = {}
        for name, candidate in CHANGES:
            ok = check_compatibility(candidate)
            matrix[mode][name] = ok
            print(f"  {mode:8} | {name:32} | {'PASS' if ok else 'FAIL'}")

    return matrix


def print_table(matrix: dict):
    print("\n=== Матрица совместимости ===")
    name_width = max(len(name) for name, _ in CHANGES)
    header = "изменение".ljust(name_width) + "".join(f" | {m:8}" for m in MODES)
    print(header)
    print("-" * len(header))
    for name, _ in CHANGES:
        row = name.ljust(name_width)
        for mode in MODES:
            row += f" | {'PASS' if matrix[mode][name] else 'FAIL':8}"
        print(row)


if __name__ == "__main__":
    result = build_matrix()
    print_table(result)

Матрица pass/fail на реальном прогоне выглядит так.

зарегистрирована схема v1: {'id': 2, 'version': 1, 'guid': 'cc7aa2de-b4c7-2886-70fb-b4766ba171a7', 'schemaType': 'AVRO', 'schema': '{"type":"record","name":"Order","fields":[{"name":"id","type":"string"},{"name":"amount","type":"double"}]}'}

режим subject выставлен: BACKWARD
  BACKWARD | поле добавлено с default         | PASS
  BACKWARD | поле добавлено без default       | FAIL
  BACKWARD | поле удалено                     | PASS
  BACKWARD | поле переименовано без alias     | FAIL
  BACKWARD | тип сужен double -> int          | FAIL
  BACKWARD | тип расширен double -> string    | FAIL

режим subject выставлен: FORWARD
  FORWARD  | поле добавлено с default         | PASS
  FORWARD  | поле добавлено без default       | PASS
  FORWARD  | поле удалено                     | FAIL
  FORWARD  | поле переименовано без alias     | FAIL
  FORWARD  | тип сужен double -> int          | PASS
  FORWARD  | тип расширен double -> string    | FAIL

режим subject выставлен: FULL
  FULL     | поле добавлено с default         | PASS
  FULL     | поле добавлено без default       | FAIL
  FULL     | поле удалено                     | FAIL
  FULL     | поле переименовано без alias     | FAIL
  FULL     | тип сужен double -> int          | FAIL
  FULL     | тип расширен double -> string    | FAIL

=== Матрица совместимости ===
изменение                     | BACKWARD | FORWARD  | FULL
--------------------------------------------------------------
поле добавлено с default      | PASS     | PASS     | PASS
поле добавлено без default    | FAIL     | PASS     | FAIL
поле удалено                  | PASS     | FAIL     | FAIL
поле переименовано без alias  | FAIL     | FAIL     | FAIL
тип сужен double -> int       | FAIL     | PASS     | FAIL
тип расширен double -> string | FAIL     | FAIL     | FAIL

Единственное изменение, прошедшее под всеми тремя режимами, это добавление поля с default-значением. Контринтуитивный результат в этой матрице, сужение типа double в int, прошло под FORWARD и провалилось под BACKWARD и FULL, что на первый взгляд странно, ведь сужение обычно считают более рискованным изменением, но соответствует направлению самой проверки FORWARD, старый читатель против новых данных, и правилам type promotion Avro, которые разрешают повышать int до double, а не наоборот.

Шаги запуска обоих демо и типовые ошибки, например реестр схем, не успевший подняться после старта Kafka, или занятые порты 9092 и 8081, разобраны в RUNBOOK.md рядом с демо-кодом.

 

Заключение

Schema Evolution работает потому, что формализует то, что раньше держалось в головах разработчиков, конкретный алгоритм сопоставления полей у Avro и явные режимы совместимости у Schema Registry. Демонстрация на живом стенде подтвердила это цифрами, поле с default безопасно везде, удаление поля без default ломает ещё не обновившихся читателей, а переименование без alias теряет данные молча. Для команд, которые выкатывают producer и consumer независимо, эти правила экономят инцидент в продакшене, для закрытой системы с синхронным деплоем это лишняя инфраструктура.

 

Референсные ссылки