pydbzengine
A Pythonic interface for the Debezium Engine, allowing you to consume database Change Data Capture (CDC) events directly in your Python applications.
Full Documentation: https://memiiso.github.io/pydbzengine
Features
- Pure Python Interface: Interact with the powerful Debezium Engine using simple Python classes and methods.
- Multi-Format Support: Unified engine supporting both JSON and Kafka Connect record formats (
jsonandconnect). - Pluggable Event Handlers: Easily create custom handlers to process CDC events according to your specific needs.
- Built-in Iceberg Handler: Stream change events directly into Apache Iceberg tables with zero boilerplate.
- Seamless Integration: Designed to work with popular Python data tools like dlt (data load tool).
- Apache Airflow Operator: Run Debezium engines directly within Airflow DAGs using the built-in
DebeziumEngineOperator(see Airflow docs). - Asynchronous & Snapshot Helpers: Run engines with time limits or terminate them automatically once initial snapshots complete using
Utils(see Helper docs). - All Debezium Connectors: Supports all standard Debezium connectors (PostgreSQL, MySQL, SQL Server, Oracle, etc.).
How it Works
This library acts as a bridge between the Python world and the Java-based Debezium Engine. It uses JPype to manage the JVM and interact with Debezium's Java classes, exposing a clean, Pythonic API so you can focus on your data logic without writing Java code.
Pre-available Data Handling Classes
pydbzengine comes with several built-in handlers. For detailed configuration and advanced usage, see the Handlers Documentation.
Apache Iceberg Handler
Stream CDC events directly into Apache Iceberg tables.
IcebergChangeHandlerV2(Recommended): Automatically infers schemas and manages table structures with native data types.IcebergChangeHandler: Appends raw change data (JSON) to source-equivalent tables using a fixed schema.
dlt (data load tool) Handler
DltChangeHandler: Integrates with thedltlibrary to load data into any supported destination (DuckDB, BigQuery, Snowflake, etc.).
Custom Handlers
BasePythonChangeHandler: Extend this class to implement your own custom processing logic in pure Python.
Installation
Prerequisites
You must have a Java Development Kit (JDK) version 17 or newer installed and available in your system's PATH.
Recommended Installation: From GitHub
You can install either the latest development version from the main branch or a specific, stable version from a release tag.
To install the latest development version:
# For core functionality
pip install "git+https://github.com/memiiso/pydbzengine.git"
# With extras (e.g., iceberg, dlt)
pip install "pydbzengine[iceberg] @ git+https://github.com/memiiso/pydbzengine.git"
pip install "pydbzengine[dlt] @ git+https://github.com/memiiso/pydbzengine.git"
pip install "pydbzengine[dev] @ git+https://github.com/memiiso/pydbzengine.git"
# To install a specific version from a release tag (e.g., 3.4.1.0):
pip install "pydbzengine @ git+https://github.com/memiiso/pydbzengine.git@3.4.1.0"
Alternative: From PyPI (Outdated Version)
An older version is available on PyPI. You can install it, but be aware that it lacks recent features and updates.
# For core functionality
pip install pydbzengine
# With extras
pip install "pydbzengine[iceberg]"
pip install "pydbzengine[dlt]"
How to Use
Consume events With custom Python consumer
- First install the packages:
pip install "pydbzengine[dev] @ git+https://github.com/memiiso/pydbzengine.git" - Second, extend
BasePythonChangeHandlerand implement your Python consuming logic. See the example below:
from typing import List
from pydbzengine import ChangeEvent, BasePythonChangeHandler, DebeziumEngine
class PrintChangeHandler(BasePythonChangeHandler):
"""
A custom change event handler class.
This class processes batches of Debezium change events received from the engine.
The `handleJsonBatch` method is where you implement your logic for consuming
and processing these events. Currently, it prints basic information about
each event to the console.
"""
def handleJsonBatch(self, records: List[ChangeEvent]):
"""
Handles a batch of Debezium change events.
This method is called by the Debezium engine with a list of ChangeEvent objects.
Change this method to implement your desired processing logic. For example,
you might parse the event data, transform it, and load it into a database or
other destination.
Args:
records: A list of ChangeEvent objects representing the changes captured by Debezium.
"""
print(f"Received {len(records)} records")
for record in records:
print(f"destination: {record.destination()}")
print(f"key: {record.key()}")
print(f"value: {record.value()}")
print("--------------------------------------")
if __name__ == "__main__":
props = {
"name": "engine",
"snapshot.mode": "initial_only",
# Add further Debezium connector configuration properties here. For example:
# "connector.class": "io.debezium.connector.mysql.MySqlConnector",
# "database.hostname": "your_database_host",
# "database.port": "3306",
}
# Create a DebeziumEngine instance (default format is "json", or specify format="connect").
engine = DebeziumEngine(properties=props, handler=PrintChangeHandler())
# Start the Debezium engine to begin consuming and processing change events.
engine.run()
Consume events to Apache Iceberg
from pyiceberg.catalog import load_catalog
from pydbzengine import DebeziumEngine
from pydbzengine.handlers.iceberg import IcebergChangeHandlerV2
conf = {
"uri": "http://localhost:8181",
# "s3.path-style.access": "true",
"warehouse": "warehouse",
"s3.endpoint": "http://localhost:9000",
"s3.access-key-id": "minioadmin",
"s3.secret-access-key": "minioadmin",
}
catalog = load_catalog(name="rest", **conf)
handler = IcebergChangeHandlerV2(
catalog=catalog,
destination_namespace=(
"iceberg",
"debezium_cdc_data",
),
)
dbz_props = {
"name": "engine",
"snapshot.mode": "always",
# ....
# Add further Debezium connector configuration properties here. For example:
# "connector.class": "io.debezium.connector.mysql.MySqlConnector",
}
engine = DebeziumEngine(properties=dbz_props, handler=handler)
engine.run()
Consume events with dlt
For the full code please see dlt_consuming.py
from pydbzengine import DebeziumEngine
from pydbzengine.helper import Utils
from pydbzengine.handlers.dlt import DltChangeHandler
import dlt
# Create a dlt pipeline and set destination. in this case DuckDb.
dlt_pipeline = dlt.pipeline(
pipeline_name="dbz_cdc_events_example",
destination="duckdb",
dataset_name="dbz_data",
)
handler = DltChangeHandler(dlt_pipeline=dlt_pipeline)
dbz_props = {
"name": "engine",
"snapshot.mode": "always",
# ....
}
engine = DebeziumEngine(properties=dbz_props, handler=handler)
# Run the Debezium engine asynchronously with a timeout.
# This runs for a limited time and then terminates automatically.
Utils.run_engine_async(engine=engine, timeout_sec=60)
Contributors
Metadata
Release files for pydbzengine 3.6.3.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| pydbzengine-3.6.3.0.tar.gz | 200.8 MB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| pydbzengine-3.6.3.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 401.8 MB
Release files / pydbzengine-3.6.3.0.tar.gz
| Download URL | pydbzengine-3.6.3.0.tar.gz |
|---|---|
| Size | 200.8 MB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
0ce5eacc9d4d041a5709df65128891e7c1942ac7a4e714c78119d0d225fb32e2
|
|
BLAKE2b-256 checksum How to use checksums |
0dedd5721cfdb636448451ef03432132f367da07f3485b38570522cf6cb50d4a
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Oct 4, 2026.
Transparency logRelease files / pydbzengine-3.6.3.0-py3-none-any.whl
| Download URL | pydbzengine-3.6.3.0-py3-none-any.whl |
|---|---|
| Size | 200.9 MB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
4c4c48d4f50c2ccc042fcb3b22698065d4c23366aab3b58af54c3e79dcf04480
|
|
BLAKE2b-256 checksum How to use checksums |
fe3341b0e258a37fd6a2370535092ba94381655ae549e38a60a737c9d28a0e58
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Oct 4, 2026.
Transparency log