Концепции#
Corax Schema Registry — это управляемый репозиторий схем, который поддерживает хранение и обмен данными как для потоковой обработки, так и для данных в состоянии покоя, таких как базы данных, файлы и другие статические хранилища данных. В этом разделе объясняются некоторые фундаментальные понятия, связанные со схемами, реестром и утилитами, а также рассматриваются детали того, как все это работает с точки зрения внутреннего устройства.
Как это работает#
Schema Registry:
предоставляет сервисный слой для метаданных;
предоставляет REST-интерфейс для хранения и извлечения схем Avro и JSON;
хранит версионную историю всех схем на основе заданной стратегии имен субъектов;
предоставляет несколько настроек совместимости и позволяет обновлять схемы в соответствии с настроенными параметрами совместимости и расширенной поддержкой этих типов схем;
предоставляет сериализаторы, подключаемые к клиентам Kafka, которые обрабатывают хранение и извлечение схем для сообщений Kafka, отправляемых в любом из поддерживаемых форматов.
Schema Registry расположен отдельно от брокеров Kafka и может быть развернут на сервере без Corax. Производители и потребители по-прежнему обращаются к брокерам Kafka для публикации и чтения данных (сообщений) в топиках. Одновременно с этим они могут обращаться к Schema Registry для отправки и получения схем, описывающих модели данных для сообщений.
Schema Registry — это распределенный слой хранения схем, который использует Kafka в качестве базового механизма хранения. Некоторые ключевые проектные решения:
Присваивает уникальный ID каждой зарегистрированной схеме. Сервер Schema Registry кеширует зарегистрированные схемы. При добавлении новой схемы ID формируется следующим образом: из кеша зарегистрированных схем определяется максимальный ID, затем к нему прибавляется
1.Kafka обеспечивает надежную работу серверной части. Функционирует как лог предзаписи изменений состояния Schema Registry и содержащихся в нем.
Schema Registry спроектирован как распределенное приложение с одним основным сервером, при этом ZooKeeper/Kafka координирует выборы основного сервера (в зависимости от конфигурации).
Схемы, субъекты и топики#
Основные понятия, используемые в Schema Registry:
Топик Kafka — способ группировки потоков сообщений по категориям. Каждое сообщение представляет собой пару «ключ-значение». Ключ, значение или пара «ключ-значение» могут быть сериализованы в формат AVRO или JSON.
Схема — стандарт для определения структуры формата данных. Имя топика может отличатся от имени схемы.
Субъект (subject) — это уникальное имя, под которым регистрируются схемы. Имя субъекта зависит от настроенной стратегии имени субъекта, которая по умолчанию настроена на получение имени субъекта из имени топика. Подробнее о стратегии имени субъекта смотрите в разделе .
Изменение стратегии имени субъекта#
Стратегию имени субъекта можно изменять для каждого конкретного топика.
Реализация интерфейса SubjectNameStrategy генерирует 2 вида идентификаторов — subject и version.
Если стратегия имен по умолчанию не подходит, реализуйте этот интерфейс самостоятельно:
subject— имя субъекта по имени топика и признаку isKey. По умолчанию дляvalueимя субъекта совпадает с именем топика, в конце добавляется-value. Пример:accounts-value.Для ключей добавляется
-key. Пример:accounts-key.version— ключ, в котором будет лежать номер версии схемы. По умолчанию[subject].version. Пример:accounts.version=2— использовать вторую версию схемыaccounts;accounts-key.version=3— использовать третью версию схемы для ключей топикаaccounts.
Сериализаторы и десериализаторы Kafka#
При передаче данных по сети или хранении их в файле нужен способ закодировать данные в байты. Раньше для каждого языка программирования использовался свой сериализатор, например сериализатор Java. Так как такой специфический формат был неудобным, был совершен переход на универсальный формат — JSON, но без строго определенной схемы формата.
Однако отсутствие строго определенного формата имеет два существенных недостатка:
Потребители данных могут не понимать производителей данных: отсутствие структуры делает потребление данных в таких форматах более сложным, поскольку поля могут быть произвольно добавлены или удалены, а данные могут быть повреждены. Этот недостаток становится тем серьезнее, чем больше приложений или команд в организации начинают потреблять поток данных: если команда, находящаяся выше по течению, может вносить произвольные изменения в формат данных по своему усмотрению, то становится очень сложно гарантировать, что все потребители, находящиеся ниже по течению, смогут и в дальнейшем интерпретировать данные. Не хватает контракта для данных между производителями и потребителями, аналогичного контракту API.
Накладные расходы и увеличение размера: универсальный формат обмена данными имеет больший размер, поскольку имена полей и информация о типе должны быть явно представлены в сериализованном формате, несмотря на то, что они идентичны во всех сообщениях.
Дальнейшее развитие форматов сериализации привело к появлению межъязыковых библиотек сериализации которые требуют формального определения структуры данных с помощью схем. К таким библиотекам относятся Avro, Thrift, Protocol Buffers и JSON Schema. Преимущество наличия схемы заключается в том, что она четко определяет структуру, тип и значение (через документацию) данных. Кроме того, с помощью схемы можно более эффективно кодировать данные.
Например, схема Avro определяет структуру данных в формате JSON. Следующая схема Avro определяет запись пользователя с двумя полями: name и favorite_number типа string и int соответственно.
{"namespace": "example.avro",
"type": "record",
"name": "user",
"fields": [
{"name": "name", "type": "string"},
{"name": "favorite_number", "type": "int"}
]
}
Можно использовать эту схему Avro, например, для сериализации Java-объекта POJO в байты и десериализации этих байтов обратно в Java-объект.
Avro требует схему не только при сериализации данных, но и при их десериализации. Поскольку схема предоставляется во время декодирования, метаданные, такие как имена полей, не нужно явно кодировать в данных. Это делает двоичное кодирование данных Avro очень компактным.
Поддерживаемые форматы Avro, JSON#
Avro был выбран в качестве формата схемы, поддерживаемого по умолчанию в Corax. Для формата Avro предусмотрены сериализаторы и десериализаторы Kafka.
Corax поддерживает форматы JSON Schema и Avro. Поддержка этих форматов сериализации не ограничивается Schema Registry, а обеспечивается во всем Corax.
Новые сериализаторы и десериализаторы Kafka доступны для JSON Schema и Avro. Сериализаторы могут автоматически регистрировать схемы при сериализации сообщения JSON-сериализуемого объекта.
Сериализаторы и десериализаторы доступны только на Java.
Schema Registry поддерживает несколько форматов одновременно. Например, в одном субъекте могут быть схемы Avro, а в другом — JSON Schema. Более того, JSON Schema имеет собственные правила совместимости, поэтому схемы Avro могут обновляться как в обратном, так и в прямом направлении.
Серверная часть Kafka#
Для хранения Schema Registry используется серверная часть Kafka. Специальный топик (по умолчанию _crx_schemas) с одной партицией, используется в качестве
высокодоступного лога предзаписи. Все схемы, субъекты, версии и идентификаторы метаданных, а также настройки совместимости
добавляются в этот лог-файл в виде сообщений. Таким образом, экземпляр Schema Registry как производит, так и потребляет
сообщения в топике _crx_schemas. Он производит сообщения в лог-файл, когда, например, регистрируются новые схемы для
субъекта или когда регистрируются обновления настроек совместимости.
Schema Registry потребляет данные из логов _crx_schemas в фоновом потоке и обновляет свои локальные кеши
при получении каждого нового сообщения _crx_schemas, чтобы отразить новую добавленную схему или настройку совместимости.
Обновление локального состояния из логов Kafka таким образом обеспечивает согласованность узлов, упорядоченность и простоту восстановления.
Реестр хранит данные вечно и возлагает на пользователя выбор какие данные очищать. На топике выставлена настройка cleanup.policy=compact,
что позволяет топику очищаться от неактуальных данных при помощи механизма Kafka compaction. Для всех записей в топике, подлежащих
удалению, создается специальная tombstone-запись, которую затем распознает Kafka.
Для очистки топика используйте одну из команд:
/cleanup/by/retention— удалить схемы, чей возраст превышает переданный retention;/cleanup/by/timestamp— удалить схемы, имеющие timestamp старше переданного;/cleanup/by/latest/version— удалить все схемы из субъектов, кроме последней.
Пример процесса очистки:
Пользователь отправил запрос
curl -X DELETE -H -d "{\"retention.ms\":84600000}" <url>/cleanup/by/retention. (удалить все схемы старше суток)Реестр получил запрос на очистку и:
Удаляет соответствующие схемы из памяти, в служебный топик записывает записи об удалении.
При завершении сегмента лога и его сжатии, сообщения о создании схем удаляются из топика
_crx_schemas.