Раздел 11. Работа с файлами, I/O и сетью (NIO.2)
Глава 3. Сериализация и форматы обмена
Apache Avro — бинарный формат сериализации с динамической схемой
Apache Avro — это формат сериализации данных и протокол удалённого вызова процедур (RPC), разработанный в рамках проекта Apache Hadoop. Avro был создан Дугом Каттингом в 2009 году как ответ на необходимость иметь компактный, быстрый и расширяемый формат для хранения и передачи больших данных в распределённых системах. На момент 2026 года актуальная версия спецификации — Avro 1.12.x, а проект остаётся одним из столпов экосистемы Apache Big Data.
В отличие от Protocol Buffers и Thrift, где схема является внешним контрактом, разделяемым между producer и consumer, Avro следует принципу self-describing (самоописывающихся) данных: каждый блок данных несёт в себе свою схему или ссылку на неё. Это фундаментальное архитектурное решение определяет все сильные и слабые стороны Avro.
Архитектура Avro: схема + данные
Avro разделяет понятия схемы (schema) и данных (datum). Схема описывает структуру данных на языке JSON. Данные сериализуются в бинарный формат, который интерпретируется исключительно через схему. Без схемы бинарный поток Avro — это бессмысленная последовательность байтов, так как в нём отсутствуют теги полей и информация о типах.
Это кардинально отличается от Protobuf, где wire format содержит теги (field numbers), и парсер может прочитать сообщение, зная только .proto файл. В Avro парсер должен иметь схему, чтобы понять, где заканчивается одно поле и начинается другое. Например, для строки Avro записывает длину в varint, затем байты. Для массива — блоки элементов с указанием размера. Но чтобы понять, что следующее значение — это строка, а не int, нужна схема.
Почему это важно
Такой подход даёт две ключевые возможности:
Компактность: в бинарном потоке нет никаких метаданных на уровне полей — ни имён, ни тегов, ни типов. Это делает Avro ещё компактнее Protobuf для многих сценариев.
Динамическая типизация: любой потребитель, получивший данные вместе со схемой, может их десериализовать без предварительно сгенерированных классов.
Но есть и цена: схема должна быть доступна при чтении. В долгосрочном хранении (HDFS, S3) схема обычно записывается в заголовок файла. В потоковой передаче (Kafka) схема регистрируется в Schema Registry и передаётся по ссылке (ID схемы), а не целиком.
Схема Avro
Схема Avro — это JSON-документ, описывающий тип данных. Avro поддерживает примитивные типы, сложные типы и логические типы.
Примитивные типы
Сложные типы
Пример схемы
Ключевые элементы:
#Java #для_новичков #beginner #IO #NIO #Serialize #Avro
Глава 3. Сериализация и форматы обмена
Apache Avro — бинарный формат сериализации с динамической схемой
Apache Avro — это формат сериализации данных и протокол удалённого вызова процедур (RPC), разработанный в рамках проекта Apache Hadoop. Avro был создан Дугом Каттингом в 2009 году как ответ на необходимость иметь компактный, быстрый и расширяемый формат для хранения и передачи больших данных в распределённых системах. На момент 2026 года актуальная версия спецификации — Avro 1.12.x, а проект остаётся одним из столпов экосистемы Apache Big Data.
В отличие от Protocol Buffers и Thrift, где схема является внешним контрактом, разделяемым между producer и consumer, Avro следует принципу self-describing (самоописывающихся) данных: каждый блок данных несёт в себе свою схему или ссылку на неё. Это фундаментальное архитектурное решение определяет все сильные и слабые стороны Avro.
Self-describing data — это данные, которые содержат внутри себя или в непосредственной близости полное описание своей структуры (схему). Это позволяет любому потребителю, не имеющему предварительного знания о формате, корректно интерпретировать содержимое.
Big Data — это термин, обозначающий наборы данных, объём, скорость поступления и разнообразие которых настолько велики, что традиционные инструменты обработки не справляются с ними. Hadoop — это фреймворк для распределённой обработки Big Data на кластерах commodity-оборудования.
Архитектура Avro: схема + данные
Avro разделяет понятия схемы (schema) и данных (datum). Схема описывает структуру данных на языке JSON. Данные сериализуются в бинарный формат, который интерпретируется исключительно через схему. Без схемы бинарный поток Avro — это бессмысленная последовательность байтов, так как в нём отсутствуют теги полей и информация о типах.
Это кардинально отличается от Protobuf, где wire format содержит теги (field numbers), и парсер может прочитать сообщение, зная только .proto файл. В Avro парсер должен иметь схему, чтобы понять, где заканчивается одно поле и начинается другое. Например, для строки Avro записывает длину в varint, затем байты. Для массива — блоки элементов с указанием размера. Но чтобы понять, что следующее значение — это строка, а не int, нужна схема.
Почему это важно
Такой подход даёт две ключевые возможности:
Компактность: в бинарном потоке нет никаких метаданных на уровне полей — ни имён, ни тегов, ни типов. Это делает Avro ещё компактнее Protobuf для многих сценариев.
Динамическая типизация: любой потребитель, получивший данные вместе со схемой, может их десериализовать без предварительно сгенерированных классов.
Но есть и цена: схема должна быть доступна при чтении. В долгосрочном хранении (HDFS, S3) схема обычно записывается в заголовок файла. В потоковой передаче (Kafka) схема регистрируется в Schema Registry и передаётся по ссылке (ID схемы), а не целиком.
Schema Registry — это сервис (например, Confluent Schema Registry), который хранит версии схем Avro, Protobuf и JSON Schema. Producer регистрирует схему и получает уникальный ID. Consumer получает ID вместе с сообщением и запрашивает схему из реестра. Это избавляет от необходимости передавать полную схему в каждом сообщении.
Схема Avro
Схема Avro — это JSON-документ, описывающий тип данных. Avro поддерживает примитивные типы, сложные типы и логические типы.
Примитивные типы
null — отсутствие значения.boolean — true/false.int — 32-битное целое со знаком.long — 64-битное целое со знаком.float — 32-битное IEEE 754 float.double — 64-битное IEEE 754 double.bytes — последовательность байтов.string — строка в UTF-8.Сложные типы
record — именованная коллекция полей (аналог struct или class).enum — именованное множество значений.array — упорядоченная коллекция однотипных элементов.map — ассоциативный массив со строковыми ключами.union — объединение типов, записывается как JSON-массив. Например, ["null", "string"] означает nullable string.fixed — массив байтов фиксированного размера.Пример схемы
{
"type": "record",
"name": "Book",
"namespace": "com.example.library",
"doc": "Описание книги в библиотеке",
"fields": [
{
"name": "title",
"type": "string",
"doc": "Название книги"
},
{
"name": "author",
"type": "string"
},
{
"name": "year",
"type": ["null", "int"],
"default": null
},
{
"name": "price",
"type": "double",
"default": 0.0
},
{
"name": "tags",
"type": {
"type": "array",
"items": "string"
},
"default": []
},
{
"name": "metadata",
"type": {
"type": "map",
"values": "string"
},
"default": {}
}
]
}Ключевые элементы:
type: "record" — объявляет именованную запись.name и namespace — полное имя типа, аналог package + class в Java.doc — документация, которая может генерироваться в Javadoc при кодогенерации.fields — массив полей. Каждое поле имеет name, type и опционально default, doc, order (для сортировки).type: ["null", "int"] — union type, реализующий nullable int. Порядок в union важен: при сериализации Avro записывает индекс типа (varint), затем значение. null имеет индекс 0, int — индекс 1.Union type (объединение типов) — это тип данных, значение которого может принадлежать одному из нескольких указанных типов. В Avro union записывается как массив типов, и при сериализации перед значением записывается индекс выбранного типа.
#Java #для_новичков #beginner #IO #NIO #Serialize #Avro
👍4
Динамическая типизация: GenericRecord vs SpecificRecord
Avro предоставляет два основных API для работы с данными: Generic и Specific.
Generic API
Generic API не требует кодогенерации. Данные представляются как
В этом примере
Specific API
Specific API требует кодогенерации: из схемы Avro генерируются Java-классы через
Specific API быстрее Generic, так как доступ к полям идёт через
Версионирование: Reader Schema и Writer Schema
Avro имеет наиболее элегантную среди бинарных форматов систему schema evolution (эволюции схемы). Ключевая концепция — разделение схемы записи (writer schema) и схемы чтения (reader schema).
#Java #для_новичков #beginner #IO #NIO #Serialize #Avro
Avro предоставляет два основных API для работы с данными: Generic и Specific.
Generic API
Generic API не требует кодогенерации. Данные представляются как
GenericRecord — динамическая структура, аналогичная Map<String, Object>, но типобезопасная на уровне схемы.import org.apache.avro.Schema;
import org.apache.avro.generic.GenericData;
import org.apache.avro.generic.GenericRecord;
import org.apache.avro.generic.GenericDatumWriter;
import org.apache.avro.generic.GenericDatumReader;
import org.apache.avro.io.DatumWriter;
import org.apache.avro.io.DatumReader;
import org.apache.avro.io.Encoder;
import org.apache.avro.io.Decoder;
import org.apache.avro.io.EncoderFactory;
import org.apache.avro.io.DecoderFactory;
import java.io.ByteArrayOutputStream;
import java.io.ByteArrayInputStream;
public class AvroGenericExample {
// Схема как JSON-строка
private static final String BOOK_SCHEMA = "{"
+ "\"type\":\"record\","
+ "\"name\":\"Book\","
+ "\"namespace\":\"com.example\","
+ "\"fields\":["
+ " {\"name\":\"title\", \"type\":\"string\"},"
+ " {\"name\":\"author\", \"type\":\"string\"},"
+ " {\"name\":\"year\", \"type\":[\"null\",\"int\"], \"default\":null},"
+ " {\"name\":\"price\", \"type\":\"double\", \"default\":0.0}"
+ "]}";
public byte[] serializeBook() throws Exception {
// Парсинг схемы из JSON-строки
Schema schema = new Schema.Parser().parse(BOOK_SCHEMA);
// Создание GenericRecord — динамического контейнера данных
// GenericRecord хранит значения полей в массиве Object[]
GenericRecord book = new GenericData.Record(schema);
book.put("title", "Clean Code");
book.put("author", "Robert C. Martin");
book.put("year", 2008); // Avro автоматически оборачивает в union
book.put("price", 42.50);
// GenericDatumWriter — сериализатор для GenericRecord
// DatumWriter — это интерфейс, абстрагирующий запись Avro-данных
DatumWriter<GenericRecord> writer = new GenericDatumWriter<>(schema);
ByteArrayOutputStream out = new ByteArrayOutputStream();
// BinaryEncoder — кодировщик в бинарный формат Avro
// EncoderFactory — фабрика для создания encoder-ов (бинарных, JSON и др.)
Encoder encoder = EncoderFactory.get().binaryEncoder(out, null);
// Сериализация: writer читает поля из GenericRecord по схеме
// и записывает их в encoder в бинарном формате
writer.write(book, encoder);
encoder.flush(); // сброс буфера encoder в поток
return out.toByteArray();
}
public GenericRecord deserializeBook(byte[] data) throws Exception {
Schema schema = new Schema.Parser().parse(BOOK_SCHEMA);
// GenericDatumReader — десериализатор
DatumReader<GenericRecord> reader = new GenericDatumReader<>(schema);
ByteArrayInputStream in = new ByteArrayInputStream(data);
// BinaryDecoder — декодировщик из бинарного формата
Decoder decoder = DecoderFactory.get().binaryDecoder(in, null);
// Десериализация: reader читает байты через decoder
// и строит GenericRecord по схеме
return reader.read(null, decoder);
}
}
В этом примере
GenericRecord — это реализация интерфейса IndexedRecord, которая хранит значения полей в массиве Object[]. Поле title имеет позицию 0, author — 1, year — 2 и т.д. Метод put(String name, Object value) ищет позицию поля по имени через схему и записывает значение в массив. Это медленнее прямого доступа к полю, но даёт полную динамичность.Specific API
Specific API требует кодогенерации: из схемы Avro генерируются Java-классы через
avro-tools или Maven-плагин. Сгенерированные классы наследуют org.apache.avro.specific.SpecificRecordBase и предоставляют типизированные геттеры и сеттеры.// Сгенерированный класс (упрощённо)
public class Book extends org.apache.avro.specific.SpecificRecordBase {
private CharSequence title;
private CharSequence author;
private Integer year;
private double price;
// Avro требует конструктор по умолчанию
public Book() {}
// put(int field, Object value) — индексированный доступ
public void put(int field, Object value) {
switch (field) {
case 0: title = (CharSequence) value; break;
case 1: author = (CharSequence) value; break;
case 2: year = (Integer) value; break;
case 3: price = (Double) value; break;
default: throw new AvroRuntimeException("Bad index");
}
}
public Object get(int field) {
switch (field) {
case 0: return title;
case 1: return author;
case 2: return year;
case 3: return price;
default: throw new AvroRuntimeException("Bad index");
}
}
// Типизированные геттеры/сеттеры
public CharSequence getTitle() { return title; }
public void setTitle(CharSequence value) { this.title = value; }
}
Specific API быстрее Generic, так как доступ к полям идёт через
switch по индексу, а не через поиск по имени в HashMap. Однако он требует перекомпиляции при изменении схемы.Версионирование: Reader Schema и Writer Schema
Avro имеет наиболее элегантную среди бинарных форматов систему schema evolution (эволюции схемы). Ключевая концепция — разделение схемы записи (writer schema) и схемы чтения (reader schema).
Writer schema — это схема, которая использовалась при сериализации данных. Reader schema — это схема, которую ожидает потребитель при десериализации. Avro гарантирует совместимость, если reader schema может быть получена из writer schema через набор правил преобразования.
#Java #для_новичков #beginner #IO #NIO #Serialize #Avro
👍4
Правила резолюции схем
Avro определяет строгие правила, как reader schema может отличаться от writer schema:
Добавление поля с default: writer записывает без нового поля, reader читает и использует default. Это backward compatible.
Удаление поля с default: writer записывает с полем, reader игнорирует его. Это forward compatible.
Изменение типа: разрешено только между совместимыми типами (например, int -> long). Avro выполняет implicit promotion.
Изменение порядка полей: разрешено — поля идентифицируются по имени, а не по позиции.
Изменение имени поля: эквивалентно удалению старого и добавлению нового. Для сохранения совместимости используются алиасы (
При десериализации Avro сопоставляет поля writer и reader по имени. Поле
Implicit promotion (неявное повышение типа) — это правило Avro, позволяющее читать значение одного типа как другой без потери данных. Например,
Сжатие и блоки
Avro поддерживает сжатие на уровне блоков данных. При записи в файл данные группируются в блоки фиксированного размера (по умолчанию), и каждый блок сжимается независимо.
Codec (кодек) — это алгоритм сжатия/распаковки данных. В Avro кодек применяется к блокам записей, а не к каждой записи отдельно. Это повышает эффективность сжатия за счёт большего объёма данных на входе кодека.
Структура Avro-файла
Avro-файл (Object Container File) имеет следующую структуру:
Magic (4 байта):
Metadata (map): содержит
Sync marker (16 байт): случайная последовательность, разделяющая блоки.
Blocks: последовательность блоков, каждый из которых содержит:
Количество записей в блоке (long).
Размер сжатых данных (long).
Сжатые данные (сериализованные записи).
Sync marker.
Такая структура позволяет эффективно разбивать файлы на части (split) в Hadoop MapReduce: каждый mapper может начать чтение с ближайшего sync marker.
Использование в Kafka
Avro широко используется как формат сообщений в Apache Kafka, особенно в экосистеме Confluent. Интеграция работает через Schema Registry:
В сообщении Kafka с Avro первые 5 байт — служебные: 1 байт magic (
Это означает, что схема не передаётся в каждом сообщении — только 4-байтовый ID. Это критично для производительности: передача полной JSON-схемы (1–5 КБ) в каждом сообщении сделала бы Kafka непригодной для high-throughput сценариев.
Путь байтов в памяти JVM при работе с Avro
1. Парсинг схемы: от JSON к объекту Schema
Когда вы вызываете
JSON-строка схемы — это объект
Результат парсинга — объект
Полное имя типа (
Список полей (
Для каждого поля: имя, позиция, тип, default value, aliases.
Схема хранится в виде дерева объектов в куче.
Объект
2. Создание GenericRecord
При создании
Аллоцируется объект
Ссылку на
Массив
Массив
При вызове
Имя поля
Значение
Utf8 — это класс Avro, реализующий
Если используется
#Java #для_новичков #beginner #IO #NIO #Serialize #Avro
Avro определяет строгие правила, как reader schema может отличаться от writer schema:
Добавление поля с default: writer записывает без нового поля, reader читает и использует default. Это backward compatible.
Удаление поля с default: writer записывает с полем, reader игнорирует его. Это forward compatible.
Изменение типа: разрешено только между совместимыми типами (например, int -> long). Avro выполняет implicit promotion.
Изменение порядка полей: разрешено — поля идентифицируются по имени, а не по позиции.
Изменение имени поля: эквивалентно удалению старого и добавлению нового. Для сохранения совместимости используются алиасы (
aliases).// Writer schema (старая версия)
{
"type": "record",
"name": "Book",
"fields": [
{"name": "title", "type": "string"},
{"name": "author", "type": "string"}
]
}
// Reader schema (новая версия)
{
"type": "record",
"name": "Book",
"fields": [
{"name": "title", "type": "string"},
{"name": "author", "type": "string"},
{"name": "year", "type": ["null", "int"], "default": null}
]
}
При десериализации Avro сопоставляет поля writer и reader по имени. Поле
year отсутствует в writer schema, поэтому reader использует default: null. Это работает автоматически на уровне DatumReader, без участия прикладного кода.Implicit promotion (неявное повышение типа) — это правило Avro, позволяющее читать значение одного типа как другой без потери данных. Например,
int может быть прочитан как long, float как double. Обратное направление запрещено, так как может привести к потере точности.Сжатие и блоки
Avro поддерживает сжатие на уровне блоков данных. При записи в файл данные группируются в блоки фиксированного размера (по умолчанию), и каждый блок сжимается независимо.
import org.apache.avro.file.DataFileWriter;
import org.apache.avro.file.CodecFactory;
public void writeCompressedFile(Schema schema, List<GenericRecord> records) throws Exception {
// DataFileWriter — высокоуровневый writer для Avro-файлов с заголовком
DataFileWriter<GenericRecord> writer = new DataFileWriter<>(new GenericDatumWriter<>(schema));
// Установка кодека сжатия
// Snappy — быстрый кодек от Google, оптимизированный на скорость
// Deflate — стандартный zlib-сжатие
// Zstandard (zstd) — современный кодек с лучшим соотношением скорость/сжатие
writer.setCodec(CodecFactory.snappyCodec());
// Схема записывается в заголовок файла
writer.create(schema, new File("books.avro"));
for (GenericRecord record : records) {
writer.append(record);
}
writer.close();
}
Codec (кодек) — это алгоритм сжатия/распаковки данных. В Avro кодек применяется к блокам записей, а не к каждой записи отдельно. Это повышает эффективность сжатия за счёт большего объёма данных на входе кодека.
Структура Avro-файла
Avro-файл (Object Container File) имеет следующую структуру:
Magic (4 байта):
Obj\x01 — идентификатор формата.Metadata (map): содержит
avro.schema (JSON-схема в виде строки) и avro.codec (имя кодека).Sync marker (16 байт): случайная последовательность, разделяющая блоки.
Blocks: последовательность блоков, каждый из которых содержит:
Количество записей в блоке (long).
Размер сжатых данных (long).
Сжатые данные (сериализованные записи).
Sync marker.
Такая структура позволяет эффективно разбивать файлы на части (split) в Hadoop MapReduce: каждый mapper может начать чтение с ближайшего sync marker.
Split — это разбиение большого файла на части для параллельной обработки. В Hadoop InputFormat определяет, как файл разбивается на splits, которые затем обрабатываются отдельными mapper-ами.
Использование в Kafka
Avro широко используется как формат сообщений в Apache Kafka, особенно в экосистеме Confluent. Интеграция работает через Schema Registry:
import io.confluent.kafka.serializers.KafkaAvroSerializer;
import io.confluent.kafka.serializers.KafkaAvroDeserializer;
import org.apache.kafka.clients.producer.ProducerRecord;
public class KafkaAvroProducer {
public void sendBook(KafkaProducer<String, GenericRecord> producer, GenericRecord book) {
// KafkaAvroSerializer автоматически:
// 1. Регистрирует схему в Schema Registry (если ещё не зарегистрирована)
// 2. Получает schema ID
// 3. Записывает в сообщение: [magic byte (1)][schema ID (4)][payload]
ProducerRecord<String, GenericRecord> record =
new ProducerRecord<>("books-topic", book.get("title").toString(), book);
producer.send(record);
}
}
В сообщении Kafka с Avro первые 5 байт — служебные: 1 байт magic (
0x00), 4 байта ID схемы в Schema Registry (big-endian int). Остальное — бинарный payload Avro. Consumer получает ID, запрашивает схему из реестра и десериализует payload.Это означает, что схема не передаётся в каждом сообщении — только 4-байтовый ID. Это критично для производительности: передача полной JSON-схемы (1–5 КБ) в каждом сообщении сделала бы Kafka непригодной для high-throughput сценариев.
Throughput — это пропускная способность системы, измеряемая количеством операций или объёмом данных, обрабатываемых за единицу времени. В Kafka throughput измеряется в сообщениях в секунду или мегабайтах в секунду.
Путь байтов в памяти JVM при работе с Avro
1. Парсинг схемы: от JSON к объекту Schema
Когда вы вызываете
new Schema.Parser().parse(jsonString), происходит следующее:JSON-строка схемы — это объект
String в куче (Heap), обычно в Young Generation (Eden). Если схема загружается из файла, предварительно создаётся byte[] с содержимым файла, который декодируется в String через InputStreamReader.Schema.Parser использует Jackson (да, внутри Avro используется Jackson для парсинга JSON-схемы) для разбора JSON. Это создаёт временные объекты: JsonNode, ArrayNode, ObjectNode, Iterator. Все они аллоцируются в Eden.Результат парсинга — объект
Schema (или Schema.RecordSchema, Schema.ArraySchema и т.д.). Этот объект содержит:Полное имя типа (
name + namespace).Список полей (
List<Schema.Field>).Для каждого поля: имя, позиция, тип, default value, aliases.
Схема хранится в виде дерева объектов в куче.
Объект
Schema обычно кэшируется приложением (singleton) и живёт в Old Generation (Tenured), так как используется многократно. Промежуточные объекты Jackson (дерево JSON) становятся мусором и собираются при следующей Minor GC.2. Создание GenericRecord
GenericRecord record = new GenericData.Record(schema);
record.put("title", "Clean Code");
При создании
GenericData.Record:Аллоцируется объект
Record в Eden. Record содержит:Ссылку на
Schema (уже существует в Old Gen).Массив
Object[] values размером, равным количеству полей. Этот массив создаётся в Eden.Массив
int[] fieldFlags для отслеживания установленных полей.При вызове
put("title", value):Имя поля
"title" — строка в куче. Avro ищет позицию поля через Schema.getField(String). Schema кэширует Map<String, Field> (обычно HashMap), поэтому поиск выполняется быстро, но создаёт временный объект Integer для хэш-кода (автобоксинг).Значение
"Clean Code" — строка Java. Для поля типа string Avro ожидает CharSequence, а не обязательно java.lang.String. Это позволяет использовать Utf8 — оптимизированную реализацию CharSequence от Avro, которая хранит строку как byte[] в UTF-8, избегая двойного хранения (UTF-16 в Java String + UTF-8 в Avro).Utf8 — это класс Avro, реализующий
CharSequence и хранящий строку как изменяемый массив байтов в UTF-8. Это экономит память при сериализации, так как не требуется перекодирование из UTF-16 Java-строки в UTF-8 байты.Если используется
Utf8 вместо String, аллокация String избегается. Однако если прикладной код передаёт String, Avro может либо сконвертировать её в Utf8 (создав новый объект), либо оставить как String (в зависимости от настроек GenericData).#Java #для_новичков #beginner #IO #NIO #Serialize #Avro
👍3
3. Сериализация: от GenericRecord к byte[]
Вызов
Читает значение из
Определяет тип поля из схемы.
Вызывает соответствующий метод
Для строк (
Для union-типов (
После завершения
Временные объекты сериализации — внутренние структуры
4. Десериализация: от byte[] к GenericRecord
Вызов
BinaryDecoder читает байты из InputStream. Для каждого поля схемы:
Читает значение согласно типу поля.
Для строк: читает длину (varint), затем выделяет byte[] нужного размера, читает байты, оборачивает в Utf8.
Для union: читает индекс типа (varint), затем значение соответствующего типа.
Записывает значение в Record.values[pos].
Важная особенность: если reader schema отличается от writer schema, GenericDatumReader выполняет schema resolution (разрешение схемы) на лету. Это означает, что для каждого поля reader сопоставляет его с полем writer по имени, проверяет совместимость типов, применяет promotion (например, int -> long) и default values. Schema resolution требует дополнительных вычислений и создаёт временные объекты (например, ResolvingDecoder), но не требует аллокаций на каждое поле — логика реализована через state machine в decoder.
После десериализации BinaryDecoder становится мусором. Если decoder создан с DecoderFactory.get().binaryDecoder(in, reuse), он может быть переиспользован, избегая аллокации.
5. Работа GC с Avro-объектами
Avro создаёт значительно больше временных объектов, чем Protobuf, но меньше, чем Jackson JSON:
Schema: хранится в Old Generation как singleton. Размер объекта Schema зависит от сложности схемы: простая схема с 5 полями — несколько сотен байт, сложная вложенная схема — несколько килобайт.
GenericRecord: создаётся для каждой записи. Содержит
Utf8: Avro предпочитает
DataFileWriter / DataFileReader: при работе с файлами Avro использует буферизацию.
Object reuse: Avro предоставляет механизмы повторного использования объектов для снижения давления на GC. При чтении файла можно передать существующий
6. Avro в Kafka: путь байтов
При работе с Kafka и Confluent Serializer:
Producer вызывает
Сериализатор извлекает схему из
Проверяет кэш схем (локальный
Если нет — отправляет HTTP-запрос в Schema Registry для регистрации. Получает ID (int).
Сериализует record в
Формирует итоговое сообщение:
Этот
На стороне Consumer:
Читает magic byte и schema ID.
Проверяет локальный кэш схем по ID. Если нет — запрашивает схему из Schema Registry по HTTP.
Использует схему для создания
Десериализует payload в
Кэширование схем на клиенте критично: без него каждое сообщение порождало бы HTTP-запрос к реестру. Кэш — это
Сравнение Avro и Protobuf
Кодогенерация
Protobuf требует кодогенерации: .proto файл компилируется в Java-классы через protoc. Avro предоставляет выбор: Generic API (без кодогенерации) и Specific API (с кодогенерацией через
Оверхед схемы
В Protobuf схема существует только как внешний контракт. В бинарном сообщении нет схемы — только теги полей. В Avro бинарное сообщение без схемы нечитаемо. Поэтому Avro-файлы содержат схему в заголовке, а Kafka-сообщения содержат ID схемы из реестра.
При записи в файл: оверхед схемы амортизируется — одна схема на миллионы записей.
При передаче по сети через Schema Registry: оверхед — 5 байт на сообщение (magic + ID).
Без Schema Registry: пришлось бы передавать полную JSON-схему в каждом сообщении, что сделало бы Avro непригодным для messaging.
Типизация
Protobuf — статически типизирован через сгенерированные классы. Avro Generic — динамически типизирован через
Schema evolution
Avro имеет более мощную и гибкую систему schema evolution благодаря разделению reader/writer schema. Protobuf поддерживает forward/backward compatibility через правила добавления/удаления полей, но не имеет встроенного механизма разрешения различий между схемами на уровне runtime — это должно быть реализовано прикладным кодом.
Производительность
Protobuf быстрее Avro по нескольким причинам:
Сгенерированный код Protobuf содержит жёстко закодированные инструкции сериализации (switch по тегу). Avro Generic использует рефлексию-подобный обход схемы.
Protobuf не требует schema resolution при чтении (если reader/writer schema совпадают). Avro всегда выполняет resolution, даже когда схемы идентичны.
Protobuf использует immutable generated объекты с zero-copy строками. Avro Generic использует mutable
Однако Avro Specific API сопоставим по производительности с Protobuf, а в некоторых сценариях (большие файлы с блоковым сжатием) может быть быстрее за счёт эффективной интеграции с Hadoop/S3.
Экосистема
Protobuf: доминирует в gRPC, микросервисах, мобильной разработке. Поддержка в Kubernetes (etcd), Envoy, многих облачных провайдерах.
Avro: доминирует в Hadoop, Spark, Kafka, data lakes (Delta Lake, Iceberg). Интеграция с Hive, Presto/Trino, Flink.
#Java #для_новичков #beginner #IO #NIO #Serialize #Avro
Вызов
writer.write(record, encoder):GenericDatumWriter получает схему и обходит поля в порядке их объявления в схеме. Для каждого поля:Читает значение из
Record.values[pos] через record.get(pos).Определяет тип поля из схемы.
Вызывает соответствующий метод
Encoder: writeString(), writeInt(), writeDouble() и т.д.BinaryEncoder (по умолчанию) пишет данные в буфер. Внутренне Avro использует BufferedBinaryEncoder, который аккумулирует данные в byte[] буфере (обычно 2 КБ) перед записью в целевой OutputStream. Этот буфер создаётся один раз и переиспользуется.Для строк (
writeString): Avro записывает длину в varint, затем байты UTF-8. Если значение — Utf8, байты копируются напрямую. Если String, выполняется перекодирование UTF-16 -> UTF-8, что создаёт временный byte[].Для union-типов (
["null", "int"]): сначала записывается индекс выбранного типа (varint: 0 для null, 1 для int), затем само значение. Для null записывается только индекс 0.После завершения
encoder.flush() копирует оставшиеся данные из внутреннего буфера в OutputStream. Возвращается byte[] (если использовался ByteArrayOutputStream).Временные объекты сериализации — внутренние структуры
BinaryEncoder, возможно временные byte[] для строк — создаются в Eden. BinaryEncoder может быть переиспользован, если передать его вторым аргументом в EncoderFactory.get().binaryEncoder(out, reuse).4. Десериализация: от byte[] к GenericRecord
Вызов
reader.read(null, decoder):GenericDatumReader создаёт новый GenericData.Record (аллокация в Eden + Object[] для значений).BinaryDecoder читает байты из InputStream. Для каждого поля схемы:
Читает значение согласно типу поля.
Для строк: читает длину (varint), затем выделяет byte[] нужного размера, читает байты, оборачивает в Utf8.
Для union: читает индекс типа (varint), затем значение соответствующего типа.
Записывает значение в Record.values[pos].
Важная особенность: если reader schema отличается от writer schema, GenericDatumReader выполняет schema resolution (разрешение схемы) на лету. Это означает, что для каждого поля reader сопоставляет его с полем writer по имени, проверяет совместимость типов, применяет promotion (например, int -> long) и default values. Schema resolution требует дополнительных вычислений и создаёт временные объекты (например, ResolvingDecoder), но не требует аллокаций на каждое поле — логика реализована через state machine в decoder.
После десериализации BinaryDecoder становится мусором. Если decoder создан с DecoderFactory.get().binaryDecoder(in, reuse), он может быть переиспользован, избегая аллокации.
5. Работа GC с Avro-объектами
Avro создаёт значительно больше временных объектов, чем Protobuf, но меньше, чем Jackson JSON:
Schema: хранится в Old Generation как singleton. Размер объекта Schema зависит от сложности схемы: простая схема с 5 полями — несколько сотен байт, сложная вложенная схема — несколько килобайт.
GenericRecord: создаётся для каждой записи. Содержит
Object[] values. Для потоковой обработки (Kafka consumer) каждое сообщение порождает новый GenericRecord. При обработке 10 000 сообщений/сек это 10 000 объектов GenericRecord + 10 000 массивов Object[] в Eden в секунду. Minor GC запускается часто, но эти объекты короткоживущие и собираются быстро.Utf8: Avro предпочитает
Utf8 вместо String для строковых полей. Utf8 — mutable объект, содержащий byte[]. В потоковой обработке GenericDatumReader может переиспользовать один и тот же экземпляр Utf8, если decoder настроен на reuse. Это снижает давление на GC.DataFileWriter / DataFileReader: при работе с файлами Avro использует буферизацию.
DataFileWriter хранит блок записей в памяти до достижения порога синхронизации (sync interval, обычно несколько мегабайт). Этот буфер — byte[] в куче, который может попасть в Old Generation, если файл большой.Object reuse: Avro предоставляет механизмы повторного использования объектов для снижения давления на GC. При чтении файла можно передать существующий
GenericRecord в reader.read(existingRecord, decoder) — в этом случае decoder перезапишет значения в существующем объекте вместо создания нового. Это критично для high-throughput обработки.// Переиспользование GenericRecord для снижения аллокаций
GenericRecord reuse = new GenericData.Record(schema);
while (fileReader.hasNext()) {
// read перезаписывает reuse вместо создания нового объекта
GenericRecord record = fileReader.next(reuse);
process(record);
}
6. Avro в Kafka: путь байтов
При работе с Kafka и Confluent Serializer:
Producer вызывает
KafkaAvroSerializer.serialize(topic, record).Сериализатор извлекает схему из
GenericRecord.getSchema().Проверяет кэш схем (локальный
Map<Schema, Integer>) — если схема уже зарегистрирована, использует кэшированный ID.Если нет — отправляет HTTP-запрос в Schema Registry для регистрации. Получает ID (int).
Сериализует record в
byte[] через GenericDatumWriter + BinaryEncoder.Формирует итоговое сообщение:
[magic(1)][schemaId(4)][payload(N)].Этот
byte[] передаётся в Kafka Producer, который копирует его в ByteBuffer (send buffer) и отправляет брокеру.На стороне Consumer:
KafkaAvroDeserializer получает сообщение.Читает magic byte и schema ID.
Проверяет локальный кэш схем по ID. Если нет — запрашивает схему из Schema Registry по HTTP.
Использует схему для создания
GenericDatumReader.Десериализует payload в
GenericRecord.Кэширование схем на клиенте критично: без него каждое сообщение порождало бы HTTP-запрос к реестру. Кэш — это
ConcurrentHashMap в Old Generation, живущий всё время жизни consumer/producer.Сравнение Avro и Protobuf
Кодогенерация
Protobuf требует кодогенерации: .proto файл компилируется в Java-классы через protoc. Avro предоставляет выбор: Generic API (без кодогенерации) и Specific API (с кодогенерацией через
avro-tools или Maven-плагин). Это делает Avro более гибким для динамических сценариев, где схема известна только во время выполнения.Оверхед схемы
В Protobuf схема существует только как внешний контракт. В бинарном сообщении нет схемы — только теги полей. В Avro бинарное сообщение без схемы нечитаемо. Поэтому Avro-файлы содержат схему в заголовке, а Kafka-сообщения содержат ID схемы из реестра.
При записи в файл: оверхед схемы амортизируется — одна схема на миллионы записей.
При передаче по сети через Schema Registry: оверхед — 5 байт на сообщение (magic + ID).
Без Schema Registry: пришлось бы передавать полную JSON-схему в каждом сообщении, что сделало бы Avro непригодным для messaging.
Типизация
Protobuf — статически типизирован через сгенерированные классы. Avro Generic — динамически типизирован через
GenericRecord. Avro Specific — статически типизирован, но сгенерированные классы менее удобны, чем в Protobuf (нет Builder pattern, поля — package-private или с простыми геттерами/сеттерами).Schema evolution
Avro имеет более мощную и гибкую систему schema evolution благодаря разделению reader/writer schema. Protobuf поддерживает forward/backward compatibility через правила добавления/удаления полей, но не имеет встроенного механизма разрешения различий между схемами на уровне runtime — это должно быть реализовано прикладным кодом.
Производительность
Protobuf быстрее Avro по нескольким причинам:
Сгенерированный код Protobuf содержит жёстко закодированные инструкции сериализации (switch по тегу). Avro Generic использует рефлексию-подобный обход схемы.
Protobuf не требует schema resolution при чтении (если reader/writer schema совпадают). Avro всегда выполняет resolution, даже когда схемы идентичны.
Protobuf использует immutable generated объекты с zero-copy строками. Avro Generic использует mutable
GenericRecord с Object[].Однако Avro Specific API сопоставим по производительности с Protobuf, а в некоторых сценариях (большие файлы с блоковым сжатием) может быть быстрее за счёт эффективной интеграции с Hadoop/S3.
Экосистема
Protobuf: доминирует в gRPC, микросервисах, мобильной разработке. Поддержка в Kubernetes (etcd), Envoy, многих облачных провайдерах.
Avro: доминирует в Hadoop, Spark, Kafka, data lakes (Delta Lake, Iceberg). Интеграция с Hive, Presto/Trino, Flink.
#Java #для_новичков #beginner #IO #NIO #Serialize #Avro
👍4