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

Репликация данных в режиме реального времени с помощью Debezium и Python

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

Хотя Debezium сам по себе является проектом на базе Java, инженерии данных все же намного удобнее работать с Python. В этой статье мы поговорим о том, как можно использовать Debezium в среде Python с помощью движка pydbzengine. Мы также рассмотрим вопрос о том, как лучше всего использовать эти технологии для создания надежного и масштабируемого решения CDC.

 

Конвейер CDC на базе Python

CDC с помощью Debezium и загрузка данных в DuckDB с помощью DLT

В этой статье мы поговорим о том, как с помощью Debezium получить данные об изменениях из PostgreSQL и загрузить их в базу данных DuckDB с помощью DLT. Все это достигается с помощью пайплайна на Python, предполагающего использование pydbzengine.

Мы также пройдемся по коду, объясняя ключевые компоненты и то, как они работают вместе.

 

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

  1. Debezium: Мощная open-source платформа, предназначенная для сбора данных о любых изменениях. Она отслеживает журналы транзакций БД и создает потоки изменений, указывая на вставки, обновления и удаления;
  2. pydbzengine: Библиотека Python, обеспечивающая удобный способ взаимодействия со встроенным движком Debezium. В разы упрощает процесс настройки и запуска Debezium в приложениях на Python;
  3. DLT: Универсальный инструмент, предназначенный для загрузки данных, который упрощает процесс извлечения и загрузки данных в различные системы. В данной статье мы используем DLT для загрузки изменений из Debezium в DuckD;
  4. DuckDB: Встраиваемая аналитическая база данных, обеспечивающая эффективную обработку данных. Оптимальный вариант для локальной разработки и тестирования;
  5. Testcontainers: Библиотека, позволяющая запускать отдельные экземпляры сервисов, например баз данных. Мы используем ее для управления базой данных PostgreSQL.

 

Разбивка кода

Предоставленный код демонстрирует полный пайплайн CDC на Python (от захвата изменений с помощью Debezium до их загрузки в DuckDB с помощью DLT). Давайте разберем ключевые этапы данного процесса:

 

1. Установка окружения

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

import os
from pathlib import Path
import dlt
import duckdb
from testcontainers.core.config import testcontainers_config
from testcontainers.core.waiting_utils import wait_for_logs
from testcontainers.postgres import PostgresContainer
from pydbzengine import DebeziumJsonEngine, Properties
from pydbzengine.debeziumdlt import DltChangeHandler
from pydbzengine.helper import Utils

# set global variables
CURRENT_DIR = Path(__file__).parent
DUCKDB_FILE = CURRENT_DIR.joinpath("dbz_cdc_events_example.duckdb")
OFFSET_FILE = CURRENT_DIR.joinpath('postgresql-offsets.dat')

# cleanup
if OFFSET_FILE.exists():
    os.remove(OFFSET_FILE)
if DUCKDB_FILE.exists():
    os.remove(DUCKDB_FILE)

def wait_for_postgresql_to_start(self) -> None:
    wait_for_logs(self, ".*database system is ready to accept connections.*")
    wait_for_logs(self, ".*PostgreSQL init process complete.*")

class DbPostgresql:
    POSTGRES_USER = "postgres"
    POSTGRES_PASSWORD = "postgres"
    POSTGRES_DBNAME = "postgres"
    POSTGRES_IMAGE = "debezium/example-postgres:3.0.0.Final"
    POSTGRES_HOST = "localhost"
    POSTGRES_PORT_DEFAULT = 5432
    CONTAINER: PostgresContainer = (PostgresContainer(image=POSTGRES_IMAGE,
                                                      port=POSTGRES_PORT_DEFAULT,
                                                      username=POSTGRES_USER,
                                                      password=POSTGRES_PASSWORD,
                                                      dbname=POSTGRES_DBNAME,
                                                      )
                                    .with_exposed_ports(POSTGRES_PORT_DEFAULT)
                                    )
    PostgresContainer._connect = wait_for_postgresql_to_start

    def start(self):
        testcontainers_config.ryuk_disabled = True
        print("Starting Postgresql Db...")
        self.CONTAINER.start()

    def stop(self):
        print("Stopping Postgresql Db...")
        self.CONTAINER.stop()

    def __exit__(self, exc_type, exc_value, traceback):
        self.stop()

 

 

2. Настройка  Debezium

Мы настраиваем Debezium, создавая объект Java Properties. Этот объект содержит все настройки , необходимые для движка Debezium. Он включает в себя сведения о подключении к базе данных, класс коннектора (PostgresConnector), хранилище смещений и преобразования. Свойство transforms используется для разворачивания сообщения Debezium, что значительно упрощает последующую обработку данных.

def debezium_engine_props(sourcedb: DbPostgresql):
    props = Properties()
    props.setProperty("name", "engine")
    props.setProperty("snapshot.mode", "initial_only")
    props.setProperty("database.hostname", sourcedb.CONTAINER.get_container_host_ip())
    props.setProperty("database.port",
                      sourcedb.CONTAINER.get_exposed_port(sourcedb.POSTGRES_PORT_DEFAULT))
    props.setProperty("database.user", sourcedb.POSTGRES_USER)
    props.setProperty("database.password", sourcedb.POSTGRES_PASSWORD)
    props.setProperty("database.dbname", sourcedb.POSTGRES_DBNAME)
    props.setProperty("connector.class", "io.debezium.connector.postgresql.PostgresConnector")
    props.setProperty("offset.storage", "org.apache.kafka.connect.storage.FileOffsetBackingStore")
    props.setProperty("offset.storage.file.filename", OFFSET_FILE.as_posix())
    props.setProperty("max.batch.size", "5")
    props.setProperty("poll.interval.ms", "10000")
    props.setProperty("converter.schemas.enable", "false")
    props.setProperty("offset.flush.interval.ms", "1000")
    props.setProperty("database.server.name", "testc")
    props.setProperty("database.server.id", "1234")
    props.setProperty("topic.prefix", "testc")
    props.setProperty("schema.whitelist", "inventory")
    props.setProperty("database.whitelist", "inventory")
    props.setProperty("table.whitelist", "inventory.*")
    props.setProperty("replica.identity.autoset.values", "inventory.*:FULL")

    # // debezium unwrap message
    props.setProperty("transforms", "unwrap")
    props.setProperty("transforms.unwrap.type", "io.debezium.transforms.ExtractNewRecordState")
    props.setProperty("transforms.unwrap.add.fields", "op,table,source.ts_ms,sourcedb,ts_ms")
    props.setProperty("transforms.unwrap.delete.handling.mode", "rewrite")

    # props.setProperty("debezium.transforms.unwrap.drop.tombstones", "true")
    return props

 

3. Реализация обработчика изменений

Класс DltChangeHandler, предоставляемый библиотекой pydbzengine, действует как своеобразный мост между Debezium и DLT. Он получает события CDC от Debezium, а затем для эффективной загрузки этих данных в выбранное место назначения (например, в DuckDB) использует конвейер DLT. По сути, это тот самый компонент, который соединяет поток данных об изменениях в режиме реального времени из Debezium с возможностями загрузки данных в DLT. Полную реализацию данного процесса можно найти в репозитории pydbzengine.

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

from pydbzengine import BasePythonChangeHandler, ChangeEvent

class MyXYZChangeHandler(BasePythonChangeHandler):
    def handleJsonBatch(self, records: List[ChangeEvent]):

        # Process your data here!
        for record in records:

            # ... your processing logic ...
            # Example: send data to another service, database, etc.

 

4. Запуск движка Debezium и конвейера DLT

Главная функция организует весь процесс в целом. Она запускает контейнер PostgreSQL, создает движок Debezium с настроенными свойствами и запускает его.

def main():
    # Start the PostgreSQL container that will serve as the replication source.
    sourcedb = DbPostgresql()
    sourcedb.start()
 
    # Get Debezium engine configuration properties
    props = debezium_engine_props(sourcedb=sourcedb)
 
    # Create a dlt pipeline to consume the change events into DuckDB.
    dlt_pipeline = dlt.pipeline(
        pipeline_name="dbz_cdc_events_example",
        destination="duckdb",
        dataset_name="dbz_data"
    )
 
    handler = DltChangeHandler(dlt_pipeline=dlt_pipeline)
    engine = DebeziumJsonEngine(properties=props, handler=handler)
 
    # Run the Debezium engine asynchronously with a timeout.  This allows the example
    # to run for a limited time and then terminate automatically.
    Utils.run_engine_async(engine=engine, timeout_sec=60)
    # engine.run()  # This would be used for synchronous execution (without timeout)
 
if __name__ == "__main__":
    main()

 

5. Запрос к базе данных DuckDB, результат

После того как движок Debezium проработает заданное время (в нашем случае это 60 секунд), мы можем подключиться к базе данных назначения (DuckDB) и отобразить загруженные в нее данные:

    con = duckdb.connect(DUCKDB_FILE.as_posix())
    result = con.sql("SHOW ALL TABLES").fetchall()
    for r in result:
        database, schema, table = r[:3]  # Extract database, schema, and table names.
        if schema == "dbz_data":  # Only show data from the schema where Debezium loaded the data.
            print(f"Data in table {table}:")
            con.sql(f"select * from {database}.{schema}.{table} limit 5").show() # Di

 

Потребляемые данные:

┌────────────────────┬────────────────────────┬────────┬───────────────────────────────┬──────────────────────────────────────────────┐
│      load_id       │      schema_name       │ status │          inserted_at          │             schema_version_hash              │
│      varchar       │        varchar         │ int64  │   timestamp with time zone    │                   varchar                    │
├────────────────────┼────────────────────────┼────────┼───────────────────────────────┼──────────────────────────────────────────────┤
│ 1738405897.413279  │ debezium_source_events │      0 │ 2025-02-01 11:31:38.086127+01 │ Q5UNIOd7gJ6ljH5qfKKcO7yWwPvNESKW+mVXJmx9geg= │
│ 1738405898.176148  │ debezium_source_events │      0 │ 2025-02-01 11:31:39.381275+01 │ OyUXGP6PvFQuUTPnPdvESnsEqpFAxivJoP+l0G6l4+M= │
│ 1738405899.4865642 │ debezium_source_events │      0 │ 2025-02-01 11:31:39.704015+01 │ jqZNcnJXF/33Va2kRWgKOZF4RnZSVgYxMDhFep8+Jg8= │
│ 1738405899.775917  │ debezium_source_events │      0 │ 2025-02-01 11:31:39.952311+01 │ jqZNcnJXF/33Va2kRWgKOZF4RnZSVgYxMDhFep8+Jg8= │
│ 1738405900.0213661 │ debezium_source_events │      0 │ 2025-02-01 11:31:40.223125+01 │ uMZY5n2NGPecXvVQIePLEg2nZQcAlkoWAXDLALKjWuQ= │
└────────────────────┴────────────────────────┴────────┴───────────────────────────────┴──────────────────────────────────────────────┘

 

Data in table _dlt_pipeline_state:

┌─────────┬────────────────┬──────────────────────┬──────────────────────┬──────────────────────┬────────────────────────────────────────────┬───────────────────┬────────────────┐
│ version │ engine_version │    pipeline_name     │        state         │      created_at      │                version_hash                │   _dlt_load_id    │    _dlt_id     │
│  int64  │     int64      │       varchar        │       varchar        │ timestamp with tim…  │                  varchar                   │      varchar      │    varchar     │
├─────────┼────────────────┼──────────────────────┼──────────────────────┼──────────────────────┼────────────────────────────────────────────┼───────────────────┼────────────────┤
│       1 │              4 │ dbz_cdc_events_exa…  │ eNp1j0FLw0AQhf/LXg…  │ 2025-02-01 11:31:3…  │ ZvlGi9hyfXjD2b0imkL9ZA7x3S1/YkmQK4QbA+Jw…  │ 1738405897.413279 │ hNbs3TIc3vRHvA │
└─────────┴────────────────┴──────────────────────┴──────────────────────┴──────────────────────┴────────────────────────────────────────────┴───────────────────┴────────────────┘

 

Data in table _dlt_version:

┌─────────┬────────────────┬──────────────────────┬──────────────────────┬──────────────────────┬─────────────────────────────────────────────────────────────────────────────────┐
│ version │ engine_version │     inserted_at      │     schema_name      │     version_hash     │                                     schema                                      │
│  int64  │     int64      │ timestamp with tim…  │       varchar        │       varchar        │                                     varchar                                     │
├─────────┼────────────────┼──────────────────────┼──────────────────────┼──────────────────────┼─────────────────────────────────────────────────────────────────────────────────┤
│       2 │             11 │ 2025-02-01 11:31:3…  │ debezium_source_ev…  │ Q5UNIOd7gJ6ljH5qfK…  │ {"version":2,"version_hash":"Q5UNIOd7gJ6ljH5qfKKcO7yWwPvNESKW+mVXJmx9geg=","e…  │
│       4 │             11 │ 2025-02-01 11:31:3…  │ debezium_source_ev…  │ OyUXGP6PvFQuUTPnPd…  │ {"version":4,"version_hash":"OyUXGP6PvFQuUTPnPdvESnsEqpFAxivJoP+l0G6l4+M=","e…  │
│       6 │             11 │ 2025-02-01 11:31:3…  │ debezium_source_ev…  │ jqZNcnJXF/33Va2kRW…  │ {"version":6,"version_hash":"jqZNcnJXF/33Va2kRWgKOZF4RnZSVgYxMDhFep8+Jg8=","e…  │
│       8 │             11 │ 2025-02-01 11:31:4…  │ debezium_source_ev…  │ uMZY5n2NGPecXvVQIe…  │ {"version":8,"version_hash":"uMZY5n2NGPecXvVQIePLEg2nZQcAlkoWAXDLALKjWuQ=","e…  │
└─────────┴────────────────┴──────────────────────┴──────────────────────┴──────────────────────┴─────────────────────────────────────────────────────────────────────────────────┘

 

Data in table testc_inventory_customers:

┌───────┬────────────┬───────────┬───────────────────────┬─────────┬─────────┬───────────┬───────────────┬───────────────┬───────────────────┬────────────────┐
│  id   │ first_name │ last_name │         email         │ deleted │   op    │   table   │ source_ts_ms  │     ts_ms     │   _dlt_load_id    │    _dlt_id     │
│ int64 │  varchar   │  varchar  │        varchar        │ varchar │ varchar │  varchar  │     int64     │     int64     │      varchar      │    varchar     │
├───────┼────────────┼───────────┼───────────────────────┼─────────┼─────────┼───────────┼───────────────┼───────────────┼───────────────────┼────────────────┤
│  1001 │ Sally      │ Thomas    │ sally.thomas@acme.com │ false   │ r       │ customers │ 1738405883186 │ 1738405896858 │ 1738405897.413279 │ KcWKrYODYJ859w │
│  1002 │ George     │ Bailey    │ gbailey@foobar.com    │ false   │ r       │ customers │ 1738405883186 │ 1738405896862 │ 1738405897.413279 │ JU6dR1S27Xt3QA │
│  1003 │ Edward     │ Walker    │ ed@walker.com         │ false   │ r       │ customers │ 1738405883186 │ 1738405896862 │ 1738405897.413279 │ 02kMVvIX2/aGGg │
│  1004 │ Anne       │ Kretchmar │ annek@noanswer.org    │ false   │ r       │ customers │ 1738405883186 │ 1738405896862 │ 1738405897.413279 │ TI7jpxl9FD2kRQ │
└───────┴────────────┴───────────┴───────────────────────┴─────────┴─────────┴───────────┴───────────────┴───────────────┴───────────────────┴────────────────┘

 

Data in table testc_inventory_geom:

┌───────┬──────────────────────────────────────────────────────────────────────┬─────────┬─────────┬─────────┬───────────────┬───────────────┬───────────────────┬────────────────┐
│  id   │                                g__wkb                                │ deleted │   op    │  table  │ source_ts_ms  │     ts_ms     │   _dlt_load_id    │    _dlt_id     │
│ int64 │                               varchar                                │ varchar │ varchar │ varchar │     int64     │     int64     │      varchar      │    varchar     │
├───────┼──────────────────────────────────────────────────────────────────────┼─────────┼─────────┼─────────┼───────────────┼───────────────┼───────────────────┼────────────────┤
│     1 │ AQEAAAAAAAAAAADwPwAAAAAAAPA/                                         │ false   │ r       │ geom    │ 1738405883186 │ 1738405896872 │ 1738405897.413279 │ 17snqevSVWL0xA │
│     2 │ AQIAAAACAAAAAAAAAAAAAEAAAAAAAADwPwAAAAAAABhAAAAAAAAAGEA=             │ false   │ r       │ geom    │ 1738405883186 │ 1738405896872 │ 1738405898.176148 │ W4kfG5n5jYhy3w │
│     3 │ AQMAAAABAAAABQAAAAAAAAAAAAAAAAAAAAAAFEAAAAAAAAAAQAAAAAAAABRAAAAAAA…  │ false   │ r       │ geom    │ 1738405883186 │ 1738405896872 │ 1738405898.176148 │ 40HrbnruXZaB/g │
└───────┴──────────────────────────────────────────────────────────────────────┴─────────┴─────────┴─────────┴───────────────┴───────────────┴───────────────────┴────────────────┘

 

Data in table testc_inventory_orders:

┌───────┬────────────┬───────────┬──────────┬────────────┬─────────┬─────────┬─────────┬───────────────┬───────────────┬────────────────────┬────────────────┐
│  id   │ order_date │ purchaser │ quantity │ product_id │ deleted │   op    │  table  │ source_ts_ms  │     ts_ms     │    _dlt_load_id    │    _dlt_id     │
│ int64 │   int64    │   int64   │  int64   │   int64    │ varchar │ varchar │ varchar │     int64     │     int64     │      varchar       │    varchar     │
├───────┼────────────┼───────────┼──────────┼────────────┼─────────┼─────────┼─────────┼───────────────┼───────────────┼────────────────────┼────────────────┤
│ 10001 │      16816 │      1001 │        1 │        102 │ false   │ r       │ orders  │ 1738405883186 │ 1738405896876 │ 1738405898.176148  │ X7ejebZDxmm+hw │
│ 10002 │      16817 │      1002 │        2 │        105 │ false   │ r       │ orders  │ 1738405883186 │ 1738405896876 │ 1738405898.176148  │ 6LU0Fe9UVE3XFQ │
│ 10003 │      16850 │      1002 │        2 │        106 │ false   │ r       │ orders  │ 1738405883186 │ 1738405896876 │ 1738405898.176148  │ 0OIBPdMzqjLh0w │
│ 10004 │      16852 │      1003 │        1 │        107 │ false   │ r       │ orders  │ 1738405883186 │ 1738405896876 │ 1738405899.4865642 │ CcY6FKlHLQ6mPg │
└───────┴────────────┴───────────┴──────────┴────────────┴─────────┴─────────┴─────────┴───────────────┴───────────────┴────────────────────┴────────────────┘

 

Data in table testc_inventory_products:

┌───────┬────────────────────┬──────────────────────────────────────┬────────┬─────────┬─────────┬──────────┬───────────────┬───────────────┬────────────────────┬────────────────┐
│  id   │        name        │             description              │ weight │ deleted │   op    │  table   │ source_ts_ms  │     ts_ms     │    _dlt_load_id    │    _dlt_id     │
│ int64 │      varchar       │               varchar                │ double │ varchar │ varchar │ varchar  │     int64     │     int64     │      varchar       │    varchar     │
├───────┼────────────────────┼──────────────────────────────────────┼────────┼─────────┼─────────┼──────────┼───────────────┼───────────────┼────────────────────┼────────────────┤
│   101 │ scooter            │ Small 2-wheel scooter                │   3.14 │ false   │ r       │ products │ 1738405883186 │ 1738405896879 │ 1738405899.4865642 │ aOf6efrtt48+1Q │
│   102 │ car battery        │ 12V car battery                      │    8.1 │ false   │ r       │ products │ 1738405883186 │ 1738405896880 │ 1738405899.4865642 │ kUuPhtKUAsTUaA │
│   103 │ 12-pack drill bits │ 12-pack of drill bits with sizes r…  │    0.8 │ false   │ r       │ products │ 1738405883186 │ 1738405896880 │ 1738405899.4865642 │ evSpPy68nldtbg │
│   104 │ hammer             │ 12oz carpenter's hammer              │   0.75 │ false   │ r       │ products │ 1738405883186 │ 1738405896880 │ 1738405899.4865642 │ lCpa9yyHSm8xqA │
│   105 │ hammer             │ 14oz carpenter's hammer              │  0.875 │ false   │ r       │ products │ 1738405883186 │ 1738405896880 │ 1738405899.775917  │ VXcDU/tw/zT2fw │
└───────┴────────────────────┴──────────────────────────────────────┴────────┴─────────┴─────────┴──────────┴───────────────┴───────────────┴────────────────────┴────────────────┘

 

Data in table testc_inventory_products_on_hand:

┌────────────┬──────────┬─────────┬─────────┬──────────────────┬───────────────┬───────────────┬────────────────────┬────────────────┐
│ product_id │ quantity │ deleted │   op    │      table       │ source_ts_ms  │     ts_ms     │    _dlt_load_id    │    _dlt_id     │
│   int64    │  int64   │ varchar │ varchar │     varchar      │     int64     │     int64     │      varchar       │    varchar     │
├────────────┼──────────┼─────────┼─────────┼──────────────────┼───────────────┼───────────────┼────────────────────┼────────────────┤
│        101 │        3 │ false   │ r       │ products_on_hand │ 1738405883186 │ 1738405896883 │ 1738405900.0213661 │ 90r3+XR7PH7y6g │
│        102 │        8 │ false   │ r       │ products_on_hand │ 1738405883186 │ 1738405896883 │ 1738405900.0213661 │ 5F+LUMVYO3I2wQ │
│        103 │       18 │ false   │ r       │ products_on_hand │ 1738405883186 │ 1738405896883 │ 1738405900.0213661 │ SguX65iX7ffyJg │
│        104 │        4 │ false   │ r       │ products_on_hand │ 1738405883186 │ 1738405896883 │ 1738405900.0213661 │ Vj/N2j0bN3ipzw │
│        105 │        5 │ false   │ r       │ products_on_hand │ 1738405883186 │ 1738405896883 │ 1738405900.0213661 │ z31M4RIQPpq3BA │
└────────────┴──────────┴─────────┴─────────┴──────────────────┴───────────────┴───────────────┴────────────────────┴────────────────┘

 

Проверьте самостоятельно.

Чтобы запустить этот процесс, Вам потребуется Docker Desktop и ключевые библиотеки Python. Вы можете установить зависимости, используя следующий код:

pip install pydbzengine[dev]
python dlt_consuming.py

 

Ключевые выводы

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

Использование Python и pydbzengine позволяет легко интегрировать Debezium в существующие рабочие процессы, осуществляемые на Python.

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

 

Подведение итогов

С помощью Debezium pydbzengine позволяет максимально просто организовать пайплайн с низкой задержкой. Это полностью open-source проект, использующий лицензию Apache 2.0. pydbzengine – проект «зеленый», в котором еще много чего предстоит подтянуть. Тестируйте его, присылайте свои комментарии и вопросы, будем улучшать данный проект вместе!

 

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

← Предыдущая статья
CI/CD и инциденты в аналитике
Следующая статья →
Современный стек данных слишком сложен … и 70% лидеров и практиков в области данных с этим полностью согласны!

Решения

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

Клиенты
  • Компания "Норникель" - лидер горно-металлургической отрасли в России и мире. Она производит металлы, необходимые для развития экологичной экономики и транспорта.

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

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

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

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