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 на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс Современная архитектура хранилища данных » Data Lake с помощью Debezium, Kafka Connect и Apache Iceberg sink

Data Lake с помощью Debezium, Kafka Connect и Apache Iceberg sink

Что такое Apache Iceberg?

Apache Iceberg - это открытый формат таблиц, предназначенный для огромных аналитических наборов данных, который можно использовать с многими мощнейшими движками обработки Big Data, такими как Apache Spark, Trino, PrestoDB, Flink и Hive. Эта технология может использоваться не только в пакетной обработке, но и быть отличным инструментом для сбора данных в режиме реального времени, которые поступают из действий пользователей, метрик, журналов, из CDC или других источников. Apache Iceberg предоставляет механизмы для изоляции чтения-записи и уплотнения данных из коробки, что позволяет избежать проблем при работе с маленькими файлами.

Также стоит отметить и то, что Apache Iceberg можно использовать с любым облачным провайдером или собственным решением, поддерживающим метахранилище Apache Hive и blob-хранилище.

 

Kafka Connect , Apache Iceberg sink

Компания GetInData создала Apache Iceberg sink, который можно развернуть на экземпляре Kafka Connect. Вы можете найти репозиторий и пакет на нашем GitHub.

Apache Iceberg sink был создан на базе memiiso/debezium-server-iceberg, который в свою очередь был создан для автономного использования наряду с Debezium Server.

Формат данных, используемый Apache Iceberg, должен представлять табличные данные и их схему, поэтому для сбора данных об изменениях мы использовали формат, созданный Debezium. 

 

 

Пример CDC

Давайте попробуем использовать Apache Kafka sink для репликации базы данных PostgreSQL с помощью Debezium для захвата всех изменений и передачи их в таблицу Apache Iceberg.

Мы запустим экземпляр Kafka Connect, на котором развернем источник Debezium,  а также наш Apache Iceberg sink. Для связи между ними будет использоваться топик Kafka, а sink будет записывать данные в бакет S3 и метаданные в Amazon Glue. Позже для чтения и отображения данных мы будем использовать Amazon Athena.

 

Шаг 1: запуск Kafka Connect

Сначала пройдите аутентификацию и сохраните учетные данные AWS в файле, например ~/.aws/confi

[default]
region = eu-west-1
aws_access_key_id=\*\**
aws_secret_access_key=\*\**

 

Загрузите sink отсюда. Например,~/Downloads/kafka-connect-iceberg-sink-0.1.3-shaded.jar

Для подключения Kafka мы будем использовать образ докера от Debezium, который поставляется с исходными пакетами Debezium. Мы смонтируем наш Apache Iceberg sink в каталог плагинов Kafka Connect и добавим файл с учетными данными AWS.

docker run -it --name connect --net=host -p 8083:8083 \
-e GROUP_ID=1 \
-e CONFIG_STORAGE_TOPIC=my-connect-configs \
-e OFFSET_STORAGE_TOPIC=my-connect-offsets \
-e BOOTSTRAP_SERVERS=localhost:9092 \
-e CONNECT_TOPIC_CREATION_ENABLE=true \
-v ~/.aws/config:/kafka/.aws/config \
-v ~/Downloads/kafka-connect-iceberg-sink-0.1.3-shaded.jar:/kafka/connect/kafka-connect-iceberg-sink-0.1.3-shaded.jar \
debezium/connect

 

Шаг 2: Чтение данных из источника PostgreSQL

Одна из возможностей Debezium по считыванию данных из PostgreSQL - это работа в качестве реплики базы данных. Чтобы Debezium работал корректно, нам необходимо увеличить объем информации, хранящейся в журнале опережающей записи. Для этого нам нужно настроить уровень wal_level на логический.

Запустим PostgreSQL на Docker:

docker run -d --name postgres -e POSTGRES_PASSWORD=postgres \
  -p 5432:5432 postgres -c wal_level=logical

 

Нам также понадобится экземпляр Kafka. 

Если функция автоматического создания топиков не включена, нужно создать топик, который будет использоваться для связи между источником Debezium и нашим Apache Iceberg sink. Имя топика состоит из логического имени, которое мы присвоим источнику Debezium, имени схемы базы данных и имени таблицы.

kafka-topics.sh --bootstrap-server localhost:9092 --create --topic postgres.public.dbz_test --partitions 1 --replication-factor 1

 

Теперь нам нужно развернуть источник на Kafka Connect. Мы можем сделать это с помощью POST-запроса, содержащего его конфигурацию.

curl -X POST -H "Content-Type: application/json" \
    -d '{
  "name": "postgres-connector", 
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "localhost",
    "database.port": "5432",
    "database.user": "postgres",
    "topic.prefix": "postgres",
    "database.password": "postgres",
    "database.dbname" : "postgres",
    "database.server.name": "postgres",
    "slot.name": "debezium",
    "plugin.name": "pgoutput",
    "table.include.list": "public.dbz_test",
    "transforms" : "unwrap",
    "transforms.unwrap.type":"io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.add.fields":"op,table,lsn,source.ts_ms,db",
    "transforms.unwrap.drop.tombstones":"true",
    "transforms.unwrap.delete.handling.mode":"rewrite",
    "drop.tombstones": "true"
  }
}' \
    <http://localhost:8083/connectors>

 

Шаг 3: Apache Iceberg Sink

Для нашего Apache Iceberg sink нам понадобится бакет  S3, например gid-streaminglabs-eu-west-1, и база данных в Amazon Glue, например gid_streaminglabs_eu_west_1_dbz.

Поскольку у нас уже есть готовый экземпляр Kafka Connect, включая учетные данные AWS, и пакет с sink, осталось только развернуть его. По аналогии с источником PostgreSQL мы сделаем это с помощью POST-запроса.

curl -X POST -H "Content-Type: application/json" \
    -d '{
  "name": "iceberg-sink",
  "config": {
    "connector.class": "com.getindata.kafka.connect.iceberg.sink.IcebergSink",
    "topics": "postgres.public.dbz_test",
    "upsert": true,
    "upsert.keep-deletes": true,
    "table.auto-create": true,
    "table.write-format": "parquet",
    "table.namespace": "gid_streaminglabs_eu_west_1_dbz",
    "table.prefix": "debeziumcdc_",
    "iceberg.catalog-impl": "org.apache.iceberg.aws.glue.GlueCatalog",
    "iceberg.warehouse": "s3a://gid-streaminglabs-eu-west-1/dbz_iceberg/gl_test",
    "iceberg.fs.defaultFS": "s3a://gid-streaminglabs-eu-west-1/dbz_iceberg/gl_test",
    "iceberg.com.amazonaws.services.s3.enableV4": true,
    "iceberg.com.amazonaws.services.s3a.enableV4": true,
    "iceberg.fs.s3a.aws.credentials.provider": "com.amazonaws.auth.DefaultAWSCredentialsProviderChain",
    "iceberg.fs.s3a.path.style.access": true,
    "iceberg.fs.s3a.impl": "org.apache.hadoop.fs.s3a.S3AFileSystem"
  }
}' \
    <http://localhost:8083/connectors>

 

Шаг 4: Проверка

Теперь мы можем открыть клиент psql и создать несколько таблиц:

psql -U postgres -h localhost
create table dbz_test (timestamp bigint, id int PRIMARY KEY, value int);
insert into dbz_test values(1, 1, 1);
insert into dbz_test values(2, 2, 2);
alter table dbz_test add test varchar(30);
insert into dbz_test values(3, 3, 3, 'aaa');
delete from dbz_test where id = 1;
update dbz_test set value = 1 where id = 2;

 

Затем перейдите в Amazon Athena и выполните следующий запрос:

select * from debeziumcdc_postgres_public_dbz_test order by timestamp desc;

 

Вот, собственно говоря, и все.

 

 

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

← Предыдущая статья
Оптимизация работы с базами данных: Моделирование данных с помощью Dbeaver
Следующая статья →
Apache Iceberg: Формат открытых таблиц для Data Lakehouse и потоковой передачи данных

Решения

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

Клиенты
  • «Балтийский лизинг» — первая компания в России, получившая лицензию № 0001 от Министерства экономики РФ на лизинговую деятельность, лицензия зарегистрирована 2 сентября 1996 года. «Балтийский лизинг» работает на российском рынке 33 года: компания представлена 79 филиалами по всей стране, сегодня в штате более 1300 сотрудников. За последние десять лет компания профинансировала имущество для 80 000 клиентов.

  • Группа компаний «Галакс» ведет свою деятельность с 2005 года, являясь в те годы дистрибьютором известных международных марок в ряде крупнейших торговых сетей России в сегменте аудио и видео аксессуаров. Активно работая в этом направлении и приобретая ценный опыт, начали создавать собственные торговые марки «GAL» и «VIXTER»

  • ЭГИС - международная фармацевтическая компания, основанная в 1907 году в Венгрии. Компания имеет представительства более чем в 60 странах мира, в том числе в России. Компания ЭГИС является одним из ведущих производителей дженерических лекарственных средств в Центральной и Восточной Европе. Её деятельность охватывает все звенья производственно-сбытовой фармацевтической цепочки.

  • Компания «Бизон-Трейд» является официальным дилером ведущих мировых производителей сельскохозяйственной техники (Fendt, Valtra, Lemken и др.) на Юге России. Входит в состав агрохолдинга «Бизон», основанного в 1994 году. Имеет 8 филиалов в Краснодарском и Ставропольском краях, Ростовской области.

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