Концепции#

Corax Schema Registry — это управляемый репозиторий схем, который поддерживает хранение и обмен данными как для потоковой обработки, так и для данных в состоянии покоя, таких как базы данных, файлы и другие статические хранилища данных. В этом разделе объясняются некоторые фундаментальные понятия, связанные со схемами, реестром и утилитами, а также рассматриваются детали того, как все это работает с точки зрения внутреннего устройства.

Как это работает#

Schema Registry:

  • предоставляет сервисный слой для метаданных;

  • предоставляет REST-интерфейс для хранения и извлечения схем Avro и JSON;

  • хранит версионную историю всех схем на основе заданной стратегии имен субъектов;

  • предоставляет несколько настроек совместимости и позволяет обновлять схемы в соответствии с настроенными параметрами совместимости и расширенной поддержкой этих типов схем;

  • предоставляет сериализаторы, подключаемые к клиентам Kafka, которые обрабатывают хранение и извлечение схем для сообщений Kafka, отправляемых в любом из поддерживаемых форматов.

Schema Registry расположен отдельно от брокеров Kafka и может быть развернут на сервере без Corax. Производители и потребители по-прежнему обращаются к брокерам Kafka для публикации и чтения данных (сообщений) в топиках. Одновременно с этим они могут обращаться к Schema Registry для отправки и получения схем, описывающих модели данных для сообщений.

SR-diagram

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, но без строго определенной схемы формата.

Однако отсутствие строго определенного формата имеет два существенных недостатка:

  1. Потребители данных могут не понимать производителей данных: отсутствие структуры делает потребление данных в таких форматах более сложным, поскольку поля могут быть произвольно добавлены или удалены, а данные могут быть повреждены. Этот недостаток становится тем серьезнее, чем больше приложений или команд в организации начинают потреблять поток данных: если команда, находящаяся выше по течению, может вносить произвольные изменения в формат данных по своему усмотрению, то становится очень сложно гарантировать, что все потребители, находящиеся ниже по течению, смогут и в дальнейшем интерпретировать данные. Не хватает контракта для данных между производителями и потребителями, аналогичного контракту API.

  2. Накладные расходы и увеличение размера: универсальный формат обмена данными имеет больший размер, поскольку имена полей и информация о типе должны быть явно представлены в сериализованном формате, несмотря на то, что они идентичны во всех сообщениях.

Дальнейшее развитие форматов сериализации привело к появлению межъязыковых библиотек сериализации которые требуют формального определения структуры данных с помощью схем. К таким библиотекам относятся 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 — удалить все схемы из субъектов, кроме последней.

Пример процесса очистки:

  1. Пользователь отправил запрос curl -X DELETE -H -d "{\"retention.ms\":84600000}" <url>/cleanup/by/retention. (удалить все схемы старше суток)

  2. Реестр получил запрос на очистку и:

    1. Удаляет соответствующие схемы из памяти, в служебный топик записывает записи об удалении.

    2. При завершении сегмента лога и его сжатии, сообщения о создании схем удаляются из топика _crx_schemas.