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 на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс по ClickHouse » Руководство по работе с движком Kafka в ClickHouse

Руководство по работе с движком Kafka в ClickHouse

Что делать, если Вы – новичок, и Вам нужна помощь в настройке Kafka и ClickHouse, так как Вы делаете это впервые? Тогда эта статья – именно для Вас!

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

 

Подготовка к работе

Приступая к изучению материала, мы предполагаем, что у Вас уже установлены и стабильно работают Kafka и ClickHouse. Кроме того, для удобства рекомендуем использовать Kubernetes. Наиболее подходящая версия Kafka - Confluent 5.4.0, установленная с помощью Kafka helm chart с тремя брокерами Kafka. Рнкомендуемая версия ClickHouse - 20.4.2, установленная на одном узле с помощью ClickHouse Kubernetes Operator.  Руководство по установке Confluent Kafka без использования Kubernetes доступно по ссылке, руководство по  ClickHouse без Kubernetes Вы найдете здесь. 

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

 

 

Обзор интеграции Kafka-ClickHouse

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

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

Следующая диаграмма иллюстрирует описанную выше схему работы Kafka.

 

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

 

Также есть возможность писать сообщения из ClickHouse обратно в Kafka. Поток сообщений в этом случае гораздо проще - просто вставьте их в таблицу Kafka.

 

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

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

kafka-topics \
--bootstrap-server kafka:9092 \
--topic readings \
--create --partitions 6 \
--replication-factor 2

 

Убедитесь в том, что топик был успешно создан. 

kafka-topics --bootstrap-server kafka:9092 --describe readings

 

Если Вы все сделали правильно, то увидите следующее:

Topic: readings    PartitionCount: 6    ReplicationFactor: 2    Configs:
    Topic: readings    Partition: 0    Leader: 0    Replicas: 0,2    Isr: 0,2
    Topic: readings    Partition: 1    Leader: 2    Replicas: 2,1    Isr: 2,1
    Topic: readings    Partition: 2    Leader: 1    Replicas: 1,0    Isr: 1,0
    Topic: readings    Partition: 3    Leader: 0    Replicas: 0,1    Isr: 0,1
    Topic: readings    Partition: 4    Leader: 2    Replicas: 2,0    Isr: 2,0
    Topic: readings    Partition: 5    Leader: 1    Replicas: 1,2    Isr: 1,2

 

На этом этапе мы готовы к дальнейшей работе со стороны Kafka. Теперь давайте переключим свое внимание на  ClickHouse.

 

Установка движка Kafka

Для того, чтобы в ClickHouse прочесть сообщения  из топика Kafka, необходимо иметь 3 вещи:

  1. Целевую таблицу MergeTree, которая примет данные;
  2. Таблицу движка Kafka, которая необходима для того, чтобы топик выглядел точно также  как таблица ClickHouse;
  3. Материализованное представление для автоматического перемещения данных из Kafka в целевую таблицу.

 

Давайте рассмотрим каждый пункт по порядку.

Сначала определим целевую таблицу MergeTree. Для этого заходим в ClickHouse и выполняем следующий SQL-запрос с целью создать таблицу:

CREATE TABLE readings (
    readings_id Int32 Codec(DoubleDelta, LZ4),
    time DateTime Codec(DoubleDelta, LZ4),
    date ALIAS toDate(time),
    temperature Decimal(5,2) Codec(T64, LZ4)
) Engine = MergeTree
PARTITION BY toYYYYMM(time)
ORDER BY (readings_id, time);

 

Далее нам нужно создать таблицу, использующую движок Kafka для подключения к топику и чтения данных.  Движок будет считывать данные с брокера на хосте kafka, используя топик'readings' и имя группы консюмеров  'readings consumer_group1'. Формат входных данных -CSV.  Обратите внимание на то, что мы опускаем колонку 'date', поскольку она будет автоматически заполняться из столбца 'time'.

CREATE TABLE readings_queue (
    readings_id Int32,
    time DateTime,
    temperature Decimal(5,2)
)
ENGINE = Kafka
SETTINGS kafka_broker_list = 'kafka-headless.kafka:9092',
         kafka_topic_list = 'readings',
         kafka_group_name = 'readings_consumer_group1',
         kafka_format = 'CSV',
         kafka_max_block_size = 1048576;

 

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

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

CREATE MATERIALIZED VIEW readings_queue_mv TO readings AS
SELECT readings_id, time, temperature
FROM readings_queue;

 

Теперь интеграция Kafka с ClickHouse завершена. Давайте протестируем, работает ли она.

 

Загрузка данных

Теперь пришло время загрузить некоторые входные данные с помощью команды kafka-console-producer. Вот пример, который загружает три записи в формате CSV.

kafka-console-producer --broker-list kafka:9092 --topic readings <<END
1,"2020-05-16 23:55:44",14.2
2,"2020-05-16 23:55:45",20.1
3,"2020-05-16 23:55:51",12.9
END

 

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

SELECT *
FROM readings

┌─readings_id─┬────────────────time─┬─temperature─┐
│           1 │ 2020-05-16 23:55:44 │       14.20 │
│           2 │ 2020-05-16 23:55:45 │       20.10 │
│           3 │ 2020-05-16 23:55:51 │       12.90 │
└─────────────┴─────────────────────┴─────────────┘

Отлично! Теперь Kafka и ClickHouse полностью интегрированы. 

 

Чтение сообщений из Kafka

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

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

TRUNCATE TABLE readings;

 

Перед сбросом смещений в разделах необходимо отключить опцию потребления сообщений. Для этого отсоединяем таблицу readings_queue в ClickHouse следующим образом.

DETACH TABLE readings_queue

 

Затем выполним следующую команду Kafka, которая позволит нам сбросить смещения разделов в группе консюмеров, используемой для таблицы readings_queue. (NB! Это не SQL-запрос. Команда должна выполняться  в Kafka, а не в ClickHouse)

kafka-consumer-groups --bootstrap-server kafka:9092 \
 --topic readings --group readings_consumer_group1 \
 --reset-offsets --to-earliest --execute

 

Теперь снова подключите таблицу readings_queue. Вот Вы и вернулись в ClickHouse. 

ATTACH TABLE readings_queue

 

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

 

Добавление виртуальных столбцов

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

Сначала отключим потребление сообщений, отсоединив таблицу Kafka.

DETACH TABLE readings_queue

 

Далее изменим целевую таблицу и материализованное представление с помощью следующих команд SQL. 

ALTER TABLE readings
  ADD COLUMN _topic String,
  ADD COLUMN _offset UInt64,
  ADD COLUMN _partition UInt64

DROP TABLE readings_queue_mv

CREATE MATERIALIZED VIEW readings_queue_mv TO readings AS
  SELECT readings_id, time, temperature, _topic, _offset, _partition
  FROM readings_queue;

 

Наконец, мы снова включаем потребление сообщений, повторно подключив таблицу readings_queue. 

ATTACH TABLE readings_queue

 

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

SELECT
    readings_id AS id, time, temperature AS temp,
    _topic, _offset, _partition
FROM readings

┌─id─┬────────────────time─┬──temp─┬─_topic───┬─_offset─┬─_partition─┐
│  1 │ 2020-05-16 23:55:44 │ 14.20 │ readings │       0 │          5 │
│  2 │ 2020-05-16 23:55:45 │ 20.10 │ readings │       1 │          5 │
│  3 │ 2020-05-16 23:55:51 │ 12.90 │ readings │       2 │          5 │
└────┴─────────────────────┴───────┴──────────┴─────────┴────────────┘

 

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

 

Запись из ClickHouse в Kafka

Чуть позже мы покажем Вам, как записывать сообщения из ClickHouse обратно в Kafka.  Это относительно новая функция, которая доступна начиная со сборки Altinity 19.16.18.85. 

Начнем с создания нового топика в Kafka, который будет содержать сообщения. Назовем его 'readings_high' (почему именно так, станет понятно чуть позже). 

kafka-topics \
--bootstrap-server kafka:9092 \
--topic readings_high \
--create --partitions 6 \
--replication-factor 2

 

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

CREATE TABLE readings_high_queue (
    readings_id Int32,
    time DateTime,
    temperature Decimal(5,2)
)
ENGINE = Kafka
SETTINGS kafka_broker_list = 'kafka:9092',
         kafka_topic_list = 'readings_high',
         kafka_group_name = 'readings_high_consumer_group1',
         kafka_format = 'CSV',
         kafka_max_block_size = 1048576;

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

CREATE MATERIALIZED VIEW readings_high_queue_mv TO readings_high_queue AS
SELECT readings_id, time, temperature FROM readings
WHERE toFloat32(temperature) >= 20.0

 

В отдельном окне терминала для распечатки сообщений из топика readings_high на Kafka запустите консюмера.  Это позволит Вам видеть все строки по мере того, как ClickHouse записывает их в Kafka.

kafka-console-consumer --bootstrap-server kafka:9092 --topic readings_high

 

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

kafka-console-producer --broker-list kafka:9092 --topic readings <<END
4,"2020-05-16 23:55:52",9.7
5,"2020-05-16 23:55:56",25.3
6,"2020-05-16 23:55:58",14.1
END

 

Через несколько секунд Вы увидите, как в окне, запущенном командой kafka-console-consumer, появится вторая строка. Она должна выглядеть следующим образом:

5,"2020-05-16 23:55:56",25.3

 

Что делать с ошибками

Если в процессе работы у Вас возникнут какие-либо проблемы, просмотрите журнал ClickHouse. Включите ведение журнала (или по-другому трассировку), если Вы этого еще не сделали. В журнале Вы увидите сообщения, подобные тому, что приведен ниже.

2020.05.17 07:24:20.609147 [ 64 ] {} <Debug> StorageKafka (readings_queue): Started streaming to 1 attached views

 

Любые ошибки, если таковые имеются, будут отображаться в журнале clickhouse-server.err.log.

 

Заключение

Итак, в этой статье мы поговорили о том, что движок  Kafka обеспечивает простой и достаточно мощный способ интеграции  Kafka и ClickHouse. Очевидно, что управление интеграцией - особенно в производственной системе – процесс достаточно сложный. Мы искренне надеемся на то, что данная статья поможет Вам начать работу с двумя популярнейшими решениями в области работы с данными.

 

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

← Предыдущая статья
Углубляемся в тонкости Apache Parquet с помощью ClickHouse - часть 2
Следующая статья →
Clickhouse — шардинг и репликация
Запросить видео презентацию Запросить доступ к демо стенду online Узнать стоимость лицензий

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

loading...

Решения

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

Клиенты
  • ПАО «Ростелеком» — российский провайдер цифровых услуг и сервисов. Предоставляет услуги широкополосного доступа в Интернет, интерактивного телевидения, сотовой связи, местной и дальней телефонной связи и др. Занимает лидирующие позиции на российском рынке высокоскоростного доступа в интернет, платного ТВ, хранения и обработки данных, а также кибербезопасности

  • ПАО «Транснефть» – крупнейшая российская нефтепроводная компания. «Транснефть» обеспечивает транспортировку более 85% добываемых в России нефти и нефтепродуктов.

  • Группа компаний «Невский кондитер» основана в 1996 году в Санкт-Петербурге и на сегодняшний день является одним из крупнейших производителей кондитерских изделий в России.

     

  • Novikov group – первый российский ресторанный холдинг, основанный в 1991 году. Это команда профессионалов под управлением Аркадия Новикова, реализующая широкий спектр услуг в сфере гостеприимства: от проведения event-мероприятия до управления рестораном, от установления стандартов сервиса до контроля качества готовой продукции, от построения бизнес-плана проекта до реализации франшизы.

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • 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 и политикой конфиденциальности.