BI Consult Desktop Logo BI Consult Mobile Logo
  • Russian BI Исследование российских bi
  • Перейти на Fine BI
  • Контакты
  • +7 812 334-08-01
    +7 499 608-13-06
  • Отправить сообщение
  • Главная
  • Продукты Эксперт-BI
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Сельское хозяйство
    • Энергетика
    • FMCG
    • Девелоперы
    • Маркетплейсы
    • Пищевая промышленность
    • Фармацевтика
    • Построение Data Platform
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и FP&A
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • IBP
    • ИТ (CIO)
    • Закупки
  • Платформы
    • Системы бизнес-анализа (BI)
    • Интегрированное бизнес-планирование (IBP)
    • Хранилища данных (DWH / Lakehouse)
    • Каталоги данных (Data Catalog)
    • Системы ETL и ELT
    • AI / Исскуственный интеллект
    • Шина данных (ESB)
    • Система управления мастер-данными (MDM)
    • Семантический слой
  • Услуги
    • Переход на отечественные BI и DWH системы
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений и DWH
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Курсы
    • Учебный курс Информационная грамотность (Data Literacy)
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Greenplum
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt (Data Build Tool)
  • Компания
    • Руководство
    • Новости
    • Клиенты
    • Карьера
    • Скачать
    • Контакты

BI

  • FineBI
  • FineReport
  • FineDataLink
  • FineChatBI (FineAI)
  • Коннекторы данных из 1С в BI
  • Airflow / Nifi
  • Visiology
  • PIX BI
  • Modus BI
  • Yandex.DataLens
  • Open-source BI: Superset/Metabase
  • Luxms BI
  • AW BI + Alpha BI
  • FlyBI + Форсайт. Аналитическая Платформа
  • Loginom
  • Триафлай
  • AI / Исскуственный интеллект
  • Optimacros
  • Навигатор BI
  • Семантический слой

СУБД

  • Arenadata
  • ClickHouse
  • Greenplum
  • Postgres Professional
  • TData

Другое

  • Построение Data Platform
    • Аналитическое хранилище данных
    • Data Lake и Data Engineering
    • Подробнее про Data Lake
    • Внедрение Lakehouse
      • Apache Doris
      • StarRocks
      • Trino
    • Миграция витрин из пропиетарных DWH на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс Современная архитектура хранилища данных » Введение в Apache Kafka со Spring Boot

Введение в Apache Kafka со Spring Boot

Apache Kafka - это распределенная платформа потоковой передачи данных, предназначенная для обработки больших объемов данных в режиме реального времени. Архитектура Kafka  состоит из четырех основных компонентов: продюсеров, брокеров, консюмеров и ZooKeeper.

  1. Продюсеры: Продюсеры - это клиенты, которые публикуют сообщения в топики Kafka. Ими могут быть любые приложения или системы, генерирующие данные, например датчики, веб-серверы или базы данных. Когда продюсер публикует сообщение в топике, он отправляет его брокеру Kafka.
  2. Брокеры: Брокеры - это серверы, которые образуют кластер Kafka. Они отвечают за прием сообщений от производителей, хранение сообщений в распределенном журнале фиксации и передачу сообщений консюмерам. Каждый брокер в кластере Kafka идентифицируется уникальным целым числом, называемым идентификатором брокера.
  3. Консюмеры: Консюмеры - это клиенты, которые читают сообщения из топиков Kafka. Ими могут быть любые приложения или системы, обрабатывающие данные, например, аналитические системы, поисковые системы или алгоритмы машинного обучения. Когда консюмер подписывается на топик, он сразу же получает сообщения из одного или нескольких разделов этого топика.
  4. ZooKeeper: ZooKeeper - это распределенный координационный сервис, который предназначен для управления состоянием кластера, выборов лидеров и выполнения различных административных задач. Он поддерживает список всех брокеров в кластере Kafka, а также их текущее состояние и метаданные.

 

 

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

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

Консюмеры могут подписаться на один или несколько топиков и читать сообщения из одного или нескольких разделов. Когда консюмер читает сообщения из раздела, он сохраняет свое собственное смещение, которое является позицией последнего использованного сообщения в разделе. Это позволяет консюмерам читать сообщения в своем собственном темпе и в случае сбоев возобновлять чтение с того места, на котором они остановились. Архитектура Kafka  представляет собой масштабируемую, отказоустойчивую и высокопроизводительную платформу для создания конвейеров данных, работающих в режиме реального времени, и потоковых приложений. Ее распределенный характер позволяет оперативно обрабатывать большие объемы данных и обеспечивать доступ к ним с низкой задержкой в различных приложениях и системах.

 

Разделы, Фактор репликации и период удержания

Разделы, коэффициент репликации и период удержания - важные понятия, которые определяют то, как данные хранятся, реплицируются и сохраняются в кластере Kafka.

  • Разделы: Топик Kafka делится на один или несколько разделов, которые по сути являются упорядоченными и неизменяемыми последовательностями записей. Каждый раздел может храниться на разных брокерах Kafka, что позволяет Kafka распределять нагрузку по чтению и записи данных между несколькими узлами. Это позволяет Kafka масштабироваться горизонтально и оперативно обрабатывать большие объемы данных.
  • Фактор репликации: Для обеспечения высокой доступности и отказоустойчивости Kafka реплицирует каждый раздел между несколькими брокерами. Количество реплик определяется коэффициентом репликации, который задает количество копий каждого раздела, которые необходимо поддерживать. Например, если коэффициент репликации равен 3, то каждый раздел будет реплицирован три раза через три разных брокера. Это означает, что даже если один или два брокера выйдут из строя, данные все равно будут доступны.
  • Период удержания: Топики Kafka могут быть настроены с учетом периода удержания или  периода хранения, который определяет, как долго Kafka должен хранить сообщения в топике перед их удалением. Это может быть полезно для отслеживания последних данных или для реализации политики хранения данных. Kafka поддерживает два типа периода удержания: по времени и по размеру. Политики, основанные на времени, удаляют сообщения, которые находились в топике в течение определенного времени, а политики, основанные на размере, удаляют сообщения на основе общего объема данных, хранящихся в топике.

 

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

 

Начало работы с Apache Kafka

Вот несколько шагов, которые помогут Вам быстро начать работу с Apache Kafka:

1. Скачайте и установите Apache Kafka: Вы можете скачать последнюю версию Apache Kafka с официального сайта. После загрузки следуйте инструкциям по установке.

2. Запустите сервер Kafka: После установки Kafka Вы сможете запустить сервер Kafka, выполнив следующую команду:

bin/kafka-server-start.sh config/server.properties

 

Данная команда запустит сервер Kafka  на Вашем локальном устройстве.

 

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

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

 

Эта команда создаст топик под названием my-topic на сервере Kafka, работающем с портом 9092.

 

4. Отправьте сообщения в опик: Вы можете отправлять сообщения в топик с помощью Kafka Producer API. Пример:

import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.KafkaProducer;
import java.util.Properties;

public class KafkaProducerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
        props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

        Producer<String, String> producer = new KafkaProducer<>(props);

        String message = "Hello, Kafka!";
        ProducerRecord<String, String> record = new ProducerRecord<>("my-topic", message);

        producer.send(record);
        producer.close();
    }
}

 

В этом примере мы создаем продюсера Kafka, который отправляет сообщение «Hello, Kafka!» в топик «my-topic».

 

5. Возьмите из топика сообщение и прочитайте его: Вы можете потреблять сообщения из топика, используя API консюмера Kafka. Пример:

import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import java.util.Collections;
import java.util.Properties;

public class KafkaConsumerExample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put("bootstrap.servers", "localhost:9092");
        props.put("key.deserializer", StringDeserializer.class.getName());
        props.put("value.deserializer", StringDeserializer.class.getName());
        props.put("group.id", "my-group");

        Consumer<String, String> consumer = new KafkaConsumer<>(props);
        consumer.subscribe(Collections.singleton("my-topic"));

        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(1000);
            records.forEach(record -> {
                System.out.println("Received message: " + record.value());
            });
        }
    }
}

 

В этои примере мы создали консмера,  которые читает сообщения из топика “my-topic” и выводит их на специальную консоль.

 

Пример:

Вот пример конечной точки Kafka Spring Boot,  которая передает событие в топик Kafka, как только а в базу данных добавляется новая запись о сотруднике:

Во-первых, Вам нужно добавить зависимости Kafka и Spring Data JPA в Ваш проект Spring Boot. Вы можете добавить их в файл pom.xml, если Вы используете Maven, или в файл build.gradle, если Вы используете Gradle

<!-- Kafka dependency -->
<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
    <version>2.7.4</version>
</dependency>

<!-- Spring Data JPA dependency -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-jpa</artifactId>
</dependency>

 

Создайте класс конфигурации Kafka, который определит продюсера.

@Configuration
@EnableKafka
public class KafkaConfig {

    @Bean
    public ProducerFactory<String, String> producerFactory() {
        Map<String, Object> configProps = new HashMap<>();
        configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return new DefaultKafkaProducerFactory<>(configProps);
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

 

Создайте службу Kafka producer, которая будет отправлять сообщения в топик Kafka.

@Service
public class KafkaProducerService {

    private final KafkaTemplate<String, String> kafkaTemplate;

    @Autowired
    public KafkaProducerService(KafkaTemplate<String, String> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void sendMessage(String message) {
        kafkaTemplate.send("employee-topic", message);
    }
}

 

Создайте экземпляр Spring Data JPA для записи о сотруднике.

@Entity
public class Employee {

    @Id
    @GeneratedValue(strategy = GenerationType.IDENTITY)
    private Long id;

    private String name;

    private String department;

    // getters and setters
}

 

Создайте JPA-репозиторий Spring Data для сотрудника.

@Repository
public interface EmployeeRepository extends JpaRepository<Employee, Long> {
}

 

Создайте контроллер Spring MVC, который будет обрабатывать HTTP POST-запрос для добавления новой записи о сотруднике в базу данных

@RestController
public class EmployeeController {

    private final EmployeeRepository employeeRepository;
    private final KafkaProducerService kafkaProducerService;

    @Autowired
    public EmployeeController(EmployeeRepository employeeRepository, KafkaProducerService kafkaProducerService) {
        this.employeeRepository = employeeRepository;
        this.kafkaProducerService = kafkaProducerService;
    }

    @PostMapping("/employees")
    public ResponseEntity<Employee> createEmployee(@RequestBody Employee employee) {
        Employee savedEmployee = employeeRepository.save(employee);
        kafkaProducerService.sendMessage("New employee added: " + savedEmployee.getName());
        return ResponseEntity.ok(savedEmployee);
    }
}

 

В приведенном выше примере, когда новая запись сотрудника добавляется в базу данных с помощью HTTP POST-запроса, метод createEmployee отправляет сообщение в топик Kafka с помощью kafkaProducerService. Сообщение содержит имя нового добавленного сотрудника. Вы можете настроить формат сообщения и имя топика в соответствии с Вашими требованиями.

Чтобы отправить только что добавленную запись о сотруднике в топик employee-topic, измените KafkaTemplate

@Service
public class KafkaProducerService {

    private static final String TOPIC_NAME = "employee-topic";
    private static final int NUM_PARTITIONS = 3;
    private static final short REPLICATION_FACTOR = 1;
    private static final long RETENTION_PERIOD_MS = 86400000L; // 1 day

    private final KafkaTemplate<Long, Employee> kafkaTemplate;

    @Autowired
    public KafkaProducerService(KafkaTemplate<Long, Employee> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    public void sendEmployee(Employee employee) {
        kafkaTemplate.executeInTransaction(operations -> {
            operations.send(TOPIC_NAME, employee.getId(), employee);
            return null;
        });

        AdminClient adminClient = AdminClient.create(kafkaTemplate.getProducerFactory().getConfigurationProperties());

        NewTopic newTopic = new NewTopic(TOPIC_NAME, NUM_PARTITIONS, REPLICATION_FACTOR);
        Map<String, String> config = new HashMap<>();
        config.put(TopicConfig.RETENTION_MS_CONFIG, String.valueOf(RETENTION_PERIOD_MS));
        newTopic.configs(config);

        try {
            adminClient.createTopics(Collections.singleton(newTopic)).all().get();
        } catch (InterruptedException | ExecutionException e) {
            throw new RuntimeException(e);
        } finally {
            adminClient.close();
        }
    }
}

 

В этой измененной версии KafkaProducerService для отправки записи Employee в топик Kafka «employee-topic» мы используем метод executeInTransaction , который гарантирует то, что продюсер Kafka отправит сообщение атомарно, то есть сообщение будет либо успешно отправлено, либо не отправлено вовсе.

Мы также используем класс AdminClient из клиентской библиотеки Kafka для настройки свойств топика Kafka. В этом примере мы устанавливаем количество разделов на 3, коэффициент репликации на 1 и период удержания на 1 день. Вы можете настроить эти свойства в соответствии с Вашими требованиями. Обратите внимание, что класс AdminClient должен быть закрыт после использования для того, чтобы освободить ресурсы, используемые клиентом Kafka.

 

Консюмер Kafka

Пример консюмера Kafka в Spring Boot, который читает сообщения из топика «employee-topic».:

@Service
public class KafkaConsumerService {

    private static final String TOPIC_NAME = "employee-topic";

    @KafkaListener(topics = TOPIC_NAME, groupId = "employee-consumer-group")
    public void receiveEmployee(Employee employee) {
        // do something with the received employee record
        System.out.println("Received employee: " + employee);
    }
}

 

В этом примере мы используем аннотацию @KafkaListener из библиотеки Spring Kafka для настройки метода, который будет вызываться каждый раз, когда будет получено новое сообщение из топика «employee-topic». Атрибут groupId указывает уникальный идентификатор группы консюмеров, к которой принадлежит данный потребитель. Обратите внимание, что несколько консюмеров могут принадлежать к одной группе, и Kafka будет автоматически распределять разделы топиками между ее членами.

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

Чтобы включить консюмера Kafka в Ваше приложение Spring Boot, Вам нужно настроить свойства Kafka в файле application.properties или application.yml. Пример конфигурации:

spring.kafka.consumer.bootstrap-servers=<kafka-broker-url>
spring.kafka.consumer.group-id=employee-consumer-group
spring.kafka.consumer.auto-offset-reset=earliest

 

В этой конфигурации мы устанавливаем свойство bootstrap servers на URL брокера (брокеров) Kafka, к которым должен подключаться консюмер, а также свойство group ID на то же значение, которое мы использовали в аннотации @KafkaListener. Свойство auto-offset-reset определяет поведение консюмера, когда он запускается и не имеет действительного смещения для раздела. В этом примере мы установим значение «earliest» для того, чтобы консюмер читал все сообщения в топике с самого начала.

 

 

Узнать стоимость решенияЗапросить видео презентацию

← Предыдущая статья
Kafka Connect: коннекторы, конфигурации, задачи, рабочие узлы
Следующая статья →
Подробное руководство по интеграции Kafka в приложение Spring Boot
Запросить видео презентацию Запросить доступ к демо стенду online Узнать стоимость лицензий

Задать вопрос

loading...

Решения

Анализировать ФинансыУвеличивайте ПродажиОптимальный Склад и ЛогистикаМаркетинговые Метрики

Клиенты
  • ПАО «Банк Уралсиб» (Публичное акционерное общество «Банк Уралсиб») — российский коммерческий банк. В 2020 году входил в топ-20 банков РФ по размеру активов (рэнкинг рейтингового агентства Эксперт РА), в 2021 году — в топ-25 крупнейших банков страны по расчётам агрегатора Банки.ру

  • НПФ «Будущее» — один из крупнейших негосударственных пенсионных фондов России, предоставляющий услуги по пенсионному обеспечению и накоплениям. Фонд активно внедряет цифровые технологии для повышения качества обслуживания клиентов.

  • КАМИ – компания-лидер по поставкам тяжёлых станков в России, занимающаяся продажей и обслуживанием оборудования для обработки металла и дерева, изготовления мебели и не только. На сегодняшний день в компании работают более 1300 человек, запущено 10 обучающих центров, в продаже более 7000 единиц техники. 

  • "Холодильник.ру" - крупнейший в России интернет-магазин бытовой техники и электроники. Компания была основана в 2003 году и за почти 20 лет работы завоевала лидирующие позиции на рынке онлайн ритейла. По данным исследовательского агентства Data Insight, "Холодильник.ру" входит в top-10 крупнейших интернет-магазинов России в категории "электроника и бытовая техника". Компания имеет развитую логистическую инфраструктуру и ежедневно осуществляет более 3500 доставок заказов по всей стране.

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Энергетика
    • Фармацевтика
  • Услуги
    • Переход на отечественные BI и DWH
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Техническая поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Платформы
    • FineBI
    • FineReport
    • FineDataLink
    • Коннекторы данных из 1С в BI
    • Airflow + NiFi
    • Visiology
    • Luxms BI
    • Modus BI
    • PIX BI
    • Arenadata
    • ClickHouse
    • Greenplum
    • Postgres Professional
    • Open-source BI: Superset/Metabase
    • Loginom
    • Yandex.DataLens
    • AI / Исскуственный интеллект
    • Optimacros
    • Шины данных
  • Курсы
    • Учебный курс Информационная грамотность
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt
  • Функциональные решения
    • Создание Data Lake
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и прогнозная аналитика
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • Сквозная аналитика
  • Компания
    • О нас
    • Руководство
    • Новости
    • Клиенты
    • Скачать
    • Контакты
    • Политика конфиденциальности
RutubeVkontakteLinkedInYouTube
ООО "Би Ай Консалт",
ИНН: 7811437757,
ОГРН: 1097847154184
199178, Россия,
Санкт-Петербург,
6-ая линия В.О., Д. 63, 4 этаж
Тел: +7 (812) 334-08-01
Тел: +7 (499) 608-13-06
E-mail: info@biconsult.ru

 

 

 

 

 

×

Пользуясь сайтом, вы соглашаетесь с использованием cookies и политикой конфиденциальности.