Java for Beginner
871 subscribers
1.01K photos
275 videos
14 files
1.68K links
Канал от новичков для новичков!
Изучайте Java вместе с нами!
Здесь мы обмениваемся опытом и постоянно изучаем что-то новое!

Наш YouTube канал - https://www.youtube.com/@Java_Beginner-Dev

Наш канал на RUTube - https://rutube.ru/channel/37896292/
Download Telegram
[Совет по Java #070]

Тема: PriorityQueue не гарантирует порядок при итерации — только при извлечении (poll()).

Проблема: PriorityQueue — это реализация очереди на основе бинарной кучи (heap). Она упорядочивает элементы согласно их естественному порядку или переданному компаратору, но порядок гарантируется только при извлечении элементов через poll()remove() или peek() (который показывает минимальный элемент).

Однако метод iterator() и все операции обхода (например, forEachtoArray()) не следуют порядку приоритета. Итератор проходит по внутреннему массиву кучи в том порядке, в котором элементы физически расположены в памяти, а этот порядок не является отсортированным. Многие разработчики ошибочно полагают, что итерация по PriorityQueue вернёт элементы в порядке возрастания приоритета, и используют её в циклах, не догадываясь, что получают случайный порядок. Это приводит к логическим ошибкам, особенно при выводе содержимого, сериализации или агрегации данных.

Решение: Для получения элементов в порядке приоритета всегда используйте метод poll() в цикле, который удаляет и возвращает наименьший (или наибольший, в зависимости от компаратора) элемент. Если нужно просто просмотреть элементы без удаления, скопируйте очередь в массив и отсортируйте его, либо создайте новую PriorityQueue и последовательно извлекайте элементы.

Для отладки и логирования используйте toArray() и затем сортируйте, либо преобразуйте в список и применяйте Collections.sort(). Никогда не полагайтесь на порядок итератора, так как он специфичен для реализации и может меняться между версиями JDK.
import java.util.*;

public class PriorityQueueOrder {

public static void main(String[] args) {
PriorityQueue<Integer> pq = new PriorityQueue<>();
pq.add(5);
pq.add(1);
pq.add(3);
pq.add(2);
pq.add(4);

//Итератор не гарантирует порядок
System.out.print("Итерация (неупорядоченно): ");
for (Integer i : pq) {
System.out.print(i + " "); // Может вывести 1, 2, 3, 4, 5, но это не гарантировано
}
System.out.println();

//poll() выдает элементы в правильном порядке
System.out.print("Извлечение через poll(): ");
while (!pq.isEmpty()) {
System.out.print(pq.poll() + " "); // 1, 2, 3, 4, 5
}
System.out.println();

// Восстанавливаем очередь
pq.addAll(Arrays.asList(5, 1, 3, 2, 4));

//Просмотр без удаления: копия и сортировка
List<Integer> sorted = new ArrayList<>(pq);
Collections.sort(sorted);
System.out.println("Копия + сортировка: " + sorted);

//Альтернатива: извлечение с сохранением (создаем копию)
PriorityQueue<Integer> copy = new PriorityQueue<>(pq);
List<Integer> inOrder = new ArrayList<>();
while (!copy.isEmpty()) {
inOrder.add(copy.poll());
}
System.out.println("Из копии: " + inOrder);
}
}


Объяснение:
 PriorityQueue хранит элементы в массиве, где позиция родителя всегда меньше (или больше) дочерних элементов согласно компаратору.

Это свойство кучи обеспечивает быстрый доступ к минимальному элементу за O(1) и извлечение за O(log n). Однако хранение в виде кучи не подразумевает полной сортировки массива: порядок элементов может быть любым, главное, чтобы выполнялось условие кучи (например, каждый родитель меньше своих детей).

Итератор обходит массив по индексам, следуя физическому расположению, которое не сохраняет глобальный порядок. Поэтому единственный способ получить элементы в порядке приоритета — многократно вызывать poll(), который удаляет корень и восстанавливает свойство кучи. При проектировании API, возвращающего PriorityQueue, обязательно документируйте, что итерация не гарантирует упорядоченность, и предоставляйте методы для извлечения отсортированных данных.


#Java #советы
👍4
Раздел 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.
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 не требует кодогенерации. Данные представляются как 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.
Изменение порядка полей: разрешено — поля идентифицируются по имени, а не по позиции.
Изменение имени поля: эквивалентно удалению старого и добавлению нового. Для сохранения совместимости используются алиасы (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[]

Вызов 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
[Совет по Java #071]

Тема: Collections.synchronizedList синхронизирует только отдельные методы, но не составные операции (например, итерация).

Проблема: Фабричный метод Collections.synchronizedList(List) возвращает потокобезопасную обертку, в которой все публичные методы (addgetremovesize и др.) синхронизированы на внутреннем мьютексе (обычно на самом объекте списка).

Это обеспечивает атомарность каждого отдельного вызова. Однако составные операции, состоящие из нескольких последовательных вызовов методов, не являются атомарными.

Например, проверка размера и последующее удаление элемента: if (!list.isEmpty()) list.remove(0). Между этими двумя вызовами другой поток может изменить список, и операция удалит не тот элемент или выбросит исключение. Особенно опасна итерация с помощью for-each или Iterator. Итератор, полученный из synchronizedList, не синхронизирован сам по себе, и если во время обхода другой поток модифицирует список, возникает ConcurrentModificationException.

Даже если обход происходит в одном потоке, но метод синхронизирован только на уровне методов, итерация не защищена, так как она выполняется вне синхронизированного контекста.

Решение: Для безопасного выполнения составных операций необходимо вручную синхронизироваться на том же объекте-мониторе, который используется оберткой.

Обычно это сам объект списка, возвращенный Collections.synchronizedList. Рекомендуется использовать блок synchronized(syncList) { ... }, внутри которого выполнять все необходимые вызовы, включая итерацию. Если вы используете синхронизированную обертку, всегда синхронизируйтесь на ней для составных операций. Альтернативно, для сценариев с частыми чтениями и редкими модификациями подходит CopyOnWriteArrayList, который обеспечивает безопасную итерацию без дополнительной синхронизации.

Для высококонкурентных структур, где итерация должна быть потокобезопасной, рассмотрите ConcurrentLinkedDeque или ConcurrentSkipListSet, но они не реализуют интерфейс List. В любом случае, явная синхронизация составных операций остается ответственностью разработчика.
public class SynchronizedListIteration {

public static void main(String[] args) {
List<String> list = Collections.synchronizedList(new ArrayList<>());
list.add("A");
list.add("B");
list.add("C");

//Опасная итерация без синхронизации
// Может выбросить ConcurrentModificationException, если другой поток модифицирует список
for (String s : list) {
System.out.println(s);
}

//Правильная итерация с синхронизацией на списке
synchronized (list) {
for (String s : list) {
System.out.println(s);
}
}

//Неатомарная составная операция
// Между size() и remove() другой поток может изменить список
if (!list.isEmpty()) {
list.remove(0); // Может удалить не тот элемент
}

//Правильная составная операция в синхронизированном блоке
synchronized (list) {
if (!list.isEmpty()) {
String first = list.remove(0);
System.out.println("Removed: " + first);
}
}

// Для итерации с удалением элементов также требуется синхронизация
synchronized (list) {
Iterator<String> it = list.iterator();
while (it.hasNext()) {
if (it.next().equals("B")) {
it.remove(); // Безопасно внутри блока
}
}
}
}
}


Объяснение:
 Collections.synchronizedList возвращает экземпляр внутреннего класса SynchronizedList, который делегирует все вызовы к оригинальному списку, оборачивая их в synchronized(mutex).

Мутексом по умолчанию является сам объект обертки. Поэтому, синхронизируясь на list, мы используем тот же самый монитор, что обеспечивает взаимное исключение с методами обертки. Без явной синхронизации составные операции не являются потокобезопасными, так как между отдельными методами блокировка освобождается.

Итерация особенно уязвима, потому что итератор не захватывает блокировку на весь период обхода. Синхронизация всего блока гарантирует, что никто не сможет модифицировать список во время итерации или выполнения нескольких операций.


#Java #советы
👍4
Раздел 12. Работа с файлами, I/O и сетью (NIO.2)

Глава 4. Работа с сетью (Socket, HTTP, HttpClient)


Основы сетевого взаимодействия: IP-адреса, порты, протоколы

Сетевое взаимодействие — это фундамент, на котором строится любое распределённое приложение. Будь то веб-сервис, микросервисная архитектура, потоковое видео или онлайн-игра, всё сводится к передаче данных между двумя точками через сеть.
В Java сетевая функциональность реализована в пакетах java.net и java.nio.channels, а понимание нижележащих протоколов — обязательное требование для разработчика, который проектирует производительные и надёжные системы.
Протокол — это формальный набор правил, определяющий формат, порядок и обработку данных при их передаче между устройствами. Протоколы работают на разных уровнях абстракции: от физического уровня (электрические сигналы в кабеле) до прикладного (HTTP, gRPC). В Java-разработке нас интересуют прежде всего транспортный и прикладной уровни модели TCP/IP.

TCP/IP — это стек протоколов, лежащий в основе интернета. Вместо семиуровневой модели OSI в инженерной практике используется упрощённая четырёхуровневая модель: канальный уровень (Ethernet, Wi-Fi), сетевой (IP), транспортный (TCP, UDP) и прикладной (HTTP, DNS, SSH).

IP-адрес: логический адрес в сети
IP-адрес (Internet Protocol address) — это числовой идентификатор, присваиваемый каждому устройству, подключённому к сети, работающей по протоколу IP.
IP-адрес выполняет две функции: идентификация хоста (устройства) и маршрутизация — определение пути, по которому данные должны добраться до получателя.

IPv4

IPv4 (Internet Protocol version 4) — четвёртая версия протокола, введённая в 1981 году. Адрес IPv4 представляет собой 32-битное число, разделённое на четыре октета (группы по 8 бит), записываемых в десятичной нотации через точку: 192.168.1.1.
32 бита дают теоретический лимит в 4 294 967 296 уникальных адресов. На практике адресное пространство исчерпано из-за неэффективного распределения в ранние годы интернета и роста числа устройств. Чтобы продлить жизнь IPv4, были введены технологии NAT (Network Address Translation — преобразование сетевых адресов), позволяющие множеству устройств в локальной сети использовать один внешний IP-адрес, и CIDR (Classless Inter-Domain Routing — бесклассовая междоменная маршрутизация), которая позволяет гибко делить адресное пространство на подсети произвольного размера.
В Java адрес IPv4 представляется классом Inet4Address, наследником InetAddress.
import java.net.InetAddress;
import java.net.UnknownHostException;

public class IpExample {
public static void main(String[] args) throws UnknownHostException {
// Разрешение доменного имени в IP-адрес через DNS
InetAddress address = InetAddress.getByName("www.example.com");
System.out.println(address.getHostAddress()); // 93.184.216.34
System.out.println(address instanceof java.net.Inet4Address); // true
}
}

Метод getByName() выполняет DNS-резолюцию (Domain Name System resolution — преобразование доменного имени в IP-адрес). Это блокирующий вызов: поток приостанавливается до получения ответа от DNS-сервера. В современном Java для неблокирующей работы используется java.net.InetAddress в сочетании с CompletableFuture или HttpClient.

IPv6

IPv6 (Internet Protocol version 6) — шестая версия, разработанная как преемник IPv4. Адрес IPv6 — это 128-битное число, записываемое в шестнадцатеричной нотации через двоеточия: 2001:0db8:85a3:0000:0000:8a2e:0370:7334. Для краткости ведущие нули в группе можно опускать, а одну последовательность из нулевых групп — заменять двойным двоеточием: 2001:db8:85a3::8a2e:370:7334.
128 бит дают 3.4 × 10^38 адресов — количество, достаточное для присвоения IP каждому атому на поверхности Земли. IPv6 устраняет необходимость в NAT, упрощает маршрутизацию за счёт фиксированного размера заголовка (в отличие от переменного в IPv4) и встроено поддерживает IPsec — протокол шифрования и аутентификации на сетевом уровне.
В Java IPv6 представлен классом Inet6Address. JVM автоматически предпочитает IPv6, если хост поддерживает оба стека (свойство java.net.preferIPv4Stack и java.net.preferIPv6Addresses управляют этим поведением).
InetAddress ipv6 = InetAddress.getByName("::1");  // loopback в IPv6
System.out.println(ipv6 instanceof java.net.Inet6Address); // true

Loopback (петлевой адрес) — это специальный адрес, который всегда указывает на текущий хост. В IPv4 это 127.0.0.1, в IPv6 — ::1. Пакеты, отправленные на loopback, не выходят за пределы сетевого стека ОС и не попадают в физический интерфейс.


Специальные адреса


Приватные адреса (RFC 1918): 10.0.0.0/8172.16.0.0/12192.168.0.0/16 — используются только внутри локальных сетей, не маршрутизируются в интернете.
Multicast224.0.0.0/4 в IPv4, ff00::/8 в IPv6 — адреса для групповой рассылки, когда один пакет доставляется множеству получателей.
Broadcast (широковещательный адрес): в IPv4 255.255.255.255 или адрес подсети с единицами в хостовой части — отправка пакета всем устройствам в сегменте сети. В IPv6 broadcast отсутствует, его заменяет multicast.
Маршрутизация (routing) — это процесс выбора пути для сетевого пакета от источника к получателю через промежуточные узлы (роутеры). Каждый роутер принимает решение на основе таблицы маршрутизации и заголовка IP-пакета.

Порт: точка входа в приложение
Если IP-адрес идентифицирует устройство в сети, то порт (port) идентифицирует конкретное приложение или процесс на этом устройстве. Порт — это 16-битное число от 0 до 65535, которое добавляется к IP-адресу, образуя полный адрес конечной точки связи.
Сокет (socket) — это абстракция, представляющая одну конечную точку сетевого соединения. В TCP сокет однозначно определяется четвёркой: IP-источника, порт-источника, IP-назначения, порт-назначения. В Java сокет представлен классами java.net.Socket (клиент) и java.net.ServerSocket (сервер).


Диапазоны портов


Well-known ports (0–1023): зарезервированы IANA для стандартных протоколов. Например: 22 — SSH, 80 — HTTP, 443 — HTTPS, 53 — DNS. На Unix-подобных системах привязка к этим портам требует прав суперпользователя (root).
Registered ports (1024–49151): зарегистрированы IANA для конкретных сервисов, но не требуют root. Например: 3306 — MySQL, 5432 — PostgreSQL, 8080 — альтернативный HTTP.
Dynamic / Private / Ephemeral ports (49152–65535): временные порты, которые ОС назначает клиентским приложениям при исходящих соединениях. Когда ваш браузер открывает соединение с сервером на порту 443, ОС выбирает свободный эфемерный порт из этого диапазона для локальной стороны соединения.
Эфемерный порт (ephemeral port) — это временный порт, автоматически выделяемый стеком TCP/IP операционной системы для исходящего соединения. После закрытия соединения порт переходит в состояние TIME_WAIT и освобождается через несколько минут.

import java.net.InetSocketAddress;

// InetSocketAddress объединяет IP-адрес и порт
InetSocketAddress serverAddr = new InetSocketAddress("192.168.1.100", 8080);
InetSocketAddress localAddr = new InetSocketAddress("0.0.0.0", 0); // порт 0 = эфемерный

Передача 0 в качестве порта при создании Socket или ServerSocket означает "выбрать любой свободный эфемерный порт". После привязки реальный порт можно узнать через socket.getLocalPort().


#Java #для_новичков #beginner #IO #NIO #Ip #Port #TCP #UDP
👍4
TCP: Transmission Control Protocol

TCP (Transmission Control Protocol — протокол управления передачей) — это основной транспортный протокол интернета, обеспечивающий надёжную, упорядоченную и проверенную доставку потока байтов между приложениями. TCP работает поверх IP и добавляет к нему механизмы контроля, которых у IP нет.

Трёхэтапное рукопожатие (Three-way Handshake)

Перед передачей данных TCP устанавливает соединение через процедуру, состоящую из трёх сообщений:
SYN (Synchronize — синхронизация): клиент отправляет серверу сегмент с флагом SYN и случайным начальным номером последовательности (ISN, Initial Sequence Number). ISN выбирается случайно для защиты от атак типа TCP sequence prediction.
SYN-ACK (Synchronize-Acknowledge): сервер получает SYN, выделяет ресурсы для соединения (создаёт запись в таблице соединений, буферы приёма и передачи), генерирует свой ISN и отправляет ответ с флагами SYN и ACK. Поле ACK содержит номер последовательности клиента плюс один, подтверждая получение SYN.
ACK (Acknowledge): клиент получает SYN-ACK, подтверждает получение отправкой ACK (номер последовательности сервера плюс один), и переходит в состояние ESTABLISHED. Сервер, получив ACK, также переходит в ESTABLISHED. Соединение установлено, и начинается передача данных.
Флаг в контексте TCP — это битовое поле в заголовке сегмента, управляющее поведением протокола. Флаги: SYN (начало соединения), ACK (подтверждение), FIN (завершение), RST (сброс), PSH (немедленная передача данных приложению), URG (срочные данные).

// Псевдокод трёхэтапного рукопожатия на уровне сокетов Java
// На самом деле это происходит внутри ядра ОС, Java лишь инициирует

// Клиент
Socket clientSocket = new Socket();
clientSocket.connect(serverAddress); // Внутри ядра: отправка SYN -> получение SYN-ACK -> отправка ACK

// Сервер
ServerSocket serverSocket = new ServerSocket(8080);
Socket acceptedSocket = serverSocket.accept(); // Внутри ядра: получение SYN -> отправка SYN-ACK -> получение ACK


Гарантия доставки

TCP гарантирует доставку данных через механизм подтверждений (acknowledgments, ACK). Каждый отправленный сегмент должен быть подтверждён получателем. Если отправитель не получает ACK в течение RTO (Retransmission Timeout — таймаута повторной передачи), он повторно отправляет сегмент. RTO вычисляется динамически на основе RTT (Round-Trip Time — времени оборота пакета до получателя и обратно).
Кроме того, TCP использует контрольную сумму (checksum) — 16-битное значение, вычисляемое над заголовком и данными сегмента. Если получатель вычисляет checksum, которая не совпадает с переданной, сегмент отбрасывается, и отправитель, не получив ACK, повторит передачу.

Упорядоченность
IP не гарантирует порядок доставки пакетов: они могут идти разными маршрутами и приходить вперемешку. TCP решает эту проблему через номера последовательности (sequence numbers). Каждый байт в потоке имеет порядковый номер. Получатель собирает сегменты в правильном порядке по этим номерам и передаёт упорядоченный поток приложению. Если сегмент приходит с разрывом (например, получены байты 1–1000 и 2001–3000, но 1001–2000 потеряны), получатель буферизует пришедшие данные и ждёт недостающий сегмент, не передавая приложению неполный поток.

Управление потоком (Flow Control)

Управление потоком — это механизм, предотвращающий переполнение буфера получателя. В заголовке TCP есть поле Window Size (размер окна) — количество байтов, которое получатель готов принять. Отправитель не отправляет больше данных, чем позволяет окно. Если приложение-получатель медленно читает данные из сокета, окно уменьшается до нуля (Zero Window), и отправитель приостанавливает передачу до обновления окна.

Контроль перегрузки (Congestion Control)

Контроль перегрузки — это механизм, предотвращающий перегрузку сети. В отличие от управления потоком, которое защищает получателя, контроль перегрузки защищает сеть. TCP медленно увеличивает скорость передачи (slow start), экспоненциально наращивая количество сегментов, пока не обнаружит потерю (по таймауту или дублированным ACK). При потере скорость резко снижается. Современные алгоритмы: CUBIC (по умолчанию в Linux), BBR (разработан Google, учитывает пропускную способность и RTT).

Завершение соединения: четырёхэтапное прощание
Закрытие TCP соединения требует четырёх сообщений:
FIN: инициатор (например, клиент) отправляет FIN, сообщая, что он больше не будет отправлять данные.
ACK: получатель подтверждает FIN.
FIN: получатель, завершив передачу своих данных, отправляет свой FIN.
ACK: инициатор подтверждает FIN.
После этого соединение переходит в состояние TIME_WAIT на стороне инициатора закрытия. Это состояние длится обычно 2 × MSL (Maximum Segment Lifetime, 60–120 секунд) и необходимо для надёжного закрытия: если последний ACK потеряется, получатель повторит FIN, и инициатор должен быть готов отправить ACK повторно. TIME_WAIT также предотвращает попадание задержанных сегментов из старого соединения в новое с теми же параметрами.


UDP: User Datagram Protocol


UDP (User Datagram Protocol — протокол пользовательских датаграмм) — это минималистичный транспортный протокол, работающий поверх IP. В отличие от TCP, UDP не устанавливает соединения, не гарантирует доставку, не обеспечивает упорядоченность и не контролирует перегрузку.

Структура и принцип работы
UDP-заголовок занимает всего 8 байт (против 20+ байт у TCP) и содержит четыре поля: порт источника, порт назначения, длина датаграммы и контрольную сумму. После формирования датаграммы ОС немедленно отправляет её в сеть. Нет рукопожатия, нет ожидания подтверждений, нет повторных передач.
Датаграмма (datagram) — это самостоятельная единица данных, передаваемая через сеть без установления соединения. Каждая датаграмма независима: она может идти своим маршрутом, прийти в другом порядке или не прийти вовсе.

// UDP-сервер: приём датаграмм
DatagramSocket serverSocket = new DatagramSocket(9876);
byte[] buffer = new byte[1024];
DatagramPacket packet = new DatagramPacket(buffer, buffer.length);
serverSocket.receive(packet); // Блокирующее ожидание датаграммы
String received = new String(packet.getData(), 0, packet.getLength());
// UDP-клиент: отправка датаграммы
DatagramSocket clientSocket = new DatagramSocket(); // эфемерный порт
byte[] sendData = "Hello, UDP".getBytes();
DatagramPacket sendPacket = new DatagramPacket(
sendData, sendData.length,
InetAddress.getByName("localhost"), 9876
);
clientSocket.send(sendPacket); // Отправка без установления соединения


Потеря пакетов

UDP не обнаруживает потерю пакетов на транспортном уровне. Если датаграмма потерялась в сети, получатель об этом не узнает, а отправитель не повторит передачу. Это делает UDP непригодным для передачи файлов или текстовых данных без дополнительной логики поверх протокола. Однако именно отсутствие гарантий даёт UDP преимущество в скорости и задержке (latency).
Задержка (latency) — это время, которое требуется пакету для прохождения от отправителя к получателю. В TCP задержка включает время рукопожатия, ожидание ACK, повторные передачи. В UDP задержка минимальна: данные отправляются сразу.



Преимущества UDP


Низкая задержка: отсутствие рукопожатия и ожидания ACK означает, что первый байт данных уходит в сеть немедленно.
Отсутствие head-of-line blocking: в TCP, если один сегмент потерян, все последующие сегменты задерживаются в буфере до получения недостающего. В UDP каждая датаграмма независима.
Multicast и broadcast: UDP естественно поддерживает отправку одной датаграммы множеству получателей, что невозможно в TCP (TCP — исключительно one-to-one).
Меньший оверхед: 8-байтный заголовок против 20+ байт у TCP, меньше CPU-нагрузки на обработку.

Когда что использовать
Выбор между TCP и UDP определяется требованиями приложения к надёжности, задержке и порядку данных.


#Java #для_новичков #beginner #IO #NIO #Ip #Port #TCP #UDP
👍4
HTTP и HTTPS

HTTP (HyperText Transfer Protocol) и HTTPS (HTTP Secure) работают поверх TCP. Каждый HTTP-запрос требует надёжной доставки: потеря даже одного байта HTML, CSS или JSON приведёт к повреждению документа. TCP гарантирует, что весь ответ сервера дойдёт до клиента целиком и в правильном порядке.
В HTTP/1.1 и HTTP/2 поверх одного TCP-соединения мультиплексируется множество запросов. В HTTP/3 (разработан Google как QUIC) используется UDP с собственным механизмом надёжности, реализованным в пространстве пользователя, а не в ядре ОС. Это позволяет устранить head-of-line blocking на уровне TCP и сократить задержку за счёт отсутствия рукопожатия при смене сети (например, переход с Wi-Fi на мобильный интернет).

DNS

DNS
 (Domain Name System — система доменных имён) изначально использовала UDP для запросов. UDP идеален здесь, потому что запросы и ответы обычно помещаются в одну датаграмму (до 512 байт в классическом DNS, до 4096 с EDNS). Низкая задержка критична: каждое обращение к веб-странице может требовать нескольких DNS-запросов, и добавление TCP-рукопожатия замедлило бы загрузку. Если ответ DNS превышает размер датаграммы или требуется надёжность (zone transfer), DNS переключается на TCP.

Видеозвонки и потоковое видео
Видеозвонки (WebRTC, Zoom, Skype) и потоковое видео в реальном времени (Twitch, YouTube Live) используют UDP. Задержка здесь важнее полноты: если один кадр видео потерян, лучше показать следующий кадр сразу, чем ждать повторную передачу старого. Повторная передача привела бы к "заморозке" изображения, что хуже, чем кратковременное артефактное изображение. Протоколы вроде RTP (Real-time Transport Protocol) работают поверх UDP и добавляют нумерацию пакетов и временные метки, но не повторную передачу.
WebRTC (Web Real-Time Communication) — это открытый фреймворк для организации передачи потоковых данных, голосовых вызовов и видеоконференций между браузерами и приложениями без установки плагинов. Он использует UDP для передачи медиа-потоков через протокол SRTP (Secure RTP).


Онлайн-игры

Многопользовательские онлайн-игры используют UDP для передачи состояния игрового мира (позиции игроков, выстрелы). Потеря одного пакета с координатами не критична: следующий пакет через 16 мс (60 FPS) содержит актуальные координаты. TCP здесь непригоден, потому что задержка в 100 мс на повторную передачу старого пакета сделала бы игру неиграбельной. Однако важные команды (покупка предмета, авторизация) могут отправляться по TCP или с собственным механизмом подтверждений поверх UDP.

IoT и телеметрия
Устройства интернета вещей (IoT — Internet of Things), датчики и системы телеметрии часто используют UDP или MQTT (Message Queuing Telemetry Transport) поверх TCP. MQTT — лёгкий протокол публикации-подписки, оптимизированный для нестабильных сетей и устройств с ограниченными ресурсами. В чистом виде UDP используется, когда потеря отдельных показаний датчика допустима, а экономия энергии и трафика приоритетна.


Java NIO: эффективная работа с сетью

Классические блокирующие сокеты java.net.Socket и ServerSocket создают отдельный поток на каждое соединение. При 10 000 одновременных клиентах это означает 10 000 потоков — огромный расход памяти на стеки и переключение контекста CPU.

NIO (New I/O, пакет java.nio) решает эту проблему через неблокирующие каналы (non-blocking channels) и селекторы (selectors).
public class NioServer {
public static void main(String[] args) throws Exception {
// ServerSocketChannel — неблокирующий серверный канал
ServerSocketChannel serverChannel = ServerSocketChannel.open();
serverChannel.bind(new InetSocketAddress(8080));
serverChannel.configureBlocking(false); // неблокирующий режим

// Selector — мультиплексор, который мониторит множество каналов
Selector selector = Selector.open();
serverChannel.register(selector, SelectionKey.OP_ACCEPT);

while (true) {
// select() блокирует до появления событий на любом из зарегистрированных каналов
selector.select();
Set<SelectionKey> keys = selector.selectedKeys();
Iterator<SelectionKey> it = keys.iterator();

while (it.hasNext()) {
SelectionKey key = it.next();
it.remove(); // обязательно удаляем обработанный ключ

if (key.isAcceptable()) {
// Новое входящее соединение
SocketChannel client = serverChannel.accept();
client.configureBlocking(false);
client.register(selector, SelectionKey.OP_READ);
} else if (key.isReadable()) {
// Данные доступны для чтения
SocketChannel client = (SocketChannel) key.channel();
ByteBuffer buffer = ByteBuffer.allocate(1024);
int bytesRead = client.read(buffer);
if (bytesRead == -1) {
client.close(); // соединение закрыто клиентом
} else {
buffer.flip(); // переключаем буфер в режим чтения
// обработка данных...
}
}
}
}
}
}

Канал (Channel) в NIO — это абстракция для ввода-вывода, которая отличается от потоков (Stream) тем, что поддерживает неблокирующий режим и может работать в обе стороны (чтение и запись). SocketChannel — это канал для TCP-соединения, DatagramChannel — для UDP.

Селектор (Selector) — это механизм мультиплексирования ввода-вывода. Один поток может мониторить множество каналов, ожидая событий (готовность к чтению, записи, подключению). Это позволяет обслуживать тысячи соединений одним или несколькими потоками вместо одного потока на соединение.

SelectionKey — это токен, представляющий зарегистрированный канал в селекторе. Он хранит информацию о канале, селекторе и интересующих операциях (OP_ACCEPT, OP_CONNECT, OP_READ, OP_WRITE).



#Java #для_новичков #beginner #IO #NIO #Ip #Port #TCP #UDP
👍3
Путь байтов в памяти JVM при сетевом взаимодействии

Сетевой ввод-вывод в Java проходит через сложную цепочку от сетевой карты до кучи JVM и обратно. Понимание этого пути критично для оптимизации производительности и управления памятью.

1. От сетевой карты до ядра ОС
Когда пакет приходит на сетевой интерфейс, NIC (Network Interface Controller — сетевой контроллер) генерирует аппаратное прерывание. Драйвер сетевой карты в ядре ОС обрабатывает прерывание, копирует пакет из буфера NIC в кольцевой буфер (ring buffer) ядра — структуру данных фиксированного размера в оперативной памяти, организованную по принципу FIFO (First In, First Out — первым пришёл, первым вышел). Затем стек TCP/IP ядра обрабатывает пакет: проверяет IP-заголовок, TCP-заголовок, вычисляет checksum, обновляет состояние соединения.
Данные приложения из TCP-сегментов накапливаются в сокет-буфере (socket receive buffer) ядра — области памяти ядра, выделенной для конкретного сокета. Размер этого буфера настраивается через опции сокета SO_RCVBUF.

2. Из ядра в JVM: read() и ByteBuffer
Когда Java-приложение вызывает SocketChannel.read(buffer), происходит переход из пользовательского пространства в пространство ядра (system call). Данные копируются из сокет-буфера ядра в буфер JVM.
Здесь критично различие между двумя типами ByteBuffer:
HeapByteBuffer — буфер, размещённый в куче JVM (Heap). При чтении из канала данные сначала копируются во временный буфер в нативной памяти (выделенный JNI), а затем — в массив byte[] в куче. Это двойное копирование снижает производительность.

DirectByteBuffer — буфер, размещённый в нативной памяти (native memory, вне кучи JVM). При чтении данные копируются напрямую из сокет-буфера ядра в нативную память, минуя кучу. Это устраняет одно копирование и позволяет использовать zero-copy техники на уровне ОС.

// HeapByteBuffer: данные в куче, медленнее при I/O
ByteBuffer heapBuffer = ByteBuffer.allocate(1024);

// DirectByteBuffer: данные в нативной памяти, быстрее при I/O
ByteBuffer directBuffer = ByteBuffer.allocateDirect(1024);

Zero-copy — это техника, при которой данные не копируются между буферами в пользовательском пространстве, а передаются по ссылке или напрямую из ядра в целевой буфер. В Java FileChannel.transferTo() и transferFrom() используют zero-copy для передачи данных между файлами и сокетами без промежуточного буфера в пользовательском пространстве.



3. Работа GC с сетевыми буферами

Объект ByteBuffer в Java — это обёртка. Для HeapByteBuffer он содержит ссылку на byte[] в куче. Для DirectByteBuffer он содержит адрес в нативной памяти (поле address типа long). Сам объект DirectByteBuffer — маленький объект в куче (десятки байтов), но он ссылается на большой блок нативной памяти (мегабайты).
Cleaner — это механизм в JVM, использующий PhantomReference (фантомную ссылку) для регистрации callback-функции, которая вызывается, когда объект становится пригодным для сборки мусора. Когда DirectByteBuffer собирается GC, Cleaner вызывает нативный метод free(), который освобождает блок нативной памяти через Unsafe.freeMemory().

PhantomReference — это тип ссылки в Java, который позволяет узнать, что объект был финализирован сборщиком мусора, но ещё не удалён. В отличие от WeakReference и SoftReference, phantom reference не даёт доступа к объекту, но гарантирует, что объект уже недостижим.

Однако Cleaner срабатывает только при сборке мусора. Если приложение активно выделяет DirectByteBuffer, но GC не запускается (потому что куча не заполнена), нативная память может исчерпать лимит ОС, вызывая OutOfMemoryError: Direct buffer memory. Чтобы избежать этого, важно явно вызывать cleaner.clean() (через рефлексию) или использовать пулы буферов (например, Netty ByteBuf), которые переиспользуют выделенную память.

4. Путь отправки данных
При записи (SocketChannel.write(buffer)) путь обратный:
Данные из ByteBuffer (куча или нативная память) копируются в сокет-буфер отправки ядра (socket send buffer).
Стек TCP ядра сегментирует данные, добавляет заголовки TCP/IP, вычисляет checksum.
Сегменты попадают в очередь передачи NIC и отправляются в сеть.
При получении ACK от получателя ядро освобождает соответствующие сегменты из send buffer.
Если send buffer заполнен (получатель не успевает читать), вызов write() вернёт 0 записанных байт в неблокирующем режиме или заблокирует поток в блокирующем режиме.

5. Объекты сетевого стека в JVM
При использовании NIO с Selector в куче JVM хранятся:
Объект Selector — содержит ссылку на нативную структуру pollfd (Linux) или WSAPOLL (Windows).
Объекты SelectionKey — хранятся в HashSet внутри Selector, по одному на каждый зарегистрированный канал.
Объекты SocketChannel и ServerSocketChannel — обёртки над файловыми дескрипторами ОС.

Эти объекты лёгкие и живут в Young Generation. При закрытии канала (channel.close()) файловый дескриптор освобождается нативным кодом, а Java-объект становится мусором для следующей Minor GC. Однако если забыть закрыть канал, дескриптор утечёт на уровне ОС — это утечка файлового дескриптора (file descriptor leak), которая при масштабе приводит к исчерпанию лимита ulimit -n и невозможности открыть новые соединения.
// Правильный паттерн: try-with-resources для автоматического закрытия
// Closeable — интерфейс, требующий реализации метода close()
try (SocketChannel channel = SocketChannel.open(new InetSocketAddress("host", 80))) {
ByteBuffer buffer = ByteBuffer.allocateDirect(4096);
channel.read(buffer);
} // channel.close() вызывается автоматически, освобождая файловый дескриптор



#Java #для_новичков #beginner #IO #NIO #Ip #Port #TCP #UDP
👍3
[Совет по Java #072]

Тема: ConcurrentHashMap при итерации не бросает ConcurrentModificationException — она weakly consistent.

Проблема: В отличие от обычных коллекций (HashMapArrayList), итераторы которых являются fail-fast и выбрасывают ConcurrentModificationException при обнаружении структурных изменений во время итерации, ConcurrentHashMap использует стратегию слабой согласованности (weakly consistent).

Это означает, что итератор не гарантирует, что отразит все изменения, произошедшие во время итерации, и не бросает исключения при конкурентных модификациях. Разработчики, привыкшие к fail-fast итераторам, ожидают, что при изменении карты во время обхода возникнет исключение, и могут строить логику, полагаясь на это поведение. Однако в ConcurrentHashMap это не происходит, и отсутствие исключения может замаскировать ошибки конкурентного доступа, приводя к тому, что итерация видит частично обновленные данные или пропускает новые элементы.

Это особенно опасно, когда предполагается, что коллекция статична во время итерации, но на самом деле изменяется другими потоками.

Решение: Осознайте, что итерация по ConcurrentHashMap безопасна в многопоточной среде, но не предоставляет моментального снимка (snapshot). Итератор отражает состояние карты на момент создания итератора, но может видеть изменения, произошедшие после (в зависимости от реализации). Для получения согласованного снимка используйте entrySet().toArray() или создайте копию карты перед итерацией. Если вам нужна строгая согласованность в момент обхода, синхронизируйтесь на самой карте или используйте Collections.synchronizedMap.

Однако для большинства сценариев, где не требуется полная изоляция, поведение ConcurrentHashMap вполне приемлемо и позволяет избежать исключений. Документируйте ожидаемое поведение и не полагайтесь на отсутствие исключений как на гарантию целостности данных.
public class ConcurrentHashMapIteration {

public static void main(String[] args) throws InterruptedException {
ConcurrentHashMap<Integer, String> map = new ConcurrentHashMap<>();
map.put(1, "One");
map.put(2, "Two");
map.put(3, "Three");

//Итерация безопасна и не бросает ConcurrentModificationException
Thread modifier = new Thread(() -> {
map.put(4, "Four"); // модификация во время итерации
map.remove(2);
});

// Итератор создается ДО изменений
var iterator = map.entrySet().iterator();
modifier.start();

// Итерация продолжается, несмотря на изменения
System.out.println("Итерация (weakly consistent):");
while (iterator.hasNext()) {
Map.Entry<Integer, String> entry = iterator.next();
System.out.println(entry.getKey() + " -> " + entry.getValue());
// Может увидеть 4, может не увидеть; 2 может быть или нет
// Исключений не будет
}
modifier.join();

// Для снимка: копирование
Map<Integer, String> snapshot = new ConcurrentHashMap<>(map);
System.out.println("Снимок: " + snapshot);
}
}


Объяснение:
 ConcurrentHashMap использует внутреннюю сегментацию (bucket-level locking) и не блокирует всю структуру при модификации. Итераторы строятся на основе текущего состояния сегментов и не хранят отдельный снимок. Они не проверяют счетчик модификаций (в отличие от fail-fast итераторов), поэтому не бросают ConcurrentModificationException.

Вместо этого они гарантируют, что обход будет продолжаться без исключений, и каждый элемент будет посещен не более одного раза, но при этом элементы, добавленные или удаленные во время итерации, могут быть как видны, так и нет. Такая модель подходит для высококонкурентных систем, где важна производительность, а не строгая изоляция.


#Java #советы
👍4