Java & Spring

Java & Spring

Java и Spring

Заглушки

Заглушка/mock/эмулятор - программа написанная тестировщиком для эмуляции поведения интеграций тестируемой системы.

Зачем писать заглушки? - Эмулировать поведение интеграции, чтобы не пользоваться её экземплярам по причинам:

Нет необходимости использовать при тестах бОльшую часть интеграции

ДОРОГО (живые интеграции себе может себе позволить только сбер, газпром, ВТБ и т.д.)

Чтобы протестировать системы или их/её части отдельно друг от друга, т.е. изолированно.

Интеграция находится у другой компании, поэтому п. b либо у неё в принципе вне других компаний нет экземпляров

Gradle/Maven - Сборка ( пока только Gradle)

Gradle

Основная информация

Gradle — это инструмент для автоматизации сборки и управления зависимостями в проектах. Он используется в разработке программного обеспечения, особенно в экосистемах Java, Kotlin, Android и других JVM-языков, а также в мультиплатформенных проектах.

Основные задачи Gradle:

Сборка проекта (компиляция, упаковка в JAR/WAR/APK и т. д.).

Управление зависимостями (автоматическая загрузка библиотек из репозиториев, например, Maven Central).

Запуск тестов (интеграция с JUnit, TestNG и другими фреймворками).

Кастомизация сборки (настройка этапов сборки под конкретные нужды).

Инкрементальные сборки (ускорение процесса за счёт пересборки только измененных частей).

Интеграция с CI/CD (работа с Jenkins, GitHub Actions и другими системами).

Основные настройки содержатся в файле build.gradle. Пример:

Ключевые команды Gradle:

gradle build — собрать проект.

gradle test — запустить тесты.

gradle run — запустить приложение.

gradle clean — очистить собранные артефакты.

Альтернативы:

Maven (использует XML, менее гибкий, но проще для простых проектов).

Ant (устаревший, низкоуровневый).

Где используется?

Android-разработка (официальный инструмент сборки).

Backend на Java/Kotlin (Spring, Micronaut и др.).

Kotlin Multiplatform.

Build.gradle - подробно

Файл build.gradle (или build.gradle.kts для Kotlin DSL) — это основной конфигурационный файл в проектах, использующих Gradle. В нём определяются настройки сборки, зависимости, плагины и другие параметры проекта.

Основные секции в build.gradle:

Плагины (plugins) -Плагины расширяют функциональность Gradle (например, добавляют поддержку Java, Android, Spring и т. д.).

Настройки проекта (group, version, sourceCompatibility)

group — идентификатор группы (например, com.example).

version — версия проекта (1.0.0, 0.0.1-SNAPSHOT).

sourceCompatibility — версия Java (17, 11).

Репозитории (repositories) - Указывает, откуда Gradle будет скачивать зависимости (библиотеки).

Зависимости (dependencies) - Здесь перечисляются библиотеки, необходимые для проекта. Типы зависимостей:

implementation — основная зависимость.

compileOnly — только для компиляции (не включается в итоговый билд)

runtimeOnly — только для выполнения (не нужна при компиляции).

testImplementation — зависимости для тестов.

Задачи (tasks) - Кастомные задачи Gradle (например, копирование файлов, очистка, кастомизация сборки).

Настройка приложения (application) - Указывает главный класс для запуска (если используется плагин application).

Конфигурация сборки для подпроектов (в мультимодульных проектах). Если проект состоит из нескольких модулей, в корневом build.gradle могут быть общие настройки:

Пример полного build.gradle.kts:

Java Spring

По сути Spring Framework представляет собой просто контейнер внедрения зависимостей, с несколькими удобными слоями (например: доступ к базе данных, прокси, аспектно-ориентированное программирование, RPC, веб-инфраструктура MVC). Это все позволяет вам быстрее и удобнее создавать Java-приложения.

Основные возможности Spring Framework:

IoC (Inversion of Control) и DI (Dependency Injection)

Spring управляет зависимостями между компонентами вместо программиста.

Вместо new MyService() зависимости внедряются через аннотации (@Autowired, @Inject).

Spring MVC

Фреймворк для создания веб-приложений (аналог Jakarta EE Servlet API, но удобнее).

Основан на паттерне Model-View-Controller.

Spring Boot

Надстройка над Spring, которая упрощает настройку и запуск приложений.

Автоконфигурация, встроенный сервер (Tomcat, Jetty), application.properties/application.yml.

Spring Data

Упрощает работу с базами данных (JPA, Hibernate, MongoDB, Redis и др.).

Достаточно объявить интерфейс:

Spring Security

Фреймворк для аутентификации и авторизации (OAuth2, JWT, Basic Auth).

Spring AOP (Aspect-Oriented Programming)

Позволяет добавлять сквозную логику (логирование, транзакции) через аспекты.

Spring REST & WebFlux

Поддержка REST API (@RestController, @GetMapping и т.д.).

WebFlux — реактивный стек для асинхронных приложений.

Spring Cloud

Инструменты для микросервисов (Eureka, Config Server, Hystrix, Zuul).

Пример простого Spring Boot-приложения:

Dependency Injection (DI, Внедрение зависимостей)

Внедрение зависимостей — это один из ключевых паттернов проектирования, который активно используется в Spring Framework и других современных фреймворках. Его суть в том, что зависимости (объекты, которые использует класс) не создаются внутри класса, а "внедряются" извне.

Это делает код:

Гибким (легко менять реализации).

Тестируемым (можно подставлять mock-объекты).

Чистым (меньше связанности между компонентами)

Схема работы:

Без DI (плохой подход) - Класс сам создаёт свои зависимости:

Проблемы:

Жёсткая привязка к PaymentService.

Сложно тестировать (например, нельзя подменить PaymentService на заглушку).

С DI (правильный подход) - Зависимости передаются извне (через конструктор, сеттер или поле).

Внедрение через конструктор (наиболее предпочтительный способ)

Плюсы:

Зависимость явная (видна в сигнатуре конструктора).

Нельзя создать OrderService без PaymentService.

Подходит для неизменяемых (immutable) зависимостей (final поле).

Как Spring это обрабатывает?

Если PaymentService — это Spring-бин, Spring автоматически внедрит его:

Внедрение через сеттер (setter injection)

Когда использовать?

Если зависимость опциональная (может быть null).

Если нужно переопределить зависимость в runtime (редкий случай).

Внедрение через поле (field injection) — не рекомендуется

Минусы:

Скрытая зависимость (не видна в конструкторе/сеттере).

Сложно тестировать без Spring (нужен ReflectionTestUtils).

Нельзя сделать поле final.

Типы DI в Spring

Автоматическое связывание (Autowiring)

Spring сам находит подходящие бины и внедряет их.

Как Spring понимает, что внедрять?

По типу (UserService).

Если есть несколько реализаций — можно уточнить с помощью @Qualifier.

Ручное связывание (через конфигурацию) - Можно явно указать зависимости в Java- или XML-конфиге.

Java Config:

XML Config (устаревший способ):

Циклические зависимости (Circular Dependencies). Проблема:

Решение:

Использовать конструкторное внедрение + @Lazy.

Пересмотреть архитектуру (возможно, вынести общую логику в третий класс).

Beans

Основная информация

Bean (бин) в контексте Spring Framework — это объект, который управляется IoC-контейнером Spring. Эти объекты создаются, настраиваются и собираются Spring'ом, а затем внедряются в другие компоненты приложения через Dependency Injection (DI).

Основные характеристики бина:

Управляется Spring'ом (не создается через new, а берется из контейнера).

Может иметь зависимости (другие бины, которые внедряются автоматически).

Имеет жизненный цикл (можно задать действия при создании/уничтожении).

Может быть настроен (scope, ленивая инициализация, профили и др.)

Объявление Bean

Есть несколько способов зарегистрировать бин в Spring-контейнере.

Через аннотацию @Component и её производные - Spring автоматически сканирует классы и регистрирует их как бины, если они помечены аннотациями:

@Component – универсальный бин.

@Service – логика бизнес-уровня.

@Repository – доступ к данным (DAO, JPA).

@Controller / @RestController – веб-слой (MVC).

Через Java-конфигурацию (@Bean) - Если класс из сторонней библиотеки (например, RestTemplate), его можно зарегистрировать вручную:

Через XML (устаревший способ)

Жизненный цикл бина

Запускается Spring приложение → Запускается Spring container → Создается объект Бина → В бин внедряются зависимости → Вызывается init-метод → Использование бина → Вызывается destroy-метод → остановка Spring приложения.

Spring управляет созданием и уничтожением бинов. Можно добавить свои действия:

@PostConstruct – после создания

@PreDestroy – перед уничтожением

Scope (область видимости) бина

Определяет, как часто создаётся бин:

Как Spring находит бины?

Сканирование компонентов (@ComponentScan) - Spring ищет классы с аннотациями (@Component, @Service и др.) в указанных пакетах.

Ручное объявление (@Bean) - Если класс нельзя пометить аннотацией (например, он из библиотеки), его регистрируют в @Configuration.

Пример работы с бинами

Объявление бинов

Использование

IoC-контейнер в Spring Framework

IoC (Inversion of Control, Инверсия управления) — это принцип, при котором управление созданием и жизненным циклом объектов передается контейнеру, а не программисту.

IoC-контейнер Spring — это "движок" фреймворка, который:

Создает объекты (бины).

Настраивает их (внедряет зависимости).

Управляет их жизненным циклом.

Как работает IoC-контейнер?

Обнаружение бинов

Сканирует классы с аннотациями (@Component, @Service и др.).

Или регистрирует их вручную через @Bean.

Создание бинов

Контейнер вызывает конструкторы (или фабричные методы).

Заполняет зависимости (через DI).

Управление жизненным циклом

Вызывает методы @PostConstruct / @PreDestroy.

Уничтожает бины при завершении работы приложения.

Предоставление бинов

По запросу (через @Autowired или ApplicationContext.getBean()).

Типы IoC-контейнеров в Spring

BeanFactory (базовый контейнер)

Простейший контейнер, поддерживает только основные функции DI.

Ленивая инициализация (бины создаются при первом запросе).

ApplicationContext (расширенный контейнер)

Наследует BeanFactory, добавляет:

Поддержку AOP, событий, интернационализации.

Автоматическое сканирование компонентов (@ComponentScan).

Загрузку конфигурации из аннотаций (@Configuration).

Пример работы IoC-контейнера

Объявление бинов

2. Использование контейнера

Ключевые преимущества IoC

Уменьшение связанности (Low Coupling)

Классы не зависят от конкретных реализаций, только от интерфейсов.

Гибкость

Легко подменить реализацию (например, MockUserService для тестов).

Упрощение тестирования

Зависимости можно внедрять вручную (без Spring).

Централизованное управление

Вся конфигурация в одном месте (@Configuration или XML).

IoC vs DI

IoC — это общий принцип (инверсия управления).

DI — частный случай IoC (внедрение зависимостей).

Аналог из жизни:

Без IoC: Вы сами идёте в магазин за продуктами.

С IoC: Доставка привозит еду вам домой (контейнер "управляет" процессом).

Как настроить IoC-контейнер?

Через аннотации (современный способ)

Через XML (устаревший метод)

Spring boot

Spring Boot — это фреймворк для быстрого создания standalone-приложений на Spring с минимальной конфигурацией. Он упрощает разработку, предоставляя "из коробки" встроенные серверы (Tomcat, Jetty), автоконфигурацию и стартовые зависимости.

Зачем нужен Spring Boot?

Раньше для Spring-приложений нужно было:

❌ Ручная настройка DispatcherServlet, ViewResolver, DataSource.

❌ Конфигурация в XML или Java-классах.

❌ Сборка WAR-файлов и деплой на внешний сервер (Tomcat).

Spring Boot решает эти проблемы:

✅ Автоконфигурация (автоматически настраивает бины на основе classpath).

✅ Встроенный сервер (не нужен внешний Tomcat).

✅ Starter-зависимости (одна зависимость = весь набор библиотек).

✅ Готовые решения для мониторинга (Actuator), безопасности (Security), БД (JPA).

Основные компоненты Spring Boot

Стартеры (Starters)

Готовые наборы зависимостей для разных задач:

spring-boot-starter-web — для веб-приложений (REST, MVC).

spring-boot-starter-data-jpa — работа с БД через Hibernate.

spring-boot-starter-test — тестирование (JUnit, Mockito).

Пример pom.xml (Maven):

Автоконфигурация

Spring Boot анализирует classpath и автоматически настраивает бины:

Если есть H2 в зависимостях — настраивает in-memory БД.

Если есть spring-webmvc — создаёт DispatcherServlet.

Как отключить автоконфигурацию?

Встроенный сервер - По умолчанию используется Tomcat (порт 8080). Можно изменить в application.properties:

Файлы конфигурации

application.properties / application.yml — настройки приложения.

@ConfigurationProperties — привязка свойств к Java-классам.

Создание приложения на Spring Boot

Инициализация проекта - Через start.spring.io или IDE (IntelliJ IDEA, Eclipse).

Основной класс - Аннотация @SpringBootApplication включает:

@Configuration — класс с бинами.

@ComponentScan — сканирование компонентов.

@EnableAutoConfiguration — автоконфигурацию.

3. REST-контроллер (подробнее ниже)

Жизненный цикл приложения

Запуск:

Инициализация контекста (бины, автоконфигурация).

Запуск встроенного сервера.

Работа:

Обработка HTTP-запросов.

Завершение:

Вызов @PreDestroy методов.

Мониторинг заглушки

Micrometer — это библиотека для сбора метрик из Java-приложений. Она поддерживает экспорт данных в различные системы мониторинга, такие как Prometheus, Grafana, InfluxDB, Datadog и другие.

В Spring Boot Micrometer интегрируется "из коробки" через Spring Boot Actuator, что делает настройку очень простой.

Добавление зависимостей

Для работы Micrometer в Spring Boot нужно добавить две основные зависимости:

Maven (pom.xml)

Gradle (build.gradle)

Если нужно подключить другие системы мониторинга, можно заменить micrometer-registry-prometheus на:

micrometer-registry-influxdb (для InfluxDB)

micrometer-registry-datadog (для Datadog)

micrometer-registry-statsd (для StatsD)

Настройка Actuator и Micrometer - В application.properties или application.yml нужно разрешить доступ к endpoint'ам Actuator:

Проверка работы метрик

После запуска приложения Micrometer автоматически собирает:

JVM-метрики (память, CPU, GC, потоки).

HTTP-запросы (количество, время ответа).

Собственные кастомные метрики (если их добавить).

Доступные endpoints:

http://localhost:8080/actuator/prometheus — метрики в формате Prometheus.

http://localhost:8080/actuator/metrics — список всех метрик.

http://localhost:8080/actuator/metrics/jvm.memory.used — конкретная метрика.

Добавление своих метрик - Micrometer позволяет создавать кастомные метрики. Пример:

Счетчик (Counter)

Таймер (Timer)

Интеграция с Prometheus + Grafana

Запуск Prometheus - Добавьте конфиг prometheus.yml:

Micrometer + Spring Boot Actuator упрощают сбор метрик.

Можно легко экспортировать данные в Prometheus, InfluxDB и другие системы.

Поддержка кастомных метрик (Counter, Timer, Gauge).

Готовая интеграция с Grafana для визуализации.

Сборка и запуск

Создание JAR-файла

Docker-образ

Spring MVC

Spring MVC предполагает разработку web-приложений с использованием архитектуры Model-View-Controller.

MVC - это паттерн проектирования приложений.

Model - логика работы с данными. Это обычный класс, который хранит в себе данные, взаимодействует с БД и отдает эти данные контроллеру.

View - логика представления данных. Класс, который получает данные от контроллера и отображает их (например в браузере).

Controller - логика обработки запросов. Это класс, который обрабатывает запрос, обменивается данными, и дает представление этих данных.

@RestController

@RestController — это удобная аннотация в Spring, которая объединяет функциональность @Controller и @ResponseBody. Она предназначена для создания RESTful веб-сервисов, возвращающих данные в формате JSON/XML вместо HTML-страниц.

Основные особенности @RestController

Автоматическое преобразование возвращаемых значений в JSON/XML (через HttpMessageConverter).

Не требует @ResponseBody на каждом методе (в отличие от @Controller).

Используется для REST API (например, мобильные приложения, SPA, микросервисы).

Сравнение @Controller и @RestController

Пример с @Controller (устаревший подход для API):

То же самое с @RestController:

Как работает @RestController

Запрос приходит на endpoint (например, /user).

Spring вызывает метод контроллера.

Результат метода (объект User) автоматически конвертируется в JSON/XML.

Клиент получает ответ в формате: {"name": "John"}

Пример REST-контроллера

Mappings / Request Handlers

Ключевые аннотации для методов

Параметры методов

@PathVariable — извлечение значения из URL:

@RequestParam — параметры запроса (/users?name=John):

@RequestBody — получение JSON/XML из тела запроса:

@RequestHeader — доступ к HTTP-заголовкам:

Подключение к Kafka и отправка сообщений в Spring Boot

Apache Kafka — это распределённый потоковый брокер сообщений, используемый для обработки событий в реальном времени.

В Spring Boot работа с Kafka упрощается благодаря Spring for Apache Kafka.

Добавление зависимостей - В pom.xml (Maven) или build.gradle (Gradle) добавьте:

Настройка подключения к Kafka - В application.yml (или application.properties) укажите адрес Kafka-брокера и тему:

Отправка сообщений (Producer) - Создайте компонент для отправки сообщений в Kafka:

4. Получение сообщений (Consumer) - Для чтения сообщений из Kafka создайте @KafkaListener:

Lombok

Lombok — это библиотека Java, которая упрощает написание шаблонного кода (геттеры, сеттеры, конструкторы и т. д.) через аннотации.

При работе с DTO Lombok особенно полезен, так как позволяет сократить объем кода, сохраняя его читаемым и удобным для поддержки.

DTO

DTO (Data Transfer Object) — это шаблон проектирования, используемый для передачи данных между подсистемами приложения (например, между клиентом и сервером или между слоями серверного приложения).

Основные особенности DTO:

Только данные

DTO не содержит бизнес-логики, только поля для хранения и передачи данных.

Пример: JSON-объект, передаваемый между фронтендом и бэкендом.

Оптимизация передачи

Позволяет объединять несколько данных в один объект, уменьшая количество запросов.

Например, вместо 10 отдельных полей — один DTO с нужной информацией.

Изоляция слоев

Защищает внутреннюю структуру данных приложения.

Клиент работает с DTO, а не с сущностями базы данных (Entity) или domain-моделями.

Гибкость

Можно кастомизировать данные под конкретные запросы (например, исключать лишние поля или добавлять вычисляемые).

Пример DTO

Допустим, у нас есть User в базе данных, но клиенту не нужны все поля:

Entity (база данных)

DTO (для передачи клиенту)

Где применяется DTO?

REST API – для формата запросов/ответов.

Микросервисы – обмен данными между сервисами.

GraphQL – клиент запрашивает только нужные поля.

Базы данных – чтобы не передавать напрямую Entity.

Основные аннотации Lombok для DTO

@Data

Автоматически генерирует:

Геттеры для всех полей.

Сеттеры для всех не-final полей.

toString(), equals(), hashCode().

Конструктор со всеми обязательными полями (final или @NonNull).

Пример DTO с @Data:

Эквивалентно написанию всех геттеров, сеттеров, toString(), equals() и hashCode() вручную.

@Getter / @Setter

Если не нужны все методы из @Data, можно использовать точечно:

@NoArgsConstructor, @AllArgsConstructor, @RequiredArgsConstructor

@NoArgsConstructor – создает пустой конструктор.

@AllArgsConstructor – конструктор со всеми полями.

@RequiredArgsConstructor – конструктор только для final или @NonNull полей.

@Builder

Позволяет использовать паттерн Builder для удобного создания DTO:

Использование:

WireMock

Настроить WireMock достаточно просто. Это инструмент для мокирования HTTP-серверов, который полезен при тестировании API. Вот пошаговая инструкция:

Установка WireMock

Есть два основных способа:

Запуск как standalone-сервер (JAR)

Скачайте JAR-файл с официального сайта или через Maven:

Запустите сервер:

--port — порт (по умолчанию 8080).

--https-port — для HTTPS.

--verbose — логирование запросов.

Встроенный в Java-приложение (Dependency)

Если используете Maven/Gradle, добавьте зависимость:

Или для Gradle:

Настройка моков (stub-правила)

WireMock позволяет задавать ожидаемые запросы и ответы. Примеры:

Через REST API (после запуска сервера)

Проверка:

Через Java-код

Дополнительные настройки

Динамические ответы

Используйте response-templates для генерации ответов:

Запрос:

Запись запросов (Record & Playback)

Запустите WireMock в режиме прокси:

Все запросы к http://localhost:8080 будут записаны в папку mappings.

Сохранение моков в файлы

WireMock сохраняет моки в папках:

mappings/ — правила (JSON).

__files/ — тела ответов (если используются внешние файлы).

Пример mappings/test.json:

Файл __files/data.json:

Расширенные возможности

Задержки (delay):

Проверка заголовков:

Остановка сервера

Для standalone-режима: Ctrl + C.

В Java-коде: WireMock.shutdownServer().

✍️ БД SQL (До июля)

ООП

Объектно-ориентированное программирование (ООП) включает 3 ключевых принципа : инкапсуляция, наследование, полиморфизм ,некоторые выделяют 4 тип: абстракция.

1)Инкапсуляция — сокрытие внутренних деталей реализации и доступ к ним через методы:🡫

2)Наследование — создание новых классов на основе существующих:🡫

3)Полиморфизм — способность объектов разных классов обрабатываться единым образом через родительский тип.

В полиморфизме родительский класс можно рассматривать как каркас или шаблон, в котором определены методы, ожидаемые от всех его наследников:🡫

реализованы в родительском классе (базовая реализация)

переопределены (override) в подклассах с их собственной логикой

интерфейсы — описывают поведение, которое может быть реализовано в любом классе:

Это позволяет работать с объектами через интерфейс (или тип) родительского класса, но фактически выполнять действия, определенные в подклассе

4)Абстракция Выделение общих черт и скрытие деталей реализации. Используются абстрактные классы и интерфейсы.

Абстрактный класс— это недополненный шаблон, нельзя создать напрямую только наследовать,содержит в себе какую-то общую логику и уникальное поведение ,которое реализует наследниК:

Разница?

Методы

Коллекции

Память

Курс Java на JavaRush - список всех лекций

Устройство памяти

Java Stack – это место где хранится список вызовов методов, кто кого вызывал. Стек может переполниться и тогда будет ошибка “stack overflow”.

Локальная переменная также может быть ссылкой на объект. В этом случае ссылка (локальная переменная) хранится в стеке потоков, но сам объект хранится в куче.

Initial heap – куча, хранение всех данных приложения, вне хипа хранится код программы и вспомогательные данные. Хип место, где находятся объекты. Здесь обитает сборщик мусора. Пространство хипа делится на три зоны генерации.

Old generation - здесь хранятся объекты, которые пережили несколько проходов сборщика.

Переменные-члены объекта хранятся в куче вместе с самим объектом. Это верно как в случае, когда переменная-член имеет примитивный тип, так и в том случае, если она является ссылкой на объект.

Статические переменные класса также хранятся в куче вместе с определением класса.

Куча/Heap содержит все объекты, созданные в вашем приложении, независимо от того, какой поток создал объект. К этому относятся и обертки примитивных типов (например, Byte, Integer, Long и так далее). Неважно, был ли объект создан и присвоен локальной переменной или создан как переменная-член другого объекта, он хранится в куче.

Garbage Collection (сборка мусора)

Сборка мусора в Java — это процесс, с помощью которого программы Java автоматически управляют памятью.

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

Живые — эти объекты используются, на них ссылаются из потока/кода программы еще.

Мертвые — эти объекты больше нигде не используются, ссылок на них нет.

Garbage Collector (Сборщик мусора) находит эти неиспользуемые объекты и удаляет их, чтобы освободить память.

Для того, чтобы признать объект живым — наличия ссылок недостаточно. Все потому, что одни мертвые объекты могут ссылаться на другие мертвые объекты. Именно поэтому нужно, чтобы среди всех ссылок на объект, была хотя бы одна от “живого” объекта.

Сборщики мусора работают с концепцией GC Roots (корней сбора мусора) для того, чтоб различать живые и мертвые объекты. Есть 100% живые объекты и от них идут ссылки, которые оживляют другие объекты и так далее.

Примеры таких корней:

Классы, которые загружаются системным загрузчиком классов.

Живые потоки.

Параметры методов, которые выполняются в данный момент и локальные переменные.

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

Объекты, которые удерживаются из сборки мусора для некоторых целей.

Сборщик мусора просматривает весь граф объектов в памяти, начиная с этих корней и следуя ссылкам на другие объекты.

Этапы сборки мусора в Java

Стандартная реализация сборки мусора имеет три этапа.

1. Помечаем объекты как живые

На этом этапе сборщик мусора (GC) должен идентифицировать все живые объекты в памяти путем обхода графа объектов.

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

2. Зачистка мертвых объектов

После фазы разметки пространство памяти занято либо живыми (посещенными), либо мертвыми (не посещаемыми) объектами. Фаза зачистки освобождает фрагменты памяти, которые содержат эти мертвые объекты.

3. Компактное расположение оставшихся объектов в памяти

Не обязательно, чтобы мертвые объекты, которые были удалены на предыдущей фазе, находились рядом друг с другом. Таким образом, ты рискуешь получить фрагментированное (полупустое) пространство памяти.

Но, конечно же, предусмотрев это, есть возможность уплотнить память в момент, когда сборщик мусора удаляет мертвые объекты. Оставшиеся будут располагаться в непрерывном блоке в начале кучи.

Процесс уплотнения облегчает последовательное выделение памяти для новых объектов.

Java-сборщики мусора реализуют некоторую стратегию сбора мусора поколений, которая умеет классифицировать объекты по возрасту.

Такую необходимость (отмечать и уплотнять все объекты) в JVM можно назвать неэффективной. Так как по мере выделения большого количества объектов их список растет, что приводит к увеличению времени сбора мусора. Эмпирический анализ приложений показал, что большинство объектов в Java недолговечны.

Область памяти кучи в JVM разделена на три секции:

Молодое поколение

Вновь созданные объекты начинаются в молодом поколении. Молодое поколение далее подразделяется на две категории.

Eden space (Пространство Эдема) — все новые объекты начинают здесь, им выделяется начальная память.

Пространства выживших (FromSpace и ToSpace) — объекты перемещаются сюда из Эдема после того, как пережили один цикл сборки мусора.

Процесс, когда объекты собираются в мусор из молодого поколения, называется малым событием сборки мусора.

Когда пространство Эдема заполнено объектами, выполняется малая сборка мусора. Все мертвые объекты удаляются, а все живые — перемещаются в одно из оставшихся двух пространств. Малая GC также проверяет объекты в пространстве выживших и перемещает их в другое (следующее) пространство выживших.

Возьмем в качестве примера следующую последовательность.

В Эдеме есть объекты обоих типов (живые и мертвые).

Происходит малая GC — все мертвые объекты удаляются из Эдема. Все живые объекты перемещаются в пространство-1 (FromSpace). Эдем и пространство-2 теперь пусты.

Новые объекты создаются и добавляются в Эдем. Некоторые объекты в Эдеме и пространстве-1 становятся мертвыми.

Происходит малая GC — все мертвые объекты удаляются из Эдема и пространства-1. Все живые объекты перемещаются в пространство-2 (ToSpace). Эдем и пространство-1 пусты.

Таким образом в любое время одно из пространств для выживших всегда пусто. Когда выжившие объекты достигают определенного порога перемещения по пространствам выживших, они переходят в старшее поколение.

Для установки размера молодого поколения можно воспользоваться флагом -Xmn.

Старшее поколение

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

Пороговое значение срока службы объекта определяет, сколько циклов сборки мусора он должен пережить, прежде чем будет перемещен в старшее поколение. Процесс, когда объекты отправляются в мусор из старшего поколения, называется основным событием сборки мусора.

Для установки начального и максимального размера памяти кучи можно воспользоваться флагами -Xms и -Xmx.

Так как Java задействует сборку мусора по поколениям, то, чем больше событий сборки мусора переживает объект, тем дальше он продвигается в куче. Он начинает в молодом поколении и в конечном итоге заканчивает в штатном поколении, если проживет достаточно долго.

Чтобы понять продвижение объектов между пространствами и поколениями, рассмотрим следующий пример:

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

Как только произойдет малая сборка мусора, живые объекты из Эдема перемещаются в пространство FromSpace. Когда происходит следующая малая сборка мусора, живые объекты как из Эдема, так и из пространства перемещаются в пространство ToSpace.

Этот цикл продолжается определенное количество раз. Если объект все еще “в строю” после этого момента, следующий цикл сборки мусора переместит его в пространство старшего поколения.

Постоянное поколение и мета-пространство

Метаданные, такие как классы и методы, хранятся в постоянном поколении. JVM заполняет его во время выполнения на основе классов, используемых приложением. Классы, которые больше не используются, могут переходить из постоянного поколения в мусор.

Для установки начального и максимального размера постоянного поколения вы можете воспользоваться флагами -XX:PermGen и -XX:MaxPermGen.

Мета-пространство

Начиная с Java 8, на смену пространству постоянного поколения (PermGen) приходит пространство памяти MetaSpace. Реализация отличается от PermGen — это пространство кучи теперь изменяется автоматически.

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

Виды сборщиков мусора в Java

Serial GC

Серийный GC — это самая простая реализация GC. Она предназначена для небольших приложений, работающих в однопоточных средах. Все события сборки мусора выполняются последовательно в одном потоке. Уплотнение выполняется после каждой сборки мусора.

Запуск сборщика приводит к событию “остановки мира”, когда все приложение приостанавливает работу.

Parallel GC

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

Несколько потоков предназначаются для малой сборки мусора в молодом поколении. Единственный поток занят основной сборкой мусора в старшем поколении. Запуск параллельного GC также вызывает “остановку мира” и приложение зависает.

CMS GC

Также известен нам как параллельный сборщик низких пауз.

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

Из-за этого сборщик CMS потребляет больше ресурсов процессора, чем другие сборщики.

G1 (Garbage first) GC

G1GC был задуман как замена CMS и разрабатывался для многопоточных приложений, которые характеризуются крупным размером кучи (более 4 ГБ). Он параллелен и конкурентен, как CMS, но “под капотом” работает совершенно иначе, чем старые сборщики мусора.

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

G1 разбивает кучу на набор областей одинакового размера (в зависимости от размера кучи) и сканирует их в несколько потоков. Область во время выполнения программы может неоднократно становиться как старой, так и молодой.

После завершения этапа разметки G1 знает, в каких областях содержится больше всего мусора. Если пользователь заинтересован в минимизации пауз, G1 может выбрать только несколько областей. Если время паузы неважно для пользователя или предел этого времени установлен высокий, G1 пройдет по большему числу областей.

Поскольку G1 GC идентифицирует регионы с наибольшим количеством мусора и сначала выполняет сбор мусора в них, он и называется: “Мусор — первым”.

Помимо областей Эдема, Выживших и Старой памяти, в G1GC присутствуют еще два типа.

Humongous (Огромная) — для объектов большого размера (более 50% размера кучи).

Available (Доступная) — неиспользуемое или не выделенное пространство.

Allocation Failure - это сигнал, что нельзя выделить память без полной очистки хипа.

Evacuation Pause — это пауза в работе приложения, во время которой Garbage Collector (GC) перемещает (эвакуирует) живые объекты из одних регионов памяти в другие. "Stop The World"

Humongous объекты:

Не помещаются в Eden (новое поколение);

Сразу размещаются в Old Gen (старое поколение);

Не любят перемещения, что усложняет работу GC;

Могут вызывать Full GC (очень дорогую по времени операцию);

Занимают много подряд идущих регионов, и если таких объектов много — они могут фрагментировать память.

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

GCLocker Initiated GC - специальная GC-пауза :"GC был отложен, пока не завершится нативный код. Как только он завершился — мы запустили GC немедленно."

Shenandoah (Шандара)

Shenandoah — новый GC, выпущенный как часть JDK 12. Ключевое преимущество Shenandoah перед G1 состоит в том, что большая часть цикла сборки мусора выполняется одновременно с потоками приложений. G1 может эвакуировать области кучи только тогда, когда приложение приостановлено, а Shenandoah перемещает объекты одновременно с приложением.

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

ZGC

ZGC — еще один GC, выпущенный как часть JDK 11 и улучшенный в JDK 12.

Он предназначен для приложений, которые требуют быстродействия и низкой задержки (паузы в менее чем 10 мс) и/или задействуют очень большую кучу (несколько терабайт).

Основные цели ZGC — низкая задержка, масштабируемость и простота в применении. Для этого он позволяет Java-приложению продолжать работу, несмотря на то, что выполняются операции по сбору мусора. По умолчанию ZGC освобождает неиспользуемую память и возвращает ее в операционную систему.

Таким образом, ZGC привносит значительное улучшение по сравнению с другими традиционными GC, обеспечивая чрезвычайно низкое время паузы (обычно в пределах 2 мс).

https://habr.com/ru/articles/269621/ - цикл статей о работе каждого GC

Ссылки со всей теории

Кафка

Курс по кафке

Топик (Topic)

Сообщения в Kafka отправляются в топики (темы). Их можно представить как таблицы в базе данных. Каждая таблица в БД зачастую хранит один вид данных - например, данные по заказам, данные по пользователям и так далее. С топиками всё работает аналогично. По сути, топик логически группирует сообщения. У каждого топика есть уникальное название.

Партиция (Partition)

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

Партиции нумеруются с нуля. Количество партиций в топике задаётся при создании топика, после этого их число можно увеличить (но уменьшить число партиций уже нельзя ⚠️). Для каждого топика число партиций устанавливается отдельно. Сообщения записываются в партицию путём добавления в конец, а читаются по порядку от начала к концу.

(!) Важно: сообщения в партиции хранятся в строгом порядке их поступления. То есть партиция - это фактически очередь сообщений. При этом порядок на уровне топика НЕ гарантируется!

Благодаря партициям, в Kafka обеспечивается возможность горизонтального масштабирования. Партиции можно реплицировать - то есть размещать копии одной и той же партиции на разных брокерах. Таким образом, если какой-то брокер в кластере выйдет из строя, мы не потеряем данные, которые были в партициях, размещённых на нём (поскольку у нас будут копии этих партиций на других брокерах). Подробнее про репликацию мы поговорим в уроке "2.6 Репликация данных".

Давайте рассмотрим следующую иллюстрацию. Здесь изображён один брокер, в котором находятся два топика - order.status и user.actions. Некий сервис-продьюсер записывает данные в топик order.status, в свою очередь два сервиса-консьюмера читают данные из этого топика.

Офсет (Offsets)

Офсет (или смещение) — это порядковый номер сообщения в определённой партиции. Он позволяет консьюмерам отслеживать, какие сообщения из данной партиции ими уже были прочитаны, а какие - ещё нет. Дело в том, что Kafka не удаляет сообщения после их чтения (в отличие от RabbitMQ); вместо этого консьюмеры отслеживают последний офсет, который они обработали. Они могут перечитывать данные или начать чтение с определённого места (офсета) в партиции.

(!) Важно: офсет считается с нуля.

Пример

Предположим, что у нас есть топик "order.status", в который отправляются статусы заказов с сайта интернет-магазина. У нас имеется консьюмер, в качестве которого выступает сервис "stats-saver", который читает топик со статусами заказов и сохраняет статистику по этим заказам в базу данных.

После прочтения каждого сообщения сервис "stats-saver" будет сохранять его offset (порядковый номер внутри партиции топика) на брокере. Благодаря этому, если с сервисом что-то случится (например, на нём возникнет ошибка и он перезагрузится), он сможет обратиться к брокеру и узнать на каком offset он закончил чтение партиции. После этого сервис сможет начать чтение с того сообщения, которое сервис ещё не читал (т.е. с offset + 1).

Сообщение (message)

Используемая в Kafka единица данных называется сообщением (message). Если вы ранее работали с базами данных, то можете рассматривать сообщение как аналог строки (row) или записи (record) в таблице БД. С точки зрения Kafka сообщение представляет собой набор байтов, так что для неё содержащиеся в нём данные не имеют формата и какого-либо смысла.

Пример сообщения

В данном сообщении передаётся информация о стоимости акции компании "Сбербанк" на определённую дату и время:

• name - название акции

• price - стоимость акции

• time - дата и время, когда акция имела такую стоимость

{

name: "SBER",

price: 285.5,

time: "17-05-2023T15:00:00"

}

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

Пример ключа и значения:

Давайте выделим название акции в ключ сообщения. Остальные данные оставим в значении.

Ключ (key)

{

name: "SBER"

}

Значение (value)

{

price: 285.5,

time: "17-05-2023T15:00:00"

}

Помимо ключа и значения сообщения в Kafka содержат Timestamp (время создания сообщения) и Headers (заголовки - опциональный набор метаданных). Про заголовки мы ещё поговорим в одном из следующих уроков.

Форматы сообщений

В примерах выше, сообщения представлены в популярном формате JSON (JavaScript Object Notion). Помимо JSON, для передачи данных посредством Kafka, также часто используются форматы AVRO (Apache AVRO) и Protobuf (Protocol Buffers). В рамках одного топика рекомендуется использовать только один формат данных и одну структуру сообщений.

Рассмотрим отличия данных форматов:

Критерий /Формат AVRO Protobuf JSON

Тип данных бинарный бинарный текстовый

Поддержка схемы да да да

Чтение схемы да нет нет

Язык схемы* JSON Protobuf IDL JSON

Сжатие данных Высокое сжатие Высокое сжатие Без сжатия

Скорость передачи Высокая Высокая Умеренная

Поддержка языками

программирования Средняя Хорошая Отличная

Из всех указанных форматов, вероятно, вы уже хорошо знакомы с JSON. Его часто используют в REST API. Он удобен в первую очередь благодаря тому, что является человекочитаемым (текстовым). Мы можем легко взглянуть на сообщение и понять что за данные в нём содержатся. AVRO и Protobuf в свою очередь являются бинарными. То есть сообщения в этих форматах выглядят как набор байтов. Но в этом также заключается их главный плюс, сообщения в этих форматах занимают меньше места, их быстрее читать и записывать, а также быстрее передавать по сети между разными сервисами.

Не пугайтесь, если не знаете что такое схема. Это понятие мы рассмотрим на следующем шаге урока.

Схема (Scheme)

Схема позволяет описать структуру данных: какие поля в ней присутствуют, какие у них типы и т.д.

Почему это удобно? Например, мы можем задать в коде сервиса-консьюмера как будут обрабатываться те сообщения, которые не прошли проверки с помощью схемы. Например, мы можем пропускать подобные сообщения. Сценарий обработки может зависеть от конкретных требований к сервису.

Такое определение может показаться довольно сложным. Давайте рассмотрим конкретные примеры.

JSON

Стандартом описания схемы для данных в формате JSON является JSON Schema.

Предположим, у нас есть топик user.data, в котором пересылаются следующие данные о пользователях социальной сети в формате JSON:

• id - уникальный идентификатор пользователя

• first_name - имя пользователя

• last_name - фамилия пользователя

{

"id": 210700286,

"first_name": "Artem",

"last_name": "Sidorov"

}

JSON Schema, описывающая данный объект данных может выглядеть так:

{

"type": "object",

"properties": {

"id": {

"type": "integer",

"description": "User ID"

},

"first_name": {

"type": "string",

"description": "User first name"

},

"last_name": {

"type": "string",

"description": "User last name"

}

},

"required": [

"id",

"first_name",

"last_name"

],

"additionalProperties": false

}

Внутри элемента properties для каждого поля данных задаются типы данных, которым должны соответствовать поля сообщения. Например, поле id должно быть типа integer (целое число), а поле first_name - типа string (строка).

Ключевое слово required задает перечень обязательных полей. Если хотя бы одно из перечисленных полей будет отсутствовать, сообщение не пройдет валидацию (проверку) по такой схеме.

Ключевое слово additionalProperties задает возможность наличия дополнительных полей у объекта. В нашем случае дополнительные поля запрещены (при их наличии объект не пройдет валидацию по схеме).

Подробнее про JSON-схемы можно узнать в официальной документации (на английском языке) - https://json-schema.org/learn/getting-started-step-by-step#define.

Protobuf

Формат Protobuf (Protocols Buffer) был придуман в Google. Данный формат использует собственный язык описания структуры данных - Protobuf IDL (Interface Description Language).

Предположим, что в топике search.queries передаются поисковые запросы пользователей к Yandex, в том числе:

• query - сам поисковый запрос

• page_number - номер поисковой страницы

• results_per_page - кол-во результатов поиска на каждой странице

Пример схемы на языке IDL:

syntax = "proto3";

message SearchRequest {

string query = 1;

int32 page_number = 2;

int32 results_per_page = 3;

}

Первая строка схемы указывает, что далее последует описание схемы на языке IDL 3-й версии.

Далее внутри message описывается структура самого сообщения. SearchRequest - это название сообщения.

Внутри указываются - название поля, его тип, а также номер. Например, поле query имеет строковый тип (string) и номер "1". Номера используются для идентификации полей в сообщении, поэтому они должны быть уникальны.

На основе схемы с помощью специального компилятора (protoc) можно сгенерировать готовые классы сообщений, методы сериализации/десереализации, а также другие вспомогательные функции для вашего языка программирования (protoc поддерживает множество языков, включая Java, C++, Go и многие другие).

Создание топика

Начнём освоение Kafka CLI с создания топика. Для этого воспользуемся утилитой kafka-topics.sh. Эта утилита также используется для удаления, изменения настроек и получения информации о топиках.

Давайте создадим топик orders.status. Для этого мы указываем следующие опции:

• create - указывает, что мы собираемся создать топик

• topic - после неё указывается название топика

• bootstrap-server - после неё указывается адрес брокера (-ов) Kafka, к которым мы будем подключаться

bin/kafka-topics.sh --create --topic orders.status --bootstrap-server localhost:9092

Поскольку мы подключаемся к брокеру Kafka, который запущен на том же самом компьютере, на котором мы создаём консьюмера, то мы указываем в качестве хоста localhost, а в качестве порта 9092 (стандартный порт, на котором по-умолчанию запускается брокер Kafka). В нашем случае у нас всего один брокер. Если у нас имеется несколько брокеров (кластер), то хосты и порты брокеров указываются через запятую.

Если всё прошло успешно, то мы должны увидеть следующее сообщение в Терминале:

Created topic orders.status.

Оно говорит нам о том, что топик orders.status был успешно создан. При этом топик будет состоять из одной партиции.

Не пугайтесь, если увидите следующее предупреждение (WARNING):

WARNING: Due to limitations in metric names, topics with a period ('.') or underscore ('_') could collide. To avoid issues it is best to use either, but not both.

Оно говорит о том, что в некоторых случаях могут быть проблемы, если вы используете в наименованиях своих топиков и точку, и нижнее подчёркивание и рекомендует использовать только один из этих символов для наименования всех своих топиков. Мы, в свою очередь, будем использовать точку.

Вывод списка топиков

Список существующих топиков на брокере (-ах) можно вывести с помощью опции list:

bin/kafka-topics.sh --list --bootstrap-server localhost:9092

Пока что у нас должен быть только один топик, который мы создали ранее - orders.status:

Указание количества партиций

При создании топика мы можем указать количество партиций с помощью опции partitions. Давайте создадим топик с двумя партициями.

bin/kafka-topics.sh --create --topic two.partitions.topic --bootstrap-server localhost:9092 --partitions 2

Если не указать опцию partitions, то будет создан топик, состоящий только из одной партиции.

Указание коэффициента репликации

При создании топика мы также можем указать коэффициент репликации. За это отвечает опция replication-factor.

Например:

bin/kafka-topics.sh --create --topic replication.factor.topic --bootstrap-server localhost:9092 --replication-factor 3

В нашем случае мы не сможем создать топик с таким коэффициентом репликации, поскольку у нас имеется всего один брокер. Если мы всё же попытаемся выполнить данную команду, то получим следующую ошибку:

Error while executing topic command : Replication factor: 3 larger than available brokers: 1.

[2024-04-15 17:00:00,000] ERROR org.apache.kafka.common.errors.InvalidReplicationFactorException: Replication factor: 3 larger than available brokers: 1.

(org.apache.kafka.tools.TopicCommand)

⚠️ Запомните: коэффициент репликации не может превышать количество брокеров в кластере!

Таким образом, для обеспечения replication-factor = 3 нам потребуется иметь 3 брокера в кластере.

Указание внутренних настроек топика

Как мы уже обсуждали ранее, у топиков, брокеров, продьюсеров и консьюмеров есть свои настройки. Чтобы задать значения для настроек при создании топика посредством Kafka CLI, нужно воспользоваться опцией config.

Давайте ограничим максимальный размер сообщения (в байтах), которое можно будет записать в топик. Это контролируется настройкой топика max.message.bytes. Давайте зададим значение в 4000 байт (это почти 4 КБ):

bin/kafka-topics.sh --create --topic one.config.topic --bootstrap-server localhost:9092 --config max.message.bytes=4000

Если требуется указать сразу несколько настроек для топика, то каждую из них нужно прописывать в виде отдельного config, то есть в следующем виде:

bin/kafka-topics.sh --create --topic multi.config.topic --bootstrap-server localhost:9092 --config max.message.bytes=4000 --config cleanup.policy=compact

Получение описания топика

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

Давайте выведем описание для топика two.partitions.topic, который мы ранее создали с 2-мя партициями.

bin/kafka-topics.sh --describe --topic two.partitions.topic --bootstrap-server localhost:9092

Результат:

Topic: two.partitions.topic TopicId: 9ZBGUW4aRnKHdN7qAL6a5Q PartitionCount: 2 ReplicationFactor: 1 Configs:

Topic: two.partitions.topic Partition: 0 Leader: 0 Replicas: 0 Isr: 0

Topic: two.partitions.topic Partition: 1 Leader: 0 Replicas: 0 Isr: 0

Что мы видим в выводе?

• Название топика (Topic)

• Уникальный идентификатор топика (TopicId)

• Количество партиций (PartitionCount)

• Коэффициент репликации (ReplicationFactor)

• Настройки топика (Configs)

В разрезе партиций выводится:

• Номер партиции (Partition)

• Номер брокера, на котором находится ведущая реплика партиции (Leader)

• Номера брокеров, на которых находятся реплики данной партиции (Replicas)

• Номера брокеров, на которых находятся синхронизированные (in-sync) реплики партиции (ReplicationFactor)

Изменение количества партиций в топике

Если мы хотим изменить количество патриций в топике, то для этого понадобится опция alter.

Давайте сделаем это для топика orders.status. Напомню, что при его создании мы не указывали количество партиций явно, следовательно сейчас топик состоит из одной партиции. Увеличим количество партиций до 5-ти:

bin/kafka-topics.sh --alter --topic orders.status --partitions 5 --bootstrap-server localhost:9092

Проверим, что их число действительно поменялось с помощью опции describe:

bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders.status

Вывод:

Topic: orders.status TopicId: 9p2uaygsQhy-oIqp-SLsPg PartitionCount: 5 ReplicationFactor: 1 Configs:

Topic: orders.status Partition: 0 Leader: 0 Replicas: 0 Isr: 0

Topic: orders.status Partition: 1 Leader: 0 Replicas: 0 Isr: 0

Topic: orders.status Partition: 2 Leader: 0 Replicas: 0 Isr: 0

Topic: orders.status Partition: 3 Leader: 0 Replicas: 0 Isr: 0

Topic: orders.status Partition: 4 Leader: 0 Replicas: 0 Isr: 0

⚠️ Важно: не забывайте, что после создания топика, число его партиций можно увеличить, но нельзя уменьшить!

Удаление топика

Для удаления топика используется опция delete:

bin/kafka-topics.sh --delete --topic orders.status --bootstrap-server localhost:9092

Можно удалить сразу несколько топиков, перечислив их через запятую:

bin/kafka-topics.sh --delete --topic one.config.topic,multi.config.topic --bootstrap-server localhost:9092

Мы можем убедиться, что топики были удалены с помощью опции list:

bin/kafka-topics.sh --list --bootstrap-server localhost:9092

Работа с продьюсерами происходит посредством утилиты kafka-console-producer.sh.

Отправка сообщения

Давайте создадим продьюсера, который будет отправлять сообщения в топик orders.status.

Поскольку ранее мы удалили данный топик, давайте сперва пересоздадим его заново (например, с 4-мя партициями):

bin/kafka-topics.sh --create --topic orders.status --partitions 4 --bootstrap-server localhost:9092

Теперь создадим продьюсера для данного топика:

bin/kafka-console-producer.sh --topic orders.status --bootstrap-server localhost:9092

После этого вы должны увидеть в консоли символ ">". Все строчки текста, которые вы напишите после него, будут отправлены в топик orders.status. Конец строки определяется нажатием кнопки "Enter".

bin/kafka-console-producer.sh --topic orders.status --bootstrap-server localhost:9092

>First message

>Learning Kafka

>Love learning

>

⚠️ Важно:

1) Обратите внимание, что по-умолчанию, сообщения отправляются с пустым ключом.

2) Если вы создадите продьюсер для несуществующего топика, то он будет мгновенно создан. При этом его количество партиций и коэффициент репликаций будут соответствовать значениям по-умолчанию, которые установлены на уровне брокера (про эти настройки мы поговорим в следующем модуле).

Отправка сообщения с ключом

Для того, чтобы отправить сообщение с ключом, нам следует задать для продьюсера ряд параметров с помощью опций property:

• parse.key (нужно ли обрабатывать ключ - да/нет)

• key.separator (символ, разделяющий ключ и значение сообщения)

Откройте новую сессию в Терминале и выполните следующую команду:

bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic topic.with.key --property parse.key=true --property key.separator=:

Отправим несколько сообщений с ключом (в качестве разделителя ключа и значения мы указали символ ":").

>first key:value1

>first:last

>

Работа с консьюмерами происходит посредством утилиты kafka-console-consumer.sh.

Получение актуальных сообщений

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

bin/kafka-console-consumer.sh --topic orders.status --bootstrap-server localhost:9092

После выполнения команды вы должны увидеть символ ">" в консоли. Однако, вы не увидите сообщений, которые мы ранее отправили в топик. Если мы создаём консьюмер вышеуказанным образом, то мы увидим только те сообщения, которые были отправлены в топик после создания консьюмера!

Получение исторических сообщений

Если мы хотим сперва прочитать исторические сообщения, а затем начать читать актуальные, следует воспользоваться опцией from-beginning.

bin/kafka-console-consumer.sh --topic orders.status --from-beginning --bootstrap-server localhost:9092

bin/kafka-console-consumer.sh --topic orders.status --from-beginning --bootstrap-server localhost:9092

First message

Learning Kafka

Love learning

Урок 4.4

Остановить консьюмер мы можем с помощью сочетания клавиш Ctrl + C. При этом будет выведено сообщение с информацией о количестве обработанных сообщений, например:

^CProcessed a total of 4 messages

⚠️ Важно:

1) По умолчанию, консьюмер не выводит ключ сообщения.

2) Обратите внимание, что мы не можем удалить топик, если его читает консьюмер (-ы).

Вывод ключа и метки времени у сообщения

Для того, чтобы помимо значения вывести ключ сообщения, нам требуется воспользоваться опцией property и указать параметр print.key, равный true (при этом также следует прописать print.value=true, чтобы вывелось и значение сообщения):

bin/kafka-console-consumer.sh --topic orders.status --property print.key=true --property print.value=true --bootstrap-server localhost:9092

Для того, чтобы вывести ещё и метку времени, требуется указать прописать print.timestamp=true:

bin/kafka-console-consumer.sh --topic orders.status --property print.key=true --property print.value=true --property print.timestamp=true --bootstrap-server localhost:9092 --from-beginning

Вывод:

CreateTime:1713977447000 null First message

CreateTime:1713977507000 null Learning Kafka

CreateTime:1713977747000 null Love learning

CreateTime:1713979847000 null Урок 4.4

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

Обратите внимание на формат времени, в котором записано CreateTime. Он может показаться несколько странным, однако это достаточно популярный формат записи даты и времени - Unix Timestamp. Он показывает количество миллисекунд (или иногда секунд) с начала эпохи Unix (то есть с 1 января 1970 года).

Батчи и батчевание

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

Давайте познакомиться с понятием "батч" и "батчевание". Батч (batch) - это группа сообщений, объединённых вместе для отправки как единое целое. Зачастую продьюсер не отправляет каждое сообщение отдельно, а сначала ждёт, пока сформируется группа из нескольких сообщений (message set), и только потом отправляет их как единое целое. Для чего это нужно? Такой подход позволяет сделать передачу данных более эффективной за счёт снижения накладных расходов на сетевые запросы.

Можно провести аналогию с автобусами. Представьте, что автобус ждёт, пока не соберётся группа пассажиров, прежде чем отправиться в путь. Это похоже на батчевание в Kafka: один автобус перевозит много пассажиров за один раз, снижая накладные расходы (на топливо и пр.) и делая поездки более эффективными. Когда автобус приезжает на конечную остановку (партицию топика), все пассажиры (сообщения) выгружаются разом.

Сверху - сообщения отправляются по отдельности (1 сообщение в одном батче),

снизу - сообщения отправляются в одном батче (4 сообщения в одном батче)

Продьюсер отправляет батч, когда его размер достигает заданного размера batch.size или по прошествии времени, установленного в linger.ms . То есть, если хоть одно из условий выполнится, то батч будет немедленно отправлен продьюсером. Обе опции являются главными настройками продьюсера, контролирующими батчевание.

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

Представим, что у нас есть топик с двумя партициями и два брокера. Каждая партиция топика расположена на одном из брокеров (партиция 0 на брокере 0, а партиция 1 на брокере 1). Сообщения, предназначенные для партиции 0, можно представить как пассажиров, которые хотят добраться до города "А", а те, что предназначены для партиции 1 - до города "В". Пассажиры (сообщения), направляющиеся в разные города (партиции), не могут ехать вместе в одном автобусе (батче), так как каждый автобус (батч) следует в конкретный город (определённую партицию конкретного брокера).

По умолчанию, batch.size равен 16 384 байтам, а linger.ms - 0 миллисекунд. На первый взгляд, кажется, что при linger.ms = 0 все сообщения будут отправляться моментально, а значит - по отдельности (1 сообщение в одном батче). Однако зачастую продьюсер собирает несколько сообщений практически одновременно, поэтому даже при linger.ms = 0 они могут быть объединены в батч.

В любом случае, рекомендуется устанавливать значение linger.ms > 0, чтобы оптимизировать использование батчевания. Для большинства приложений подходит значение от 5 до 50 миллисекунд (для конкретной системы значение подбирается экспериментальным путём). Таким образом, мы повышаем шансы того, что сообщения будут отправляться в виде батчей за счёт внедрения малозаметной задержки перед отправкой.

⚠️ Важно: если размер сообщения превышает значение batch.size, то оно будет отправлено в отдельном батче, где будет только одно это сообщение.

Хранение данных

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

Но как же с технической точки зрения реализовано хранение сообщений на брокере? В первом модуле мы обсуждали, что Kafka пишет сообщения в файл журнала (лог). Вероятно, вы уже знакомы с понятием логов или лог-файлов. Обычно под понятием "лог" (log) понимается один или несколько текстовых файлов, в которые программа автоматически записывает информацию о событиях, происходящих в системе. Например, в лог может сохраняться информация об ошибках, действиях пользователей, и так далее, в зависимости от предметной области.

В Kafka в лог-файлах хранятся сами сообщения. Однако, не стоит путать лог-файлы для хранения сообщений и лог-файлы непосредственно самой Apache Kafka, как приложения.

Пример логов некоторого приложения

Предположим, что у нас имеется один Kafka-брокер. У него есть файловая система. Где же хранятся лог-файлы с сообщениями? Мы можем узнать это, посмотрев настройки брокера. Как уже было сказано, они хранятся в файле server.properties. В нём, помимо множества других настроек, указана директория для хранения лог-файлов (log.dirs). Если бы у нас было несколько брокеров в кластере, то для каждого бы присутствовал отдельный файл настроек.

По умолчанию, директории с лог-файлами расположены внутри tmp/kafka-logs . Внутри данной директории находятся вложенные директории, названия которых строятся по следующему правилу:

-,

где - это название топика, а - номер партиции данного топика.

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

В предыдущем модуле мы создавали топик "two.partitions.topic", состоящий из 2-х партиций. Значит, внутри tmp/kafka-logs у нас будут следующие директории:

• two.partitions.topic-0

• two.partitions.topic-1

Сообщения, которые были отправлены продьюсером(-ми) в топик с номером 0, будут храниться в директории "two.partitions.topic-0", а те, которые были отправлены в топик с номером 1, - в директории "two.partitions.topic-1".

Давайте проверим действительно ли данные директории присутствуют в /tmp/kafka-logs:

Отлично, теперь давайте посмотрим что хранится внутри вложенных директорий:

Мы видим несколько файлов. Сейчас нас интересуют только три из них, с расширениями .log, .index и .timeindex. Давайте разберёмся для чего нужен каждый из них.

Файл с расширением .log содержит сами сообщения (а точнее - батчи сообщений), которые отправляют продьюсеры. Это и есть, так называемый, лог-файл. Упрощённо, можно считать, что сообщения хранятся следующим образом:

offset position CreateTime payload

0 0 1714557190000 First message

1 67 1714564086000 Learning Kafka

Как мы уже знаем, offset - это порядковый номер сообщения в партиции (в нашем случае мы рассматриваем партицию с номером 0 в топике "two.partitions.topic"). Напомню, что значения офсета, так же как и нумерация партиций, начинаются с нуля.

position указывает на место в файле, с которого начинается данное сообщение. Представим, что все сообщения в лог-файле хранятся последовательно, при это каждое занимает определённое кол-во байт. position представляет собой "байтовое смещение" от начала файла, то есть, показывает на сколько байт нужно отступить от начала файла, чтобы попасть на определённое сообщение. Например, сообщение "Learning Kafka" начинается с 67-го байта в лог-файле.

CreateTime - это TimeStamp, то есть метка даты и времени сообщения, а payload - непосредственно содержимое самого сообщения (его value).

При наличии ключа, можно было бы добавить в нашу таблицу ещё одну колонку - "key".

В Kafka CLI, который мы изучали в предыдущем модуле, имеется специальная утилита (kafka-dump-log.sh) для просмотра содержимого файлов .log, .index и .timeindex.

Синтаксис утилиты: bin/kafka-dump-log.sh --files <путь_к_файлу>.

Давайте посмотрим на содержимое одного из лог-файлов (опция --print-data-log позволяет вывести в том числе сам payload):

bin/kafka-dump-log.sh --files /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.log --print-data-log

Вывод:

Dumping /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.log

Log starting offset: 0

baseOffset: 0 lastOffset: 0 count: 1 baseSequence: 0 lastSequence: 0 producerId: 0 producerEpoch: 0 partitionLeaderEpoch: 0 isTransactional: false isControl: false deleteHorizonMs: OptionalLong.empty position: 0 CreateTime: 1714557190000 size: 81 magic: 2 compresscodec: none crc: 1236464849 isvalid: true

| offset: 0 CreateTime: 1714557190000 keySize: -1 valueSize: 13 sequence: 0 headerKeys: [] payload: First message

baseOffset: 1 lastOffset: 1 count: 1 baseSequence: 1 lastSequence: 1 producerId: 0 producerEpoch: 0 partitionLeaderEpoch: 0 isTransactional: false isControl: false deleteHorizonMs: OptionalLong.empty position: 81 CreateTime: 1714564086000 size: 82 magic: 2 compresscodec: none crc: 3590501564 isvalid: true

| offset: 1 CreateTime: 1714564086000 keySize: -1 valueSize: 14 sequence: 1 headerKeys: [] payload: Learning Kafka

Во второй строке вывода вы видим, что в данном лог-файле содержатся сообщения, начиная с offset = 0.

Далее идут сами сообщения, сгруппированные по батчам. В нашем случае у нас имеются два батча, в каждом присутствует только 1 сообщение. Параметры батча следуют до символа "|", параметры сообщений, входящих в батч, идут после данного символа.

Рассмотрим основные параметры батча:

• baseOffset и lastOffset указывают на офсеты первого и последнего сообщения, входящих в данный батч;

• count - количество сообщений в данном батче;

• size - размер батча в байтах;

• compresscodec - алгоритм, который использовался для сжатия данного батча (соответствует compression.type)

Для сообщения имеется в том числе информация о размере его ключа и значения в байтах (keySize и valueSize) и массив заголовков headerKeys. Что такое заголовки мы рассмотрим в уроке 5.4.

Стоит отметить, что CreateTime хранится в формате Unix Epoch Time (в миллисекундах). Оно определяет кол-во миллисекунд, прошедших с "1 января 1970 года, 00:00:00" по UTC (Всемирному координированному времени) до нужного момента времени. В привычном нам виде у сообщения с offset = 0 метка даты и времени может быть представлена как: "1 мая 2024 года, 09:53:10" (по UTC).

Для перевода можно использовать один из онлайн-сервисов, например, epochconverter.com. Вставляем в поле значение CreateTime в миллисекундах и нажимаем "Timestamp to Human Date":

Стоит отметить, что временные метки в системе Unix могут быть также в секундах, микросекундах и наносекундах. В нашем случае сервис корректно определил, что метка указана в миллисекундах, и вывел дату и время в привычном нам формате, по UTC, и ниже - в моём часовом поясе (GMT+3).

Файл с расширением .index (индексный файл) состоит из пар соответствия offset (смещения) и position (позиции сообщения в лог-файле).

Файл .index содержит записи для сообщений из лог-файла (.log). Он необходим для того, чтобы консьюмеры могли быстро находить позицию сообщения в лог-файле по его офсету.

Ранее мы уже говорили, что консьюмеры обычно сохраняют (коммитят) номер офсета последнего успешно прочитанного сообщения. Это делается для того, чтобы в случае сбоя или перезагрузки сервиса (который выступает консьюмером), он мог продолжить чтение сообщений с того, которое ещё не было прочитано и обработано.

Здесь и пригождается индексный файл .index. Консьюмер, используя сохранённый офсет, находит в индексном файле позицию (position) первого необработанного сообщения и начинает читать сообщения из лог-файла, начиная с найденной позиции.

Без данного файла консьюмеру бы пришлось полностью читать лог-файл (.log), начиная с самого начала, что было бы крайне неэффективно.

Содержимое индексного файла .index можно представить следующим образом:

offset position

0 0

1 67

Теперь, давайте посмотрим на содержимое одного из .index-файлов для нулевой партиции топика "two.partitions.topic":

bin/kafka-dump-log.sh --files /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.index

Вывод:

Dumping /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.index

offset: 0 position: 0

Вывод чем-то похож на нашу таблицу. Но почему здесь всего одна строка?

Всё дело в том, что записи в индексную таблицу вносятся не для каждого сообщения из лог-файла, а раз в определённое время. Это контролируется настройкой брокера log.index.interval.bytes. Сохранять пару (offset, position) для каждого приходящего сообщения слишком затратно, это создаёт дополнительную нагрузку на диск.

Настройка log.index.interval.bytes определяет количество байт данных в лог-файле (.log), после которых Apache Kafka будет добавлять новую запись в индексный файл .index. Чем меньше это значение, тем чаще будут добавляться записи в файл и тем быстрее будет происходить поиск по офсету (но при этом возрастёт нагрузка на диск, а сам индексный файл будет очень быстро расти в объёме). Чем значение больше - тем медленнее будет происходить поиск нужной позиции, но при этом нагрузка на диск будет также меньше. В данном вопросе требуется баланс.

По умолчанию, значение log.index.interval.bytes равно 4096 байтам. То есть новая записи будет добавляться в индексный файл .index после каждых 4096 байт данных, записанных в .log-файл. Зачастую, значение по умолчанию покрывает все стандартные сценарии, так что можете сильно не волноваться по поводу него.

Файл с расширением .timeindex содержит пары, состоящие из CreateTime (метки даты и времени сообщения) и offset (смещения). Этот файл также называют "индексным".

Иногда возникает необходимость прочитать сообщения, начиная с определённой временной метки, например, с "10 января 2024 года 10:00:00". Индексный файл .timeindex необходим для того, чтобы консьюмеры могли быстро определить, с какой позиции в лог-файле следует начать чтение, если известна только временная метка (то есть, CreateTime сообщения).

Процесс выглядит следующим образом: выполняется поиск ближайшей временной метки в файле .timeindex, которая больше требуемой. Из найденной строки берётся offset. Далее происходит поиск позиции сообщения на основе этого офсета в индексном файле .index . После этого консьюмер может начать чтение лог-файла с нужной позиции.

Содержимое индексного файла .timeindex можно представить следующим образом:

CreateTime offset

1714557190000 0

1714564086000 1

Опять же, давайте посмотрим на содержимое реального .timeindex-файла:

bin/kafka-dump-log.sh --files /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.timeindex

Вывод:

Dumping /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.timeindex

timestamp: 0 offset: 0

Found timestamp mismatch in :/tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.timeindex

Index timestamp: 0, log timestamp: 1714557190000

В выводе фигурирует пара (timestamp, offset), где оба параметра равны нулю. Ниже мы также видим сообщение о том, что найдено несоответствие между .log и .timeindex-файлами: для offset = 0 в лог-файле проставлено другое значение timestamp (CreateTime), а не 0. Такое бывает, пока в .timeindex-файл не будет добавлена первая запись.

По факту, это никак не влияет на дальнейшую работу. Если мы отправим побольше сообщений в топик, то через время увидим несколько другую картину:

Dumping /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.timeindex

timestamp: 1715542926000 offset: 60

Частота добавления пар в файл .timeindex контролируется той же настройкой, что и для файла .index - log.index.interval.bytes. Она едина для двух индексных файлов.

Заголовки (headers)

В модуле 2. Основы Apache Kafka мы затронули тему такого опционального элемента сообщений, как "заголовки".

Заголовки (англ. headers) дают возможность добавлять некоторые метаданные, не добавляя никакой дополнительной информации ни в ключ, ни в значение самого сообщения (то есть заголовки - это абсолютно отдельный элемент сообщения).

Заголовки часто используются для:

• указания источника данных (например, названия приложения или сервиса, который отправил сообщение)

• "отлавливания" нужных сообщений из топика на основе информации из заголовка

Заголовки хороши тем, что консьюмер может прочитать их без необходимости обрабатывать (парсить) ключ и значение сообщения. Они задаются в виде упорядоченного списка пар (ключ, значение), например:

country : France,

privacy-level : 15

В данном случае "country" - это ключ заголовка, а "France" - его значение. Аналогично и с "privacy-level".

Ключи всегда являются строками, а значения могут быть любыми объектами (например, в формате JSON, Avro и так далее), точно так же, как и в случае со значениями сообщений (value).

Рассмотрим следующую ситуацию: у нас есть топик orders.status, куда отправляются статусы всех пользовательских заказов. Предположим, что у нас есть две системы. Одной - нужны все статусы заказов. Другая - должна отправлять пуш-уведомление на телефон пользователя, когда статус заказа становится "Доставлен в пункт выдачи" (delivered). Иными словами, второй системе интересны только сообщения, в которых статус заказа "delivered". Если статус заказа будет содержаться в значении сообщения, то консьюмеру второй системы придётся сначала распарсить это значение, а затем уже принять решение о дальнейшей обработке сообщения на основе статуса заказа. Хорошим решением может быть дублирование статуса заказа в заголовке - таким образом консьюмер второй системы будет фильтровать сообщения, поступающие из топика, по заголовкам и сразу пропускать те, которые для него не предназначены. При этом первая система, который нужны все сообщения, может в принципе не смотреть на заголовки и сразу обрабатывать все поступающие в топик сообщения.

Стоит упомянуть, что Kafka CLI, к сожалению, не позволяет отправлять сообщения с заголовками.

zookeeper что это

Топик (Topic)

Сообщения в Kafka отправляются в топики (темы). Их можно представить как таблицы в базе данных. Каждая таблица в БД зачастую хранит один вид данных - например, данные по заказам, данные по пользователям и так далее. С топиками всё работает аналогично. По сути, топик логически группирует сообщения. У каждого топика есть уникальное название.

Партиция (Partition)

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

Партиции нумеруются с нуля. Количество партииций в топике задаётся при создании топика, после этого их число можно увеличить (но уменьшить число партиций уже нельзя ⚠️). Для каждого топика число партиций устаналивается отдельно. Сообщения записываются в партицию путём добавления в конец, а читаются по порядку от начала к концу.

(!) Важно: сообщения в партиции хранятся в строгом порядке их поступления. То есть партиция - это фактически очередь сообщений. При этом порядок на уровне топика НЕ гарантируется!

Благодаря партициям, в Kafka обеспечивается возможность горизонтального масштабирования. Партиции можно реплицировать - то есть размещать копии одной и той же партиции на разных брокерах. Таким образом, если какой-то брокер в кластере выйдет из строя, мы не потеряем данные, которые были в партициях, размещённых на нём (поскольку у нас будут копии этих партиций на других брокерах). Подробнее про репликацию мы поговорим в уроке "2.6 Репликация данных".

Давайте рассмотрим следующую иллюстрацию. Здесь изображён один брокер, в котором находятся два топика - order.status и user.actions. Некий сервис-продьюсер записывает данные в топик order.status, в свою очередь два сервиса-консьюмера читают данные из этого топика.

Офсет (Offsets)

Офсет (или смещение) — это порядковый номер сообщения в определённой партиции. Он позволяет консьюмерам отслеживать, какие сообщения из данной партиции ими уже были прочитаны, а какие - ещё нет. Дело в том, что Kafka не удаляет сообщения после их чтения (в отличие от RabbitMQ); вместо этого консьюмеры отслеживают последний офсет, который они обработали. Они могут перечитывать данные или начать чтение с определённого места (офсета) в партиции.

(!) Важно: офсет считается с нуля.

Пример

Предположим, что у нас есть топик "order.status", в который отправляются статусы заказов с сайта интернет-магазина. У нас имеется консьюмер, в качестве которого выступает сервис "stats-saver", который читает топик со статусами заказов и сохраняет статистику по этим заказам в базу данных.

После прочтения каждого сообщения сервис "stats-saver" будет сохранять его offset (порядковый номер внутри партиции топика) на брокере. Благодаря этому, если с сервисом что-то случится (например, на нём возникнет ошибка и он перезагрузится), он сможет обратиться к брокеру и узнать на каком offset он закончил чтение партиции. После этого сервис сможет начать чтение с того сообщения, которое сервис ещё не читал (т.е. с offset + 1).

Сообщение (message)

Используемая в Kafka единица данных называется сообщением (message). Если вы ранее работали с базами данных, то можете рассматривать сообщение как аналог строки (row) или записи (record) в таблице БД. С точки зрения Kafka сообщение представляет собой набор байтов, так что для неё содержащиеся в нём данные не имеют формата и какого-либо смысла.

Пример сообщения

В данном сообщении передаётся информация о стоимости акции компании "Сбербанк" на определённую дату и время:

· name - название акции

· price - стоимость акции

· time - дата и время, когда акция имела такую стоимость

{

name: "SBER",

price: 285.5,

time: "17-05-2023T15:00:00"

}

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

Пример ключа и значения:

Давайте выделим название акции в ключ сообщения. Остальные данные оставим в значении.

Ключ (key)

{

name: "SBER"

}

Значение (value)

{

price: 285.5,

time: "17-05-2023T15:00:00"

}

Помимо ключа и значения сообщения в Kafka содержат Timestamp (время создания сообщения) и Headers (заголовки - опциональный набор метаданных). Про заголовки мы ещё поговорим в одном из следующих уроков.

Форматы сообщений

В примерах выше, сообщения представлены в популярном формате JSON (JavaScript Object Notion). Помимо JSON, для передачи данных посредством Kafka, также часто используются форматы AVRO (Apache AVRO) и Protobuf (Protocol Buffers). В рамках одного топика рекомендуется использовать только один формат данных и одну структуру сообщений.

Рассмотрим отличия данных форматов:

Из всех указанных форматов, вероятно, вы уже хорошо знакомы с JSON. Его часто используют в REST API. Он удобен в первую очередь благодаря тому, что является человекочитаемым (текстовым). Мы можем легко взглянуть на сообщение и понять что за данные в нём содержатся. AVRO и Protobuf в свою очередь являются бинарными. То есть сообщения в этих форматах выглядят как набор байтов. Но в этом также заключается их главный плюс, сообщения в этих форматах занимают меньше места, их быстрее читать и записывать, а также быстрее передавать по сети между разными сервисами.

Не пугайтесь, если не знаете что такое схема. Это понятие мы рассмотрим на следующем шаге урока.

Схема (Scheme)

Схема позволяет описать структуру данных: какие поля в ней присутствуют, какие у них типы и т.д.

Почему это удобно? Например, мы можем задать в коде сервиса-консьюмера как будут обрабатываться те сообщения, которые не прошли проверки с помощью схемы. Например, мы можем пропускать подобные сообщения. Сценарий обработки может зависеть от конкретных требований к сервису.

Такое определение может показаться довольно сложным. Давайте рассмотрим конкретные примеры.

JSON

Стандартом описания схемы для данных в формате JSON является JSON Schema.

Предположим, у нас есть топик user.data, в котором пересылаются следующие данные о пользователях социальной сети в формате JSON:

· id - уникальный идентификатор пользователя

· first_name - имя пользователя

· last_name - фамилия пользователя

{

"id": 210700286,

"first_name": "Artem",

"last_name": "Sidorov"

}

JSON Schema, описывающая данный объект данных может выглядеть так:

{

"type": "object",

"properties": {

"id": {

"type": "integer",

"description": "User ID"

},

"first_name": {

"type": "string",

"description": "User first name"

},

"last_name": {

"type": "string",

"description": "User last name"

}

},

"required": [

"id",

"first_name",

"last_name"

],

"additionalProperties": false

}

Внутри элемента properties для каждого поля данных задаются типы данных, которым должны соответствовать поля сообщения. Например, поле id должно быть типа integer (целое число), а поле first_name - типа string (строка).

Ключевое слово required задает перечень обязательных полей. Если хотя бы одно из перечисленных полей будет отсутствовать, сообщение не пройдет валидацию (проверку) по такой схеме.

Ключевое слово additionalProperties задает возможность наличия дополнительных полей у объекта. В нашем случае дополнительные поля запрещены (при их наличии объект не пройдет валидацию по схеме).

Подробнее про JSON-схемы можно узнать в официальной документации (на английском языке) - https://json-schema.org/learn/getting-started-step-by-step#define.

Protobuf

Формат Protobuf (Protocols Buffer) был придуман в Google. Данный формат использует собственный язык описания структуры данных - Protobuf IDL (Interface Description Language).

Предположим, что в топике search.queries передаются поисковые запросы пользователей к Yandex, в том числе:

· query - сам поисковый запрос

· page_number - номер поисковой страницы

· results_per_page - кол-во результатов поиска на каждой странице

Пример схемы на языке IDL:

syntax = "proto3";

message SearchRequest {

string query = 1;

int32 page_number = 2;

int32 results_per_page = 3;

}

Первая строка схемы указывает, что далее последует описание схемы на языке IDL 3-й версии.

Далее внутри message описывается структура самого сообщения. SearchRequest - это название сообщения.

Внутри указываются - название поля, его тип, а также номер. Например, поле query имеет строковый тип (string) и номер "1". Номера используются для идентификации полей в сообщении, поэтому они должны быть уникальны.

На основе схемы с помощью специального компилятора (protoc) можно сгенерировать готовые классы сообщений, методы сериализации/десереализации, а также другие вспомогательные функции для вашего языка программирования (protoc поддерживает множество языков, включая Java, C++, Go и многие другие).

Создание топика

Начнём освоение Kafka CLI с создания топика. Для этого воспользуемся утилитой kafka-topics.sh. Эта утилита также используется для удаления, изменения настроек и получения информации о топиках.

Давайте создадим топик orders.status. Для этого мы указываем следующие опции:

· create - указывает, что мы собираемся создать топик

· topic - после неё указывается название топика

· bootstrap-server - после неё указывается адрес брокера (-ов) Kafka, к которым мы будем подключаться

bin/kafka-topics.sh --create --topic orders.status --bootstrap-server localhost:9092

Поскольку мы подключаемся к брокеру Kafka, который запущен на том же самом компьютере, на котором мы создаём консьюмера, то мы указываем в качестве хоста localhost, а в качестве порта 9092 (стандартный порт, на котором по-умолчанию запускается брокер Kafka). В нашем случае у нас всего один брокер. Если у нас имеется несколько брокеров (кластер), то хосты и порты брокеров указываются через запятую.

Если всё прошло успешно, то мы должны увидеть следующее сообщение в Терминале:

Created topic orders.status.

Оно говорит нам о том, что топик orders.status был успешно создан. При этом топик будет состоять из одной партиции.

Не пугайтесь, если увидите следующее предупреждение (WARNING):

WARNING: Due to limitations in metric names, topics with a period ('.') or underscore ('_') could collide. To avoid issues it is best to use either, but not both.

Оно говорит о том, что в некоторых случаях могут быть проблемы, если вы используете в наименованиях своих топиков и точку, и нижнее подчёркивание и рекомендует использовать только один из этих символов для наименования всех своих топиков. Мы, в свою очередь, будем использовать точку.

Вывод списка топиков

Список существующих топиков на брокере (-ах) можно вывести с помощью опции list:

bin/kafka-topics.sh --list --bootstrap-server localhost:9092

Пока что у нас должен быть только один топик, который мы создали ранее - orders.status:

Указание количества партиций

При создании топика мы можем указать количество партиций с помощью опции partitions. Давайте создадим топик с двумя партициями.

bin/kafka-topics.sh --create --topic two.partitions.topic --bootstrap-server localhost:9092 --partitions 2

Если не указать опцию partitions, то будет создан топик, состоящий только из одной партиции.

Указание коэффициента репликации

При создании топика мы также можем указать коэффициент репликации. За это отвечает опция replication-factor.

Например:

bin/kafka-topics.sh --create --topic replication.factor.topic --bootstrap-server localhost:9092 --replication-factor 3

В нашем случае мы не сможем создать топик с таким коэффициентом репликации, поскольку у нас имеется всего один брокер. Если мы всё же попытаемся выполнить данную команду, то получим следующую ошибку:

Error while executing topic command : Replication factor: 3 larger than available brokers: 1.

[2024-04-15 17:00:00,000] ERROR org.apache.kafka.common.errors.InvalidReplicationFactorException: Replication factor: 3 larger than available brokers: 1.

(org.apache.kafka.tools.TopicCommand)

⚠️ Запомните: коэффициент репликации не может превышать количество брокеров в кластере!

Таким образом, для обеспечения replication-factor = 3 нам потребуется иметь 3 брокера в кластере.

Указание внутренних настроек топика

Как мы уже обсуждали ранее, у топиков, брокеров, продьюсеров и консьюмеров есть свои настройки. Чтобы задать значения для настроек при создании топика посредством Kafka CLI, нужно воспользоваться опцией config.

Давайте ограничим максимальный размер сообщения (в байтах), которое можно будет записать в топик. Это контролируется настройкой топика max.message.bytes. Давайте зададим значение в 4000 байт (это почти 4 КБ):

bin/kafka-topics.sh --create --topic one.config.topic --bootstrap-server localhost:9092 --config max.message.bytes=4000

Если требуется указать сразу несколько настроек для топика, то каждую из них нужно прописывать в виде отдельного config, то есть в следующем виде:

bin/kafka-topics.sh --create --topic multi.config.topic --bootstrap-server localhost:9092 --config max.message.bytes=4000 --config cleanup.policy=compact

Получение описания топика

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

Давайте выведем описание для топика two.partitions.topic, который мы ранее создали с 2-мя партициями.

bin/kafka-topics.sh --describe --topic two.partitions.topic --bootstrap-server localhost:9092

Результат:

Topic: two.partitions.topic TopicId: 9ZBGUW4aRnKHdN7qAL6a5Q PartitionCount: 2 ReplicationFactor: 1 Configs:

Topic: two.partitions.topic Partition: 0 Leader: 0 Replicas: 0 Isr: 0

Topic: two.partitions.topic Partition: 1 Leader: 0 Replicas: 0 Isr: 0

Что мы видим в выводе?

· Название топика (Topic)

· Уникальный идентификатор топика (TopicId)

· Количество партиций (PartitionCount)

· Коэффициент репликации (ReplicationFactor)

· Настройки топика (Configs)

В разрезе партиций выводится:

· Номер партиции (Partition)

· Номер брокера, на котором находится ведущая реплика партиции (Leader)

· Номера брокеров, на которых находятся реплики данной партиции (Replicas)

· Номера брокеров, на которых находятся синхронизированные (in-sync) реплики партиции (ReplicationFactor)

Изменение количества партиций в топике

Если мы хотим изменить количество патриций в топике, то для этого понадобится опция alter.

Давайте сделаем это для топика orders.status. Напомню, что при его создании мы не указывали количество партиций явно, следовательно сейчас топик состоит из одной партиции. Увеличим количество партиций до 5-ти:

bin/kafka-topics.sh --alter --topic orders.status --partitions 5 --bootstrap-server localhost:9092

Проверим, что их число действительно поменялось с помощью опции describe:

bin/kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders.status

Вывод:

Topic: orders.status TopicId: 9p2uaygsQhy-oIqp-SLsPg PartitionCount: 5 ReplicationFactor: 1 Configs:

Topic: orders.status Partition: 0 Leader: 0 Replicas: 0 Isr: 0

Topic: orders.status Partition: 1 Leader: 0 Replicas: 0 Isr: 0

Topic: orders.status Partition: 2 Leader: 0 Replicas: 0 Isr: 0

Topic: orders.status Partition: 3 Leader: 0 Replicas: 0 Isr: 0

Topic: orders.status Partition: 4 Leader: 0 Replicas: 0 Isr: 0

⚠️ Важно: не забывайте, что после создания топика, число его партиций можно увеличить, но нельзя уменьшить!

Удаление топика

Для удаления топика используется опция delete:

bin/kafka-topics.sh --delete --topic orders.status --bootstrap-server localhost:9092

Можно удалить сразу несколько топиков, перечислив их через запятую:

bin/kafka-topics.sh --delete --topic one.config.topic,multi.config.topic --bootstrap-server localhost:9092

Мы можем убедиться, что топики были удалены с помощью опции list:

bin/kafka-topics.sh --list --bootstrap-server localhost:9092

Работа с продьюсерами происходит посредством утилиты kafka-console-producer.sh.

Отправка сообщения

Давайте создадим продьюсера, который будет отправлять сообщения в топик orders.status.

Поскольку ранее мы удалили данный топик, давайте сперва пересоздадим его заново (например, с 4-мя партициями):

bin/kafka-topics.sh --create --topic orders.status --partitions 4 --bootstrap-server localhost:9092

Теперь создадим продьюсера для данного топика:

bin/kafka-console-producer.sh --topic orders.status --bootstrap-server localhost:9092

После этого вы должны увидеть в консоли символ ">". Все строчки текста, которые вы напишите после него, будут отправлены в топик orders.status. Конец строки определяется нажатием кнопки "Enter".

bin/kafka-console-producer.sh --topic orders.status --bootstrap-server localhost:9092

>First message

>Learning Kafka

>Love learning

>

⚠️ Важно:

1) Обратите внимание, что по-умолчанию, сообщения отправляются с пустым ключом.

2) Если вы создадите продьюсер для несуществующего топика, то он будет мгновенно создан. При этом его количество партиций и коэффициент репликаций будут соответствовать значениям по-умолчанию, которые установлены на уровне брокера (про эти настройки мы поговорим в следующем модуле).

Отправка сообщения с ключом

Для того, чтобы отправить сообщение с ключом, нам следует задать для продьюсера ряд параметров с помощью опций property:

· parse.key (нужно ли обрабатывать ключ - да/нет)

· key.separator (символ, разделяющий ключ и значение сообщения)

Откройте новую сессию в Терминале и выполните следующую команду:

bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic topic.with.key --property parse.key=true --property key.separator=:

Отправим несколько сообщений с ключом (в качестве разделителя ключа и значения мы указали символ ":").

>first key:value1

>first:last

>

Работа с консьюмерами происходит посредством утилиты kafka-console-consumer.sh.

Получение актуальных сообщений

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

bin/kafka-console-consumer.sh --topic orders.status --bootstrap-server localhost:9092

После выполнения команды вы должны увидеть символ ">" в консоли. Однако, вы не увидите сообщений, которые мы ранее отправили в топик. Если мы создаём консьюмер вышеуказанным образом, то мы увидим только те сообщения, которые были отправлены в топик после создания консьюмера!

Получение исторических сообщений

Если мы хотим сперва прочитать исторические сообщения, а затем начать читать актуальные, следует воспользоваться опцией from-beginning.

bin/kafka-console-consumer.sh --topic orders.status --from-beginning --bootstrap-server localhost:9092

bin/kafka-console-consumer.sh --topic orders.status --from-beginning --bootstrap-server localhost:9092

First message

Learning Kafka

Love learning

Урок 4.4

Остановить консьюмер мы можем с помощью сочетания клавиш Ctrl + C. При этом будет выведено сообщение с информацией о количестве обработанных сообщений, например:

^CProcessed a total of 4 messages

⚠️ Важно:

1) По умолчанию, консьюмер не выводит ключ сообщения.

2) Обратите внимание, что мы не можем удалить топик, если его читает консьюмер (-ы).

Вывод ключа и метки времени у сообщения

Для того, чтобы помимо значения вывести ключ сообщения, нам требуется воспользоваться опцией property и указать параметр print.key, равный true (при этом также следует прописать print.value=true, чтобы вывелось и значение сообщения):

bin/kafka-console-consumer.sh --topic orders.status --property print.key=true --property print.value=true --bootstrap-server localhost:9092

Для того, чтобы вывести ещё и метку времени, требуется указать прописать print.timestamp=true:

bin/kafka-console-consumer.sh --topic orders.status --property print.key=true --property print.value=true --property print.timestamp=true --bootstrap-server localhost:9092 --from-beginning

Вывод:

CreateTime:1713977447000 null First message

CreateTime:1713977507000 null Learning Kafka

CreateTime:1713977747000 null Love learning

CreateTime:1713979847000 null Урок 4.4

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

Обратите внимание на формат времени, в котором записано CreateTime. Он может показаться несколько странным, однако это достаточно популярный формат записи даты и времени - Unix Timestamp. Он показывает количество миллисекунд (или иногда секунд) с начала эпохи Unix (то есть с 1 января 1970 года).

Батчи и батчевание

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

Давайте познакомиться с понятием "батч" и "батчевание". Батч (batch) - это группа сообщений, объединённых вместе для отправки как единое целое. Зачастую продьюсер не отправляет каждое сообщение отдельно, а сначала ждёт, пока сформируется группа из нескольких сообщений (message set), и только потом отправляет их как единое целое. Для чего это нужно? Такой подход позволяет сделать передачу данных более эффективной за счёт снижения накладных расходов на сетевые запросы.

Можно провести аналогию с автобусами. Представьте, что автобус ждёт, пока не соберётся группа пассажиров, прежде чем отправиться в путь. Это похоже на батчевание в Kafka: один автобус перевозит много пассажиров за один раз, снижая накладные расходы (на топливо и пр.) и делая поездки более эффективными. Когда автобус приезжает на конечную остановку (партицию топика), все пассажиры (сообщения) выгружаются разом.

Сверху - сообщения отправляются по отдельности (1 сообщение в одном батче),

снизу - сообщения отправляются в одном батче (4 сообщения в одном батче)

Продьюсер отправляет батч, когда его размер достигает заданного размера batch.size или по прошествии времени, установленного в linger.ms . То есть, если хоть одно из условий выполнится, то батч будет немедленно отправлен продьюсером. Обе опции являются главными настройками продьюсера, контролирующими батчевание.

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

Представим, что у нас есть топик с двумя партициями и два брокера. Каждая партиция топика расположена на одном из брокеров (партиция 0 на брокере 0, а партиция 1 на брокере 1). Сообщения, предназначенные для партиции 0, можно представить как пассажиров, которые хотят добраться до города "А", а те, что предназначены для партиции 1 - до города "В". Пассажиры (сообщения), направляющиеся в разные города (партиции), не могут ехать вместе в одном автобусе (батче), так как каждый автобус (батч) следует в конкретный город (определённую партицию конкретного брокера).

По умолчанию, batch.size равен 16 384 байтам, а linger.ms - 0 миллисекунд. На первый взгляд, кажется, что при linger.ms = 0 все сообщения будут отправляться моментально, а значит - по отдельности (1 сообщение в одном батче). Однако зачастую продьюсер собирает несколько сообщений практически одновременно, поэтому даже при linger.ms = 0 они могут быть объединены в батч.

В любом случае, рекомендуется устанавливать значение linger.ms > 0, чтобы оптимизировать использование батчевания. Для большинства приложений подходит значение от 5 до 50 миллисекунд (для конкретной системы значение подбирается экспериментальным путём). Таким образом, мы повышаем шансы того, что сообщения будут отправляться в виде батчей за счёт внедрения малозаметной задержки перед отправкой.

⚠️ Важно: если размер сообщения превышает значение batch.size, то оно будет отправлено в отдельном батче, где будет только одно это сообщение.

Хранение данных

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

Но как же с технической точки зрения реализовано хранение сообщений на брокере? В первом модуле мы обсуждали, что Kafka пишет сообщения в файл журнала (лог). Вероятно, вы уже знакомы с понятием логов или лог-файлов. Обычно под понятием "лог" (log) понимается один или несколько текстовых файлов, в которые программа автоматически записывает информацию о событиях, происходящих в системе. Например, в лог может сохраняться информация об ошибках, действиях пользователей, и так далее, в зависимости от предметной области.

В Kafka в лог-файлах хранятся сами сообщения. Однако, не стоит путать лог-файлы для хранения сообщений и лог-файлы непосредственно самой Apache Kafka, как приложения.

Пример логов некоторого приложения

Предположим, что у нас имеется один Kafka-брокер. У него есть файловая система. Где же хранятся лог-файлы с сообщениями? Мы можем узнать это, посмотрев настройки брокера. Как уже было сказано, они хранятся в файле server.properties. В нём, помимо множества других настроек, указана директория для хранения лог-файлов (log.dirs). Если бы у нас было несколько брокеров в кластере, то для каждого бы присутствовал отдельный файл настроек.

По умолчанию, директории с лог-файлами расположены внутри tmp/kafka-logs . Внутри данной директории находятся вложенные директории, названия которых строятся по следующему правилу:

-,

где - это название топика, а - номер партиции данного топика.

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

В предыдущем модуле мы создавали топик "two.partitions.topic", состоящий из 2-х партиций. Значит, внутри tmp/kafka-logs у нас будут следующие директории:

· two.partitions.topic-0

· two.partitions.topic-1

Сообщения, которые были отправлены продьюсером(-ми) в топик с номером 0, будут храниться в директории "two.partitions.topic-0", а те, которые были отправлены в топик с номером 1, - в директории "two.partitions.topic-1".

Давайте проверим действительно ли данные директории присутствуют в /tmp/kafka-logs:

Отлично, теперь давайте посмотрим что хранится внутри вложенных директорий:

Мы видим несколько файлов. Сейчас нас интересуют только три из них, с расширениями .log, .index и .timeindex. Давайте разберёмся для чего нужен каждый из них.

Файл с расширением .log содержит сами сообщения (а точнее - батчи сообщений), которые отправляют продьюсеры. Это и есть, так называемый, лог-файл. Упрощённо, можно считать, что сообщения хранятся следующим образом:

Как мы уже знаем, offset - это порядковый номер сообщения в партиции (в нашем случае мы рассматриваем партицию с номером 0 в топике "two.partitions.topic"). Напомню, что значения офсета, так же как и нумерация партиций, начинаются с нуля.

position указывает на место в файле, с которого начинается данное сообщение. Представим, что все сообщения в лог-файле хранятся последовательно, при это каждое занимает определённое кол-во байт. position представляет собой "байтовое смещение" от начала файла, то есть, показывает на сколько байт нужно отступить от начала файла, чтобы попасть на определённое сообщение. Например, сообщение "Learning Kafka" начинается с 67-го байта в лог-файле.

CreateTime - это TimeStamp, то есть метка даты и времени сообщения, а payload - непосредственно содержимое самого сообщения (его value).

При наличии ключа, можно было бы добавить в нашу таблицу ещё одну колонку - "key".

В Kafka CLI, который мы изучали в предыдущем модуле, имеется специальная утилита (kafka-dump-log.sh) для просмотра содержимого файлов .log, .index и .timeindex.

Синтаксис утилиты: bin/kafka-dump-log.sh --files <путь_к_файлу>.

Давайте посмотрим на содержимое одного из лог-файлов (опция --print-data-log позволяет вывести в том числе сам payload):

bin/kafka-dump-log.sh --files /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.log --print-data-log

Вывод:

Dumping /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.log

Log starting offset: 0

baseOffset: 0 lastOffset: 0 count: 1 baseSequence: 0 lastSequence: 0 producerId: 0 producerEpoch: 0 partitionLeaderEpoch: 0 isTransactional: false isControl: false deleteHorizonMs: OptionalLong.empty position: 0 CreateTime: 1714557190000 size: 81 magic: 2 compresscodec: none crc: 1236464849 isvalid: true

| offset: 0 CreateTime: 1714557190000 keySize: -1 valueSize: 13 sequence: 0 headerKeys: [] payload: First message

baseOffset: 1 lastOffset: 1 count: 1 baseSequence: 1 lastSequence: 1 producerId: 0 producerEpoch: 0 partitionLeaderEpoch: 0 isTransactional: false isControl: false deleteHorizonMs: OptionalLong.empty position: 81 CreateTime: 1714564086000 size: 82 magic: 2 compresscodec: none crc: 3590501564 isvalid: true

| offset: 1 CreateTime: 1714564086000 keySize: -1 valueSize: 14 sequence: 1 headerKeys: [] payload: Learning Kafka

Во второй строке вывода вы видим, что в данном лог-файле содержатся сообщения, начиная с offset = 0.

Далее идут сами сообщения, сгруппированные по батчам. В нашем случае у нас имеются два батча, в каждом присутствует только 1 сообщение. Параметры батча следуют до символа "|", параметры сообщений, входящих в батч, идут после данного символа.

Рассмотрим основные параметры батча:

· baseOffset и lastOffset указывают на офсеты первого и последнего сообщения, входящих в данный батч;

· count - количество сообщений в данном батче;

· size - размер батча в байтах;

· compresscodec - алгоритм, который использовался для сжатия данного батча (соответствует compression.type)

Для сообщения имеется в том числе информация о размере его ключа и значения в байтах (keySize и valueSize) и массив заголовков headerKeys. Что такое заголовки мы рассмотрим в уроке 5.4.

Стоит отметить, что CreateTime хранится в формате Unix Epoch Time (в миллисекундах). Оно определяет кол-во миллисекунд, прошедших с "1 января 1970 года, 00:00:00" по UTC (Всемирному координированному времени) до нужного момента времени. В привычном нам виде у сообщения с offset = 0 метка даты и времени может быть представлена как: "1 мая 2024 года, 09:53:10" (по UTC).

Для перевода можно использовать один из онлайн-сервисов, например, epochconverter.com. Вставляем в поле значение CreateTime в миллисекундах и нажимаем "Timestamp to Human Date":

Стоит отметить, что временные метки в системе Unix могут быть также в секундах, микросекундах и наносекундах. В нашем случае сервис корректно определил, что метка указана в миллисекундах, и вывел дату и время в привычном нам формате, по UTC, и ниже - в моём часовом поясе (GMT+3).

Файл с расширением .index (индексный файл) состоит из пар соответствия offset (смещения) и position (позиции сообщения в лог-файле).

Файл .index содержит записи для сообщений из лог-файла (.log). Он необходим для того, чтобы консьюмеры могли быстро находить позицию сообщения в лог-файле по его офсету.

Ранее мы уже говорили, что консьюмеры обычно сохраняют (коммитят) номер офсета последнего успешно прочитанного сообщения. Это делается для того, чтобы в случае сбоя или перезагрузки сервиса (который выступает консьюмером), он мог продолжить чтение сообщений с того, которое ещё не было прочитано и обработано.

Здесь и пригождается индексный файл .index. Консьюмер, используя сохранённый офсет, находит в индексном файле позицию (position) первого необработанного сообщения и начинает читать сообщения из лог-файла, начиная с найденной позиции.

Без данного файла консьюмеру бы пришлось полностью читать лог-файл (.log), начиная с самого начала, что было бы крайне неэффективно.

Содержимое индексного файла .index можно представить следующим образом:

Теперь, давайте посмотрим на содержимое одного из .index-файлов для нулевой партиции топика "two.partitions.topic":

bin/kafka-dump-log.sh --files /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.index

Вывод:

Dumping /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.index

offset: 0 position: 0

Вывод чем-то похож на нашу таблицу. Но почему здесь всего одна строка?

Всё дело в том, что записи в индексную таблицу вносятся не для каждого сообщения из лог-файла, а раз в определённое время. Это контролируется настройкой брокера log.index.interval.bytes. Сохранять пару (offset, position) для каждого приходящего сообщения слишком затратно, это создаёт дополнительную нагрузку на диск.

Настройка log.index.interval.bytes определяет количество байт данных в лог-файле (.log), после которых Apache Kafka будет добавлять новую запись в индексный файл .index. Чем меньше это значение, тем чаще будут добавляться записи в файл и тем быстрее будет происходить поиск по офсету (но при этом возрастёт нагрузка на диск, а сам индексный файл будет очень быстро расти в объёме). Чем значение больше - тем медленнее будет происходить поиск нужной позиции, но при этом нагрузка на диск будет также меньше. В данном вопросе требуется баланс.

По умолчанию, значение log.index.interval.bytes равно 4096 байтам. То есть новая записи будет добавляться в индексный файл .index после каждых 4096 байт данных, записанных в .log-файл. Зачастую, значение по умолчанию покрывает все стандартные сценарии, так что можете сильно не волноваться по поводу него.

Файл с расширением .timeindex содержит пары, состоящие из CreateTime (метки даты и времени сообщения) и offset (смещения). Этот файл также называют "индексным".

Иногда возникает необходимость прочитать сообщения, начиная с определённой временной метки, например, с "10 января 2024 года 10:00:00". Индексный файл .timeindex необходим для того, чтобы консьюмеры могли быстро определить, с какой позиции в лог-файле следует начать чтение, если известна только временная метка (то есть, CreateTime сообщения).

Процесс выглядит следующим образом: выполняется поиск ближайшей временной метки в файле .timeindex, которая больше требуемой. Из найденной строки берётся offset. Далее происходит поиск позиции сообщения на основе этого офсета в индексном файле .index . После этого консьюмер может начать чтение лог-файла с нужной позиции.

Содержимое индексного файла .timeindex можно представить следующим образом:

Опять же, давайте посмотрим на содержимое реального .timeindex-файла:

bin/kafka-dump-log.sh --files /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.timeindex

Вывод:

Dumping /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.timeindex

timestamp: 0 offset: 0

Found timestamp mismatch in :/tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.timeindex

Index timestamp: 0, log timestamp: 1714557190000

В выводе фигурирует пара (timestamp, offset), где оба параметра равны нулю. Ниже мы также видим сообщение о том, что найдено несоответствие между .log и .timeindex-файлами: для offset = 0 в лог-файле проставлено другое значение timestamp (CreateTime), а не 0. Такое бывает, пока в .timeindex-файл не будет добавлена первая запись.

По факту, это никак не влияет на дальнейшую работу. Если мы отправим побольше сообщений в топик, то через время увидим несколько другую картину:

Dumping /tmp/kafka-logs/two.partitions.topic-0/00000000000000000000.timeindex

timestamp: 1715542926000 offset: 60

Частота добавления пар в файл .timeindex контролируется той же настройкой, что и для файла .index - log.index.interval.bytes. Она едина для двух индексных файлов.

Заголовки (headers)

В модуле 2. Основы Apache Kafka мы затронули тему такого опционального элемента сообщений, как "заголовки".

Заголовки (англ. headers) дают возможность добавлять некоторые метаданные, не добавляя никакой дополнительной информации ни в ключ, ни в значение самого сообщения (то есть заголовки - это абсолютно отдельный элемент сообщения).

Заголовки часто используются для:

· указания источника данных (например, названия приложения или сервиса, который отправил сообщение)

· "отлавливания" нужных сообщений из топика на основе информации из заголовка

Заголовки хороши тем, что консьюмер может прочитать их без необходимости обрабатывать (парсить) ключ и значение сообщения. Они задаются в виде упорядоченного списка пар (ключ, значение), например:

country : France,

privacy-level : 15

В данном случае "country" - это ключ заголовка, а "France" - его значение. Аналогично и с "privacy-level".

Ключи всегда являются строками, а значения могут быть любыми объектами (например, в формате JSON, Avro и так далее), точно так же, как и в случае со значениями сообщений (value).

Рассмотрим следующую ситуацию: у нас есть топик orders.status, куда отправляются статусы всех пользовательских заказов. Предположим, что у нас есть две системы. Одной - нужны все статусы заказов. Другая - должна отправлять пуш-уведомление на телефон пользователя, когда статус заказа становится "Доставлен в пункт выдачи" (delivered). Иными словами, второй системе интересны только сообщения, в которых статус заказа "delivered". Если статус заказа будет содержаться в значении сообщения, то консьюмеру второй системы придётся сначала распарсить это значение, а затем уже принять решение о дальнейшей обработке сообщения на основе статуса заказа. Хорошим решением может быть дублирование статуса заказа в заголовке - таким образом консьюмер второй системы будет фильтровать сообщения, поступающие из топика, по заголовкам и сразу пропускать те, которые для него не предназначены. При этом первая система, который нужны все сообщения, может в принципе не смотреть на заголовки и сразу обрабатывать все поступающие в топик сообщения.

Стоит упомянуть, что Kafka CLI, к сожалению, не позволяет отправлять сообщения с заголовками.

zookeeper что это

Черновик Звягин

Нагрузочное тестирование (НТ) — это оценка производительности программного продукта (системы) путем увеличения нагрузки на него из вне с целью установления

соответствия ожидаемого результата поведения системы (то, которое необходимо на основании требований) и фактического (то, которое происходит на момент проведения тестирования). Помогает выявить слабые места, предотвратить сбои во время пиковой активности системы. Простыми словами, НТ - проверка работоспособности системы под различными нагрузками на нее и анализ ее поведения.