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). Оба заменяют полную схему компактным идентификатором.
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) схемы.
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.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 — централизованного сервиса хранения схем.
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.Overhead: 5 байт на сообщение (1 magic + 4 schema ID). Для 100-байтового payload — 5% overhead вместо 2000%. Schema Registry хранит схемы и выдаёт их по ID через HTTP API.
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 взаимодействуют с ним при сериализации и десериализации:
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.
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 scope | Когда использовать |
|---|---|---|
| TopicName (default) | Per topic | Один тип данных на topic. Простейший случай. |
| RecordName | Per type | Несколько типов в одном topic (event sourcing). Compatibility per type. |
| TopicRecordName | Per type per topic | Гибрид: один тип может эволюционировать по-разному в разных topics. |
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
| Аспект | GenericRecord | SpecificRecord |
|---|---|---|
| Type safety | Runtime (cast exceptions) | Compile-time |
| Schema binding | Runtime — любая схема | Compile-time — сгенерированный класс |
| Code generation | Не нужна | Обязательна (avro-tools, avro-maven-plugin) |
| Гибкость | Высокая — подходит для generic tools | Низкая — per-schema класс |
| Performance | Чуть медленнее (boxing, map lookup) | Чуть быстрее (direct field access) |
| Kafka use case | Consumer читает разные schema versions | Producer/consumer знают exact schema |
В 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 | Оба направления, все версии | Самый строгий режим |
BACKWARD vs BACKWARD_TRANSITIVE: BACKWARD проверяет совместимость только с последней зарегистрированной версией. BACKWARD_TRANSITIVE — со всеми предыдущими версиями. Transitive защищает от ситуации, когда v3 совместима с v2, но не с v1. Детальный разбор compatibility strategies — в модуле 08.
Ключевые выводы
- Single Object Encoding — spec Avro: C3 01 + 8-byte Rabin fingerprint + payload. Не требует внешнего сервиса, но нуждается в локальном реестре схем
- Confluent Wire Format — de facto стандарт для Kafka: 0x00 + 4-byte schema ID + payload. Overhead — 5 байт на сообщение
- Schema Registry кэшируется на стороне клиента — в steady state HTTP-вызовов нет, всё из кэша
- Subject naming определяет scope для compatibility checks: TopicName (default, per-topic), RecordName (per-type), TopicRecordName (per-type-per-topic)
- GenericRecord — динамический доступ без code generation. SpecificRecord — compile-time type safety с code generation
- 7 compatibility modes от NONE (без проверок) до FULL_TRANSITIVE (оба направления, все версии)