Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 05.05 · 30 мин
Продвинутый
AvroKafkaWire FormatSchema RegistrySingle Object EncodingConfluentSubject NamingGeneric Record

Kafka-интеграция и wire formats

Проблема: схема на каждое сообщение

В OCF (Object Container File) writer schema хранится один раз в header файла. Все записи в файле разделяют одну схему — overhead минимален. Но в Kafka каждое сообщение независимо: consumer может прочитать любое сообщение без чтения предыдущих.

Вложить полную JSON-схему (сотни байт — килобайты для сложных типов) в каждое сообщение (десятки — сотни байт payload) — неприемлемый overhead. Kafka-сообщение с 100 байтами данных и 2 КБ схемы — это 95% overhead.

Два решения: Single Object Encoding (спецификация Avro) и Confluent Wire Format (de facto стандарт для Kafka). Оба заменяют полную схему компактным идентификатором.

NOTE

Single Object Encoding — часть Avro spec 1.8.2+. Confluent Wire Format — проприетарное расширение Confluent Platform. Оба решают одну проблему, но несовместимы между собой: разные magic bytes, разные идентификаторы схемы.

Single Object Encoding

Avro spec определяет формат для хранения одиночных объектов вне OCF контейнера. Вместо полной схемы используется 8-байтовый fingerprint — Rabin hash от Parsing Canonical Form (PCF) схемы.

Single Object Encoding Format

C3 01 (2 байта)

Два фиксированных байта: 0xC3 0x01 (195, 1 в decimal). Маркер формата Single Object Encoding. Позволяет отличить от OCF (magic: Obj\x01) и Confluent (magic: 0x00).

Rabin Fingerprint (8 байт)

8-байтовый CRC-64-AVRO (Rabin fingerprint) от Parsing Canonical Form схемы. PCF — каноническое представление: без пробелов, алиасов, doc, defaults. Одинаковая схема → одинаковый fingerprint.

Avro Binary Payload

Avro binary payload — объект, сериализованный по обычным правилам бинарного кодирования. Без schema, без разделителей. Ридер декодирует по схеме, найденной через fingerprint.
MagicОтличает Single Object Encoding от других форматов. 0xC3 выбран так, чтобы не совпадать с первым байтом ASCII (Obj), JSON ({, [), или Confluent (0x00).
FingerprintCRC-64-AVRO polynomial: 0xC96C5795D7870F42. Вход: Parsing Canonical Form (PCF) — нормализованный JSON без whitespace, doc, aliases. Коллизии возможны, но крайне маловероятны (2^64 пространство).
PayloadСтандартное Avro binary encoding — то же, что внутри data block OCF, но ровно один объект. Без count, без size prefix.

Parsing Canonical Form (PCF) — нормализованная JSON-форма схемы для вычисления fingerprint. Убирает whitespace, doc, aliases, default — оставляет только структуру. Две эквивалентные схемы с разным форматированием дают одинаковый PCF и одинаковый fingerprint.

Ограничения:

  • Fingerprint — 8 байт (64 бита). Коллизии теоретически возможны, хотя для практических объёмов пренебрежимо маловероятны
  • Ридер должен иметь локальный реестр схем (fingerprint → schema) — spec не определяет протокол получения схемы по fingerprint
  • Не получил широкого adoption в Kafka — Confluent Wire Format доминирует

Confluent Wire Format

Confluent Platform определяет собственный wire format, ставший de facto стандартом для Avro в Kafka. Вместо fingerprint используется 4-байтовый schema ID из Schema Registry — централизованного сервиса хранения схем.

Confluent Wire Format

0x00 (1 байт)

Один байт: 0x00. Magic byte Confluent wire format. Отличает от Avro OCF (0x4F = 'O'), Single Object Encoding (0xC3), и JSON (0x7B = '{').

Schema ID (4 байта BE)

4 байта, big-endian unsigned integer. ID схемы в Confluent Schema Registry. Присваивается при первой регистрации схемы. Глобально уникален в пределах Registry.

Avro Binary Payload

Avro binary payload — стандартное бинарное кодирование без header. Ридер получает схему из Schema Registry по ID и декодирует payload.
Пример: сообщение 105 байт1 byte magic + 4 bytes schema ID + 100 bytes payload = 105 bytes. Overhead: 5 bytes (4.8%). Сравните с 2 КБ встроенной JSON-схемы.

Overhead: 5 байт на сообщение (1 magic + 4 schema ID). Для 100-байтового payload — 5% overhead вместо 2000%. Schema Registry хранит схемы и выдаёт их по ID через HTTP API.

TIP

Schema ID — глобально инкрементальный в пределах Schema Registry. Первая зарегистрированная схема получает ID = 1, вторая — ID = 2 и т.д. Это не hash, а автоинкрементный ключ в хранилище (Kafka topic _schemas). 4 байта дают до ~4.3 миллиарда схем.

Schema Registry: Produce / Consume Flow

Schema Registry — центральный компонент Kafka-экосистемы для управления схемами. Producer и consumer взаимодействуют с ним при сериализации и десериализации:

Schema Registry: Produce и Consume

Produce Flow

Producer: serialize record

Producer знает схему записи (compile-time или runtime). При первой отправке — регистрирует схему в Schema Registry для соответствующего subject.

Registry: register/lookup schema → ID

POST /subjects/{subject}/versions — регистрация или lookup существующей схемы. Registry проверяет compatibility, присваивает ID если новая. Возвращает schema ID. Producer кэширует ID локально.

Serialize: 0x00 + ID + binary

Serializer: 0x00 + 4-byte schema ID + Avro binary payload. KafkaAvroSerializer делает это автоматически.

Kafka Broker: store bytes

Kafka broker получает байтовое сообщение. Broker не знает про Avro и Schema Registry — для него это просто bytes. Схема не передаётся через Kafka.

Consume Flow

Consumer: read bytes from Kafka

Consumer читает raw bytes из Kafka. Первый байт 0x00 → Confluent wire format. Следующие 4 байта → schema ID.

Registry: fetch schema by ID

GET /schemas/ids/{id} — получить writer schema по ID. Consumer кэширует ответ (ID → schema) для последующих сообщений. Один HTTP-запрос на каждый уникальный schema ID.

Resolution: writer + reader schemas

Schema resolution: writer schema (из Registry) + reader schema (из consumer config). ResolvingDecoder строит план чтения. Кэшируется для пары (writer, reader).

Deserialized record

Десериализованная запись: GenericRecord (динамический) или конкретный POJO (если использовался SpecificRecord с code generation). Доступна для бизнес-логики consumer.

Кэширование: Producer кэширует (subject, schema) → ID после первой регистрации. Consumer кэширует ID → schema после первого fetch. В steady state Schema Registry не участвует в пути данных — все вызовы идут из кэша. HTTP overhead только при cold start или появлении нового schema ID.

WARNING

Kafka broker не знает про Schema Registry, Avro или schema validation. Для брокера сообщение — просто byte[]. Если producer отправит невалидные байты (без magic byte, с несуществующим schema ID) — broker примет их. Ошибка обнаружится только при десериализации на consumer. Schema Registry enforcement — на стороне клиента (serializer/deserializer).

Subject Naming Strategies

Schema Registry организует схемы по subjects — логическим именам, привязанным к потокам данных. Subject определяет scope для compatibility checks: новая версия схемы проверяется на совместимость с предыдущими версиями того же subject.

TopicNameStrategy (default):

Subject для key: {topic-name}-key
Subject для value: {topic-name}-value

Пример: topic "orders"
 → orders-key (схема ключа)
 → orders-value (схема значения)

RecordNameStrategy:

Subject = полное имя типа (namespace + name)

Пример: com.example.Order → "com.example.Order"
 — один subject для всех topics, где используется этот тип

TopicRecordNameStrategy:

Subject = {topic-name}-{record-name}

Пример: orders + com.example.Order → "orders-com.example.Order"
 — один subject per type per topic
Subject Naming Strategies
TopicNameStrategy (default)Простейший вариант: один subject на topic + key/value. Подходит когда в topic один тип записей. Compatibility проверяется per-topic.
RecordNameStrategySubject = полное имя record (namespace.name). Один subject для типа across all topics. Подходит для event sourcing: разные типы событий в одном topic.
TopicRecordNameStrategyГибрид: subject = topic + record name. Один тип может эволюционировать по-разному в разных topics. Максимальная гранулярность compatibility checks.
СтратегияSubject scopeКогда использовать
TopicName (default)Per topicОдин тип данных на topic. Простейший случай.
RecordNamePer typeНесколько типов в одном topic (event sourcing). Compatibility per type.
TopicRecordNamePer type per topicГибрид: один тип может эволюционировать по-разному в разных topics.
TIP

TopicNameStrategy запрещает писать разные типы записей в один topic — вторая регистрация другой схемы под тот же subject вызовет compatibility check, который скорее всего провалится. Для event sourcing (разные типы событий в одном topic) используйте RecordNameStrategy.

Generic vs Specific Records

Avro предлагает два стиля работы с записями на стороне клиента:

GenericRecord — динамический доступ по имени поля:

GenericRecord user = new GenericData.Record(schema);
user.put("name", "Alice");
String name = (String) user.get("name"); // без type safety

SpecificRecord — code generation из schema, строго типизированные классы:

// Сгенерированный класс из avro-maven-plugin
User user = User.newBuilder()
 .setName("Alice")
 .setId(42L)
 .build();
String name = user.getName(); // compile-time type check
АспектGenericRecordSpecificRecord
Type safetyRuntime (cast exceptions)Compile-time
Schema bindingRuntime — любая схемаCompile-time — сгенерированный класс
Code generationНе нужнаОбязательна (avro-tools, avro-maven-plugin)
ГибкостьВысокая — подходит для generic toolsНизкая — per-schema класс
PerformanceЧуть медленнее (boxing, map lookup)Чуть быстрее (direct field access)
Kafka use caseConsumer читает разные schema versionsProducer/consumer знают exact schema
NOTE

В Kafka-экосистеме GenericRecord используется чаще: consumer не знает заранее все версии схем, которые встретит в topic. SpecificRecord удобнее для producer (одна известная схема) и для stream processing (Kafka Streams, ksqlDB), где тип фиксирован.

Compatibility Modes — обзор

Schema Registry поддерживает семь режимов совместимости. Каждый определяет, какие изменения допустимы при регистрации новой версии схемы:

РежимПроверкаОписание
BACKWARD (default)Новая читает старыеДобавление поля с default, удаление поля
FORWARDСтарая читает новыеУдаление поля с default, добавление поля
FULLОба направленияТолько с default: добавление и удаление
NONE проверкиЛюбые изменения — на ответственности разработчика
BACKWARD_TRANSITIVEНовая читает все старыеКак BACKWARD, но проверка со всеми версиями, не только последней
FORWARD_TRANSITIVEВсе старые читают новуюКак FORWARD, но transitive
FULL_TRANSITIVEОба направления, все версииСамый строгий режим
NOTE

BACKWARD vs BACKWARD_TRANSITIVE: BACKWARD проверяет совместимость только с последней зарегистрированной версией. BACKWARD_TRANSITIVE — со всеми предыдущими версиями. Transitive защищает от ситуации, когда v3 совместима с v2, но не с v1. Детальный разбор compatibility strategies — в модуле 08.

Ключевые выводы

  1. Single Object Encoding — spec Avro: C3 01 + 8-byte Rabin fingerprint + payload. Не требует внешнего сервиса, но нуждается в локальном реестре схем
  2. Confluent Wire Format — de facto стандарт для Kafka: 0x00 + 4-byte schema ID + payload. Overhead — 5 байт на сообщение
  3. Schema Registry кэшируется на стороне клиента — в steady state HTTP-вызовов нет, всё из кэша
  4. Subject naming определяет scope для compatibility checks: TopicName (default, per-topic), RecordName (per-type), TopicRecordName (per-type-per-topic)
  5. GenericRecord — динамический доступ без code generation. SpecificRecord — compile-time type safety с code generation
  6. 7 compatibility modes от NONE (без проверок) до FULL_TRANSITIVE (оба направления, все версии)
Kafka: Avro deep-dive и Schema Registry Kafka: compatibility modes (BACKWARD/FORWARD/FULL) Debezium: Avro в Schema Registry для CDC

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 6