Skip to main content

Python client helpers for the ducklake-cdc DuckDB extension

Project description

ducklake-cdc-client

Python client helpers for the ducklake_cdc DuckDB community extension.

The package gives you two layers:

  • CDCClient: a direct Python wrapper over the extension's SQL table functions.
  • DMLConsumer / DDLConsumer: durable consumers that yield batches and commit only after your code has processed them.

Install

pip install ducklake-cdc-client

The package uses ducklake-client for DuckLake connections. CDCClient installs and loads the DuckDB community extension on first use:

INSTALL ducklake_cdc FROM community;
LOAD ducklake_cdc;

Batch iteration

from ducklake_client import DiskStorage, DuckDBCatalog, DuckLake
from ducklake_cdc_client import DMLConsumer

with DuckLake(
    catalog=DuckDBCatalog("metadata.ducklake"),
    storage=DiskStorage("data"),
) as lake:
    with DMLConsumer(
        lake,
        "orders-consumer",
        table="main.orders",
        mode="changes",
    ) as consumer:
        for batch in consumer.batches(infinite=False):
            for change in batch:
                print(change.to_dict())
            batch.commit()

batch.commit() advances the durable consumer cursor. If processing raises before that call, the same batch can be read again on the next run.

Connections, retries, and restart safety

Each high-level consumer uses one dedicated DuckDB connection for create/read/listen, heartbeats, and commit. This is required because the extension's lease belongs to the connection that acquired it. By default the client derives and owns that connection. If you pass connection= or client=, dedicate its connection to that one consumer: do not run unrelated queries on it, because cancellation calls connection.interrupt().

Known SQLite lock bursts, every observed H-022 deadlock spelling (thread::join failed, resource deadlock would occur, and resource deadlock avoided), and typed lease contention errors use the default bounded retry policy. Retry is not a recovery boundary for a poisoned DuckDB handle. After retries are exhausted, discard the consumer and its connection, open a new consumer instance, and resume the same durable consumer name. Call prewarm() on every handle immediately after loading the extension and before other catalog activity.

Lease failures are catchable as LeaseContentionError or LeaseTimeoutError; both inherit RetryableCDCError. A process supervisor should back off and reopen on a fresh dedicated connection. Reopening with on_exists="use" resumes from the last committed snapshot, so only commit after sink processing succeeds.

Graceful shutdown is bounded:

consumer.close(timeout=5.0, cancel=True, release=True)

cancel=True interrupts an active listen/read and that caller receives ConsumerCancelledError. cancel=False waits without interrupting. Either mode raises ConsumerCloseTimeoutError if the deadline expires. release=True clears a lease acquired by this instance; use release=False only when an external supervisor owns lease cleanup. Abrupt process death cannot run close(): the next process must wait for lease expiry or use lease_policy="takeover" only after it knows the previous holder is dead.

Cancellation ends the current run and is never retried in place. A lease timeout is marked retryable for supervisors but likewise does not repeat its entire wait internally; reopen after backoff instead.

Sink-driven usage

If you prefer a push style, pass sinks and let consumer.run() deliver and commit for you.

from ducklake_client import DiskStorage, DuckDBCatalog, DuckLake
from ducklake_cdc_client import DMLConsumer, StdoutSink

with DuckLake(
    catalog=DuckDBCatalog("metadata.ducklake"),
    storage=DiskStorage("data"),
) as lake:
    with DMLConsumer(
        lake,
        "orders-consumer",
        table="main.orders",
        mode="changes",
        sinks=[StdoutSink()],
    ) as consumer:
        consumer.run(infinite=False)

Demo

Run the local demo:

uv run python demo.py

The demo creates a local DuckLake catalog under .demo/, inserts one row into main.orders, prints the emitted CDC change batch, and commits it.

Test with a local extension build

The extension binary must match the Python package's DuckDB version exactly. This checkout pins DuckDB 1.5.4. After building the sibling extension repo, copy its unsigned macOS/arm64 binary into this repo's ignored local-artifact directory:

mkdir -p .local/extensions/v1.5.4/osx_arm64
cp ../ducklake-cdc-extension/build/release/extension/ducklake_cdc/ducklake_cdc.duckdb_extension \
  .local/extensions/v1.5.4/osx_arm64/
export DUCKLAKE_CDC_EXTENSION="$PWD/.local/extensions/v1.5.4/osx_arm64/ducklake_cdc.duckdb_extension"

Allow unsigned extensions when DuckDB opens, load the binary before constructing the consumer, and let the client skip its normal community-repository install:

import os

from ducklake_client import DiskStorage, DuckDBCatalog, DuckDBConfig, DuckLake
from ducklake_cdc_client import CDCClient, DMLConsumer

with DuckLake(
    catalog=DuckDBCatalog("metadata.ducklake"),
    storage=DiskStorage("data"),
    duckdb=DuckDBConfig(config={"allow_unsigned_extensions": True}),
) as lake:
    lake.connection.execute(f"LOAD '{os.environ['DUCKLAKE_CDC_EXTENSION']}'")
    client = CDCClient(lake, install_extension=False)
    with DMLConsumer(
        lake,
        "orders-consumer",
        table="main.orders",
        mode="changes",
        client=client,
    ) as consumer:
        print(consumer.client.version())

This binary is platform-specific (osx_arm64) and unsigned; rebuild it for another DuckDB version or platform instead of reusing it.

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

ducklake_cdc_client-0.6.0.tar.gz (47.4 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

ducklake_cdc_client-0.6.0-py3-none-any.whl (41.7 kB view details)

Uploaded Python 3

File details

Details for the file ducklake_cdc_client-0.6.0.tar.gz.

File metadata

  • Download URL: ducklake_cdc_client-0.6.0.tar.gz
  • Upload date:
  • Size: 47.4 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.1.0 CPython/3.13.12

File hashes

Hashes for ducklake_cdc_client-0.6.0.tar.gz
Algorithm Hash digest
SHA256 1d9a3c8237b707d28448c29ce84805da036748a2b9d8b53d78aa466a5af2c479
MD5 2096339e953a722487a347ea5ea02726
BLAKE2b-256 036c301ac04b3c6c04f215c5ca4410640685f9cc7187203736ef0dcb80b60a9e

See more details on using hashes here.

File details

Details for the file ducklake_cdc_client-0.6.0-py3-none-any.whl.

File metadata

File hashes

Hashes for ducklake_cdc_client-0.6.0-py3-none-any.whl
Algorithm Hash digest
SHA256 8d4d75cb5dc5879446b3f203e02ce75d9721f5b2c571b97460d3a04cd33ea54d
MD5 e74308616395d3295b23e42158de0d45
BLAKE2b-256 92ed59fc7f21548a03ca2fbd86292197ecfe31934af3d0405bd7843e854a9de1

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page