Skip to main content

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.

For one schema-independent cursor over every table in the catalogue, omit the table identity and use tick mode:

with DMLConsumer(
    lake,
    "catalogue-dml",
    mode="ticks",
) as consumer:
    for batch in consumer.batches(infinite=True):
        for tick in batch:
            publish_to_nats(tick.table_ids, tick.snapshot_id)
        batch.commit()

This cursor follows newly created tables and continues across DDL boundaries. Downstream systems can fan out by table_ids.

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 calls the owner-token- conditional cdc_consumer_release, so a stale close cannot clear a successor connection's lease. Use release=False only when an external supervisor owns lease cleanup. The unconditional cdc_consumer_force_release remains an operator recovery command. 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.

cdc_consumer_release requires ducklake_cdc >= 0.5.4. With an older extension, high-level close safely closes its owned connection and relies on lease expiry; it never falls back to unconditional force release. Upgrade the extension to make graceful release immediate.

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)

CDCApp.stats() includes generic operation activity for each worker: current_operation, operation_started_at, and last_operation_completed_at. Supervisors can use those timestamps to apply deployment-specific stall policy without putting timeout policy in the client.

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.5. 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.5/osx_arm64
cp ../ducklake-cdc-extension/build/release/extension/ducklake_cdc/ducklake_cdc.duckdb_extension \
  .local/extensions/v1.5.5/osx_arm64/
export DUCKLAKE_CDC_EXTENSION="$PWD/.local/extensions/v1.5.5/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.

Release files for ducklake-cdc-client 0.7.2

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for ducklake-cdc-client 0.7.2
File Size Uploaded
ducklake_cdc_client-0.7.2.tar.gz 50.1 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for ducklake-cdc-client 0.7.2
File Interpreter ABI Platform
ducklake_cdc_client-0.7.2-py3-none-any.whl Python 3 none any Details

Total release size: 93.6 kB

Release files / ducklake_cdc_client-0.7.2.tar.gz

Download URL ducklake_cdc_client-0.7.2.tar.gz
Size 50.1 kB
Tags Source
SHA-256 checksum
How to use checksums
a959043e1c8bffef4c9a42a279be96926b6e4aa58f1ebc2e17b6c0122f2920f9
BLAKE2b-256 checksum
How to use checksums
3c34a968337b0f36a1ed1803c68db60363338cfb7bbcbcf612355fe94a797977
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.1.0 CPython/3.13.14

Release files / ducklake_cdc_client-0.7.2-py3-none-any.whl

Download URL ducklake_cdc_client-0.7.2-py3-none-any.whl
Size 43.4 kB
Tags Python 3
SHA-256 checksum
How to use checksums
9deddb9f0a754343fdc60d108d0f36e51fda488d237a315db903666e0ccc7c18
BLAKE2b-256 checksum
How to use checksums
6367d5fcac07b2ed3c068714ee7eca04075e7ee8381e7237bee155e65562f795
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.1.0 CPython/3.13.14

Release history Release notifications | RSS feed

This release

0.7.2 This release

2 release files

0.7.1

2 release files

0.7.0

2 release files

0.6.1

2 release files

0.6.0

2 release files

0.5.1

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page