Schema Registry

Schema Registry

 

Schema Registry (реестр схем) это централизованный сервис хранения, версионирования и контроля совместимости схем сообщений, которые ходят между producer и consumer. Чаще всего он работает в связке с Apache Kafka и форматами Avro, JSON Schema или Protocol Buffers. Задача реестра простая на словах и неприятная на практике. Нужно сделать так, чтобы изменение структуры сообщения одной командой не уронило чужие приложения, которые эти сообщения читают.

Кафка сама по себе про схемы ничего не знает. Для брокера сообщение это массив байтов. Он одинаково спокойно пропустит и корректную запись, и мусор. Разбираться, что за байты приехали, приходится потребителю. Реестр схем закрывает ровно этот пробел: он становится общим справочником форматов для всех участников обмена.

 

Что такое Schema Registry и какую задачу он решает

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

Без реестра схема живёт в головах разработчиков и в коде двух десятков сервисов. Кто-то переименовал поле, кто-то поменял тип с int на string. Сломается это не в момент коммита, а через неделю в 3 часа ночи на стороне потребителя. Реестр переносит проверку влево: несовместимая схема просто не регистрируется, и producer падает на старте, а не отравляет топик.

Второй эффект менее очевидный, но для больших топиков считается деньгами. Схема не едет в каждом сообщении, вместо неё в payload попадает только числовой идентификатор. Для потока из миллиардов мелких событий экономия на трафике и на диске выходит заметная.

 

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

Confluent Schema Registry это отдельное stateless-приложение на JVM, которое ставится рядом с кластером Kafka. Собственной базы данных у него нет. Все схемы он пишет в специальный compacted-топик, по умолчанию _schemas, и держит их копию в памяти. Отсюда пара важных следствий. Реестр переживает перезапуск без потерь, потому что состояние восстанавливается из топика. И он масштабируется горизонтально: несколько экземпляров работают на одном топике, но записывать новые схемы имеет право только лидер.

 

Subject и стратегии именования

Схемы группируются в сущности, которые называются subject. Subject это именованная линия версий: внутри него схема живёт как v1, v2, v3, и именно к subject привязано правило совместимости. Как формируется имя, определяет стратегия именования (subject naming strategy).

  • TopicNameStrategy. Значение по умолчанию, имя складывается из названия топика и суффикса, например orders-value и orders-key. Один топик тянет одну схему значения.
  • RecordNameStrategy. Имя берётся из полного имени типа записи. Позволяет держать в одном топике события разных типов и версионировать каждый тип отдельно.
  • TopicRecordNameStrategy. Комбинация двух предыдущих: топик плюс имя типа. Разные топики с одинаковым типом события эволюционируют независимо друг от друга.

Выбор стратегии это архитектурное решение, а не настройка на потом. Если в топик планируется класть несколько типов событий, стратегию по умолчанию придётся менять сразу. Иначе проверка совместимости начнёт сравнивать между собой заведомо разные схемы и отвергать нормальные изменения.

 

REST API как единая точка входа

Всё управление реестром сводится к HTTP-запросам: зарегистрировать схему, получить её по идентификатору, посмотреть список версий, проверить совместимость, поменять уровень проверки для конкретного subject. Тот же API используют CI-пайплайны, чтобы прогонять проверку схемы до выкладки сервиса, а не в рантайме.

 

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

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

 

Как это работает под капотом

Ключ к пониманию реестра лежит в формате байтов, который сериализатор кладёт в топик. Он называется wire format и устроен так. Первый байт всегда нулевой, это магический байт-индикатор. Следом идут 4 байта с идентификатором схемы, знаковое 32-битное целое в порядке big-endian. Дальше начинается обычное бинарное представление самих данных. У Protobuf между идентификатором и полезной нагрузкой добавляются индексы сообщения, потому что в одном .proto-файле лежит несколько типов и нужно указать конкретный.

# структура по документации Confluent Platform 8.3

[0x00][schema_id: 4 байта, big-endian][полезная нагрузка]
  |            |                        |
  |            |                        L- Avro/JSON/Protobuf в бинарном виде
  |            L- идентификатор схемы из реестра
  L- магический байт, признак wire format

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

Schema Registry: путь сообщения от Producer через сериализатор и Kafka к Consumer, регистрация схемы и получение schema ID в реестре, извлечение схемы десериализатором и хранение состояния в служебном топике _schemas

 

В Confluent Platform 8.3 у этой схемы появился альтернативный вариант. Идентификатор можно передавать в заголовке записи Kafka вместо пятибайтового префикса в payload. Полезно там, где данные читают сторонние инструменты, не знающие про формат Confluent, и лишние байты в начале ломают им разбор.

 

Правила совместимости, ради которых всё затевалось

Хранение схем это половина дела. Ценность реестра в том, что он отказывается принимать изменение, которое сломает читателей. Уровень строгости задаётся типом совместимости, глобально или отдельно для каждого subject.

Тип Что гарантирует Что разрешено менять Кого обновлять первым
BACKWARD (по умолчанию) Новая схема читает данные, записанные предыдущей версией Удалять поля, добавлять опциональные с default Consumer
BACKWARD_TRANSITIVE То же, но против всех прошлых версий, а не только последней То же самое Consumer
FORWARD Старая схема читает данные, записанные новой версией Добавлять поля, удалять опциональные с default Producer
FORWARD_TRANSITIVE То же против всех прошлых версий То же самое Producer
FULL Совместимость в обе стороны с соседней версией Только опциональные поля с default Порядок не важен
FULL_TRANSITIVE В обе стороны со всеми версиями Только опциональные поля с default Порядок не важен
NONE Ничего, проверка отключена Что угодно Молитва

Практический вывод из таблицы такой: тип совместимости диктует порядок деплоя. При BACKWARD сначала выкатывают потребителей, при FORWARD сначала поставщиков. Перепутать эти два сценария означает получить именно тот инцидент, от которого защищались. Подробный разбор поведения клиентов при таком обмене есть в статье Как реестр схем помогает снизить нагрузку на запись сообщений в топики Apache Kafka.

 

Data Contracts, когда структуры мало

Совместимость проверяет форму записи, но не её содержание. Начиная с версии 7.4 в реестре появились контракты данных (data contracts), надстройка, которая довешивает к схеме метаданные, теги и правила. Правила бывают двух сортов. Доменные правила (domain rules) описываются выражениями на языке CEL и проверяют значения полей прямо в момент сериализации. Правила миграции (migration rules) пишутся на JSONata и превращают запись одной версии в запись другой при переходах UPGRADE, DOWNGRADE или UPDOWN.

Оговорка тут существенная. Контракты данных доступны в Confluent Platform Enterprise и в Confluent Cloud с пакетом Advanced Stream Governance. В открытой сборке реестра их нет, поэтому смысловую валидацию там придётся закрывать своими средствами.

Apache Kafka: администрирование кластера

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

 

Проверка совместимости на живом стенде

Теория про BACKWARD выглядит убедительно ровно до первого прогона. Поднимем брокер и реестр в Docker и посмотрим, что реестр отвечает на самом деле. Сценарий короткий: регистрируем схему заказа, пробуем добавить поле со значением по умолчанию и без него.

По традиции весь код используемый в статье выкладываем на наш GitHub репозиторий
Schema Registry 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 Registry

Скрипт compat_check.sh начинается с удаления subject, поэтому его можно гонять сколько угодно раз подряд.

#!/usr/bin/env bash
# прогон на Confluent Schema Registry 8.3.1, REST API v1
SR=http://localhost:8081
CT='Content-Type: application/vnd.schemaregistry.v1+json'

# 0. чистим subject от прошлого прогона: мягкое удаление, затем окончательное
curl -s -X DELETE $SR/subjects/orders-value > /dev/null
curl -s -X DELETE "$SR/subjects/orders-value?permanent=true" > /dev/null

# 1. регистрируем первую версию схемы заказа: идентификатор и сумма
curl -s -X POST -H "$CT" --data '{"schemaType":"AVRO","schema":"{\"type\":\"record\",\"name\":\"Order\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"amount\",\"type\":\"double\"}]}"}' \
  $SR/subjects/orders-value/versions
echo " # id схемы v1"

# 2. проверяем вторую версию ДО регистрации: добавлено поле с default
curl -s -X POST -H "$CT" --data '{"schemaType":"AVRO","schema":"{\"type\":\"record\",\"name\":\"Order\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"amount\",\"type\":\"double\"},{\"name\":\"currency\",\"type\":\"string\",\"default\":\"RUB\"}]}"}' \
  $SR/compatibility/subjects/orders-value/versions/latest
echo " # поле с default"

# 3. то же поле, но без default: реестр обязан отказать
curl -s -X POST -H "$CT" --data '{"schemaType":"AVRO","schema":"{\"type\":\"record\",\"name\":\"Order\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"amount\",\"type\":\"double\"},{\"name\":\"currency\",\"type\":\"string\"}]}"}' \
  $SR/compatibility/subjects/orders-value/versions/latest
echo " # поле без default"

# 4. регистрируем совместимую версию, теперь в subject две версии
curl -s -X POST -H "$CT" --data '{"schemaType":"AVRO","schema":"{\"type\":\"record\",\"name\":\"Order\",\"fields\":[{\"name\":\"id\",\"type\":\"string\"},{\"name\":\"amount\",\"type\":\"double\"},{\"name\":\"currency\",\"type\":\"string\",\"default\":\"RUB\"}]}"}' \
  $SR/subjects/orders-value/versions
echo " # id схемы v2"

# 5. текущий уровень проверки для subject и глобальный по умолчанию
curl -s $SR/config; echo " # глобальный уровень"

Вывод прогона на Confluent Schema Registry 8.3.1 приведён ниже, длинные тела схем сокращены.

{"version":"8.3.1","commitId":"db21878bbe09d5cfe0a4b757963e062984f34475"}
{"id":3,"version":1,"guid":"cc7aa2de-b4c7-2886-70fb-b4766ba171a7",...}  # id схемы v1
{"is_compatible":true}                                                 # поле с default
{"is_compatible":false}                                                # поле без default
{"id":4,"version":2,"guid":"3737a14e-229a-684f-52f3-48767e17a76b",...}  # id схемы v2
{"compatibilityLevel":"BACKWARD"}                                      # глобальный уровень

Два момента в этом выводе стоит отметить отдельно. Рядом с числовым id приезжает guid, глобальный идентификатор схемы: он считается от её содержимого и не зависит от конкретного реестра. А номера схем начинаются не с единицы, потому что реестр не переиспользует идентификаторы окончательно удалённых схем. Счётчик просто едет дальше.

Второй скрипт показывает тот самый префикс живьём. Он сериализует запись в формате AVRO, печатает первые байты и читает сообщение обратно из топика.

# прогон на confluent-kafka 2.15.0, Apache Kafka 4.3.0, Schema Registry 8.3.1
from confluent_kafka import Producer, Consumer
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer, AvroDeserializer
from confluent_kafka.serialization import SerializationContext, MessageField

TOPIC = "orders"
SCHEMA = """
{"type":"record","name":"Order","fields":[
 {"name":"id","type":"string"},
 {"name":"amount","type":"double"},
 {"name":"currency","type":"string","default":"RUB"}]}
"""

# клиент реестра: через него сериализатор регистрирует схему и получает её идентификатор
sr = SchemaRegistryClient({"url": "http://localhost:8081"})
serializer = AvroSerializer(sr, SCHEMA)
ctx = SerializationContext(TOPIC, MessageField.VALUE)

# сериализация превращает словарь в байты: 5 байт префикса плюс тело Avro
payload = serializer({"id": "A-1001", "amount": 149.5, "currency": "RUB"}, ctx)
print("байт всего:", len(payload))
print("префикс:", payload[:5].hex(" "))
print("магический байт:", payload[0], "| schema id:", int.from_bytes(payload[1:5], "big"))

producer = Producer({"bootstrap.servers": "localhost:9092"})
producer.produce(TOPIC, value=payload)
producer.flush()

# потребитель читает сырые байты, десериализатор сам ходит в реестр за схемой по id
consumer = Consumer({"bootstrap.servers": "localhost:9092",
                     "group.id": "wiki_demo", "auto.offset.reset": "earliest"})
consumer.subscribe([TOPIC])
message = consumer.poll(15)
deserializer = AvroDeserializer(sr, SCHEMA)
print("прочитано:", deserializer(message.value(), ctx))
consumer.close()
байт всего: 24
префикс: 00 00 00 00 04
магический байт: 0 | schema id: 4
прочитано: {'id': 'A-1001', 'amount': 149.5, 'currency': 'RUB'}

Двадцать четыре байта на всю запись, из них пять служебных. Сама схема при этом никуда не поехала. Потребитель забрал её у реестра по номеру 4 и положил в кэш.

 

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

Реестр схем не бесплатен в эксплуатации, и большая часть проблем с ним типовая.

  • Соблазн выставить NONE. Самый быстрый способ пройти проверку это её выключить. После такого реестр превращается в архив схем без единой гарантии, и потребители ломаются ровно так, как ломались бы без него.
  • Забытый default. В Avro поле без значения по умолчанию нельзя добавить обратно совместимо. Половина отказов регистрации это именно оно, а не хитрая эволюция типов.
  • Единая точка отказа на старте. Кэш спасает работающие приложения, но холодный старт потребителя при недоступном реестре заканчивается ошибкой десериализации. Реестр нужно резервировать так же, как брокеры.
  • Мусор в subject. Стратегия по умолчанию плодит по два subject на топик, и через год в реестре сотни линий версий, половина от удалённых топиков. Уборка не автоматическая, это ручной процесс.
  • Схема не проверяет смысл. Реестр гарантирует структуру, но не содержание: поле amount типа double спокойно примет отрицательную сумму. Нужны контракты данных или своя валидация.

Общий принцип тут такой. Реестр это инструмент дисциплины, а не автопилот. Он работает ровно настолько, насколько команда согласилась соблюдать выбранный уровень совместимости.

 

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

Schema Registry нужен там, где данные пересекают границу команды: несколько независимых сервисов пишут и читают одни и те же топики, схемы меняются регулярно, а договориться устно уже не выходит. Если у топика один поставщик и один потребитель в том же репозитории, реестр добавит инфраструктуры больше, чем пользы. Реализаций несколько, и различаются они в основном моделью развёртывания.

Реализация Особенности Когда брать
Confluent Schema Registry Де-факто стандарт, Avro, JSON Schema, Protobuf, хранение в топике Kafka Классический стек Kafka, on-premise или гибрид
Apicurio Registry Открытая реализация, помимо схем хранит OpenAPI и AsyncAPI, группировка артефактов Экосистема Red Hat, единый каталог схем и API-контрактов
Karapace Реализация на Python от Aiven, совместимая по API, есть backup и restore Kubernetes, потребность в открытой лицензии
AWS Glue Schema Registry Serverless, интеграция с MSK, Kinesis и IAM Инфраструктура целиком в AWS

API Confluent стал общим знаменателем. Karapace и Redpanda реализуют его совместимо, поэтому клиентский код при переезде обычно не переписывают. Разобрать сериализацию Avro и работу с реестром руками, на живом кластере, можно на курсе Apache Kafka: обучение для разработчиков.

 

Заключение

Schema Registry превращает неявную договорённость о формате сообщений в проверяемый контракт с версиями и правилами изменения. Технически он устроен несложно: топик со схемами, REST API, пятибайтовый префикс в сообщении и кэш на клиенте. Вся ценность концентрируется в правилах совместимости и в готовности команды им следовать, потому что реестр с отключённой проверкой не защищает ни от чего.

 

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