Skip to content

Integration engines

dbwarden supports 13 integration engines. Credentials use the named-collections declare-only pattern: no secret values are diffed.

Kafka

from dbwarden.databases.clickhouse import kafka

class Meta(CHTableMeta):
    ch = ch_table(
        engine=kafka(
            named_collection="kafka_prod",
            topic="events",
            format="JSONEachRow",
            group_name="dbwarden_consumer",
        ),
    )

Generated DDL:

CREATE TABLE kafka_events (
    payload String
) ENGINE = Kafka
SETTINGS kafka_named_collection = 'kafka_prod',
         kafka_topic_list = 'events',
         kafka_format = 'JSONEachRow',
         kafka_group_name = 'dbwarden_consumer';

Parameters that can be set directly (overriding the named collection):

Factory parameter DDL setting
named_collection kafka_named_collection
topic kafka_topic_list
format kafka_format
group_name kafka_group_name
num_consumers kafka_num_consumers
thread_per_consumer kafka_thread_per_consumer
handle_error_mode kafka_handle_error_mode
commit_every_batch kafka_commit_every_batch

KafkaSettings is a fully-typed TypedDict for arbitrary Kafka engine settings:

from dbwarden.databases.clickhouse import KafkaSettings

settings: KafkaSettings = {
    "kafka_max_block_size": 524288,
}

S3

from dbwarden.databases.clickhouse import s3

class Meta(CHTableMeta):
    ch = ch_table(
        engine=s3(
            named_collection="s3_prod",
            pattern="events/*.parquet",
            format="Parquet",
        ),
    )

Parameters:

Parameter DDL setting
named_collection s3_named_collection
pattern url (first positional)
format format
compression compression

S3Settings for arbitrary settings. Key naming variants (s3_*, s3queue_*) are all typed.

S3Queue

from dbwarden.databases.clickhouse import s3_queue

class Meta(CHTableMeta):
    ch = ch_table(
        engine=s3_queue(
            named_collection="s3_prod",
            pattern="incoming/*.json",
            format="JSONEachRow",
        ),
    )

S3QueueSettings covers all s3queue-specific settings.

RabbitMQ

from dbwarden.databases.clickhouse import rabbitmq

class Meta(CHTableMeta):
    ch = ch_table(
        engine=rabbitmq(
            named_collection="rabbit_prod",
            format="JSONEachRow",
        ),
    )

RabbitMQSettings, typed TypedDict.

NATS

from dbwarden.databases.clickhouse import nats

class Meta(CHTableMeta):
    ch = ch_table(
        engine=nats(
            named_collection="nats_prod",
            format="JSONEachRow",
        ),
    )

NatsSettings, typed TypedDict.

MySQL, PostgreSQL, MongoDB, Redis

from dbwarden.databases.clickhouse import mysql_engine, postgresql_engine, mongodb, redis

# MySQL engine
engine = mysql_engine(
    named_collection="mysql_prod",
    query="SELECT * FROM source_db.table",
)

# PostgreSQL engine
engine = postgresql_engine(
    named_collection="pg_prod",
    query="SELECT * FROM source_schema.source_table",
)

# MongoDB engine
engine = mongodb(
    named_collection="mongo_prod",
    collection="source_collection",
)

# Redis engine
engine = redis(
    named_collection="redis_prod",
    key="prefix:*",
)

Each has an associated *Settings TypedDict for engine-specific settings.

Additional model examples

Named collection for multi-engine reuse

# Single named collection reused by Kafka and S3 engines
named_collection(
    name="aws_prod",
    keys={
        "region": "us-east-1",
        "access_key_id": "AKIA...",
        # secret_access_key from secret store
    },
)

engine = kafka(
    named_collection="aws_prod",
    topic="events",
    format="JSONEachRow",
    group_name="ch_consumer",
)

engine2 = s3(
    named_collection="aws_prod",
    pattern="data/*.parquet",
    format="Parquet",
)

S3Queue with complex settings

engine = s3_queue(
    named_collection="aws_prod",
    pattern="incoming/*.json",
    format="JSONEachRow",
)

# With custom settings
settings: S3QueueSettings = {
    "s3queue_processing_threads": 8,
    "s3queue_polling_min_timeout_ms": 1000,
    "s3queue_polling_max_timeout_ms": 30000,
    "s3queue_tracked_files_limit": 100000,
}

PostgreSQL engine with query

engine = postgresql_engine(
    named_collection="pg_prod",
    query="SELECT id, name, created_at FROM public.users WHERE active = 1",
)

URL engine with multiple formats

# CSV
engine = url_engine(
    named_collection="http_data",
    format="CSV",
)

# With specific compression
engine = url_engine(
    named_collection="http_data",
    format="JSONEachRow",
    compression="gzip",
)

URL

from dbwarden.databases.clickhouse import url_engine

class Meta(CHTableMeta):
    ch = ch_table(
        engine=url_engine(
            named_collection="http_prod",
            format="CSV",
        ),
    )

URLSettings.

File

from dbwarden.databases.clickhouse import file_engine

class Meta(CHTableMeta):
    ch = ch_table(
        engine=file_engine(
            path="/var/lib/clickhouse/user_files/export.csv",
            format="CSV",
        ),
    )

HDFS

from dbwarden.databases.clickhouse import hdfs

class Meta(CHTableMeta):
    ch = ch_table(
        engine=hdfs(
            named_collection="hdfs_prod",
            format="Parquet",
        ),
    )

HDFSSettings.

Per-engine settings TypedDicts

Every integration engine has its own *Settings TypedDict for arbitrary settings. These are all defined in dbwarden.databases.clickhouse:

Engine TypedDict
Kafka KafkaSettings
S3 S3Settings
S3Queue S3QueueSettings
RabbitMQ RabbitMQSettings
NATS NatsSettings
MySQL MySQLSettings
PostgreSQL PostgreSQLSettings
MongoDB MongoDBSettings
Redis RedisSettings
URL URLSettings
HDFS HDFSSettings

What changes are allowed

Change Safety
Named collection swap CRITICAL (metadata only, not data)
Any setting INFO: ALTER TABLE MODIFY SETTING
Format change WARN: requires data re-ingestion
Pattern / query / key change WARN

Rollback behavior

Settings changes revert via RESET SETTING. Named collection swaps require reversing the collection reference.