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

Потоковая передача данных Postgres с помощью Apache Kafka и Debezium | ETL в режиме реального времени

Используем Apache Kafka, Debezium и Postgres правильно.

 

Сегодня мы продолжим говорить о потоковой передачи данных  Postgres в Apache Kafka. Ранее мы установили среду и применили все настройки, необходимые для потоковой передачи данных из Postgres. Мы также установили Apache Kafka, Debezium и все остальные важные компоненты, а также настроили  БД Postgres на потоковую передачу данных. В первую очередь нам нужна таблица, из которой мы будем брать данные. Во время Python ETL сессии мы загрузили несколько таблиц. Для организации потоковой передачи данных предлагаю воспользоваться таблицей FactInternetSales , содержащей транзакции по продажам. Представим, что информация о транзакции продажи приходит каждые несколько секунд, и нам надо сразу же передать ее в топик Kafka.

  • Из этой статьи Вы узнаете о том, как:
  • настроить Postgres на потоковую передачу данных
  • передавать изменения, произошедшие в БД, в топик Kafka
  • настроить коннектор потоковой передачи данных Postgres - Kafka
  • организовать потоковую передачи данных консюмера через Python Kafka Consumer

 

Если Вам удобнее работать с видео-материалом, тогда предлагаю Вашему вниманию подробное видео на YouTube .

 

 

Таблица для потоковой передачи данных

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

CREATE TABLE IF NOT EXISTS public.factinternetsales_streaming
(
    productkey bigint,
    customerkey bigint,
    salesterritorykey bigint,
    salesordernumber text COLLATE pg_catalog."default",
    totalproductcost double precision,
    salesamount double precision
)

TABLESPACE pg_default;

ALTER TABLE IF EXISTS public.factinternetsales_streaming
    OWNER to postgres;

 

Инкрементальная загрузка данных

Мы должны вставлятьстроки в таблицу по одной за раз. Чтобы не делать этого вручную, примените специальный скрипт, как это сделал я. Это будет чем-то похоже на Python ETL серии. Когда мы выполним этот скрипт, он вставит строки так, как если бы это делало приложение, записывающее данные в БД. Так мы настроим нашу таблицу на потоковую передачу данных.

engine = create_engine(f'postgresql://{uid}:{pwd}@{server}:{port}/{db}')
df = pd.read_sql('Select * from public.factinternetsales', engine)
df = df[['productkey', 'customerkey', 'salesterritorykey', 'salesordernumber', 'totalproductcost', 'salesamount']]
#
for index, row in df.head(100).iterrows():
    mod = pd.DataFrame(row.to_frame().T)     mod.to_sql(f"factinternetsales_streaming", engine, if_exists='append', index=False)
    print("Row Inserted " + mod.salesordernumber.astype(str) + ' ' + mod.salesamount.astype(str).astype(str))
    time.sleep(3)

 

Коннектор Kafka Postgres

Переходим к настройке Kafka. Для организации потоковой передачи данных из Postgres мы можем использовать Debezium (при этом никакого дополнительного кода нам не нужно). Все, что нам нужно сделать – это настроить коннектор. Конечная точка нашего Debezium API: localhost:8083. Настройка коннектора выглядит следующим образом:

{
    "name": "source-productcategory-connector",
    "config": {
        "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
        "database.hostname": "your-host-ip",
        "database.port": "5432",
        "database.user": "user",
        "database.password": "password",
        "database.dbname": "AdventureWorks",
        "plugin.name": "pgoutput",
        "database.server.name": "source",
        "key.converter.schemas.enable": "false",
        "value.converter.schemas.enable": "false",
        "transforms": "unwrap",
        "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
        "value.converter": "org.apache.kafka.connect.json.JsonConverter",
        "key.converter": "org.apache.kafka.connect.json.JsonConverter",
        "table.include.list": "public.factinternetsales_streaming",
        "slot.name" : "dbz_sales_transaction_slot"
    } }

 

Kafka Python Consumer

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

 

Если мы вернемся к нашему консюмеру, то увидим, что изменения в базе данных успешно передаются в топик. Наш консюмер получает эти изменения и сразу же отображает их. Мы успешно транслируем изменения базы данных в топик Kafka, и одновременно наш консюмер Python читает этот поток из топикаKafka. Таким образом, мы успешно транслируем данные из Postgres в Kafka с помощью Debezium.

 

Заключение

  • Мы продемонстрировали, как настроить базу данных Postgres для потоковой передачи данных в режиме реального времени.
  • Для потоковой передачи изменений базы данных в топик Kafka мы создали коннектор Kafka.
  • Для инкрементной вставки данных в базу мы создали специальный Python-скрипт.
  • Для чтения потока данных из базы данных Postgres в режиме реального времени мы создали Python Consumer.
  • Полный код Вы найдете здесь.

 

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

← Предыдущая статья
Что такое Debezium и как его применять
Следующая статья →
Планирование миллионов сообщений с помощью Kafka и Debezium

Решения

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

Клиенты
  • ГК «Агропромкомплектация-Курск» - одна из ведущих в Российской Федерации агропромышленных компаний с полным производственным циклом "от поля до прилавка". За 32 года работы на рынке компания заслуженно завоевала репутацию одного из лидеров страны в производстве свинины и молока.

  • ООО "Уральская транспортная компания" — это транспортно-логистическая компания, специализирующаяся на железнодорожных перевозках грузов, создана в 2009 году.

  • "Холодильник.ру" - крупнейший в России интернет-магазин бытовой техники и электроники. Компания была основана в 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 и политикой конфиденциальности.