Репликация данных в режиме реального времени с помощью 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.
Мы также пройдемся по коду, объясняя ключевые компоненты и то, как они работают вместе.
Основные компоненты
- Debezium: Мощная open-source платформа, предназначенная для сбора данных о любых изменениях. Она отслеживает журналы транзакций БД и создает потоки изменений, указывая на вставки, обновления и удаления;
- pydbzengine: Библиотека Python, обеспечивающая удобный способ взаимодействия со встроенным движком Debezium. В разы упрощает процесс настройки и запуска Debezium в приложениях на Python;
- DLT: Универсальный инструмент, предназначенный для загрузки данных, который упрощает процесс извлечения и загрузки данных в различные системы. В данной статье мы используем DLT для загрузки изменений из Debezium в DuckD;
- DuckDB: Встраиваемая аналитическая база данных, обеспечивающая эффективную обработку данных. Оптимальный вариант для локальной разработки и тестирования;
- 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 messageprops.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") returnprops
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 – проект «зеленый», в котором еще много чего предстоит подтянуть. Тестируйте его, присылайте свои комментарии и вопросы, будем улучшать данный проект вместе!




