dbt Adapter Reference for Confluent Cloud for Apache Flink

The dbt-confluent adapter runs dbt models as Flink SQL statements in Confluent Cloud for Apache Flink®. This page describes every profile setting, materialization, and model configuration that the adapter supports, and how the adapter behaves when you re-run, change, or rebuild a model.

If you’re new to the adapter, start with the end-to-end tutorial, Build a Streaming Pipeline with dbt and Confluent Cloud for Apache Flink.

Version compatibility

Adapter version

dbt Core

Python

Notes

0.4.x

1.11.x

3.10 to 3.13

Adds materialized_table, tableflow, statement_properties, and configuration validation.

0.3.2

1.11.x

3.10 to 3.13

Pins dbt Core to 1.11.x.

0.3.0, 0.3.1

1.11.x

3.10 to 3.13

These versions also allow dbt Core 1.12 to install, which isn’t compatible. Version 0.3.0 also allows Python 3.14, which isn’t supported. Pin dbt Core yourself with pip install "dbt-core~=1.11.0", or upgrade the adapter.

The adapter installs the dbt Core version it supports. Don’t install or upgrade dbt-core separately. To check your versions, run dbt --version.

How dbt concepts map to Flink

Confluent Cloud Flink uses different terms than the databases that dbt traditionally targets:

dbt concept

Flink SQL concept

Confluent Cloud entity

database

Catalog

Environment

schema

Database

Kafka cluster

Model

Table or view, plus the statements that maintain it

Kafka topic, except for views, ephemeral models, and faker tables

These mappings determine two required profile fields:

  • environment_id is the Confluent Cloud environment ID, in the form env-xxxxxx.

  • dbname is the name of the Kafka cluster, not the cluster ID (lkc-xxxxxx). A model-level schema configuration must also be a cluster name. The adapter uses a custom schema value as-is, instead of appending it to the target schema the way dbt does by default.

The adapter can’t create, rename, or drop Kafka clusters. Every cluster that a profile or model references must already exist.

Profile configuration

Define a connection profile in profiles.yml, which is by default at ~/.dbt/profiles.yml. dbt init creates one for you.

my_project:
  target: dev
  outputs:
    dev:
      type: confluent
      cloud_provider: aws
      cloud_region: us-east-2
      organization_id: <your-organization-id>
      environment_id: <your-environment-id>
      compute_pool_id: <your-compute-pool-id>
      dbname: <your-kafka-cluster-name>
      global_api_key: "{{ env_var('CONFLUENT_GLOBAL_API_KEY') }}"
      global_api_secret: "{{ env_var('CONFLUENT_GLOBAL_API_SECRET') }}"
      threads: 1

Setting

Description

type

Required. Always confluent.

organization_id

Required. Your Confluent Cloud organization ID.

environment_id

Required. The environment ID, for example env-xxxxxx.

dbname

Required. The name of the Kafka cluster that models are created in.

cloud_provider, cloud_region

The cloud provider and region of your Flink endpoint, for example aws and us-east-2. Required unless you set endpoint.

endpoint

A full Flink endpoint URL, for a private network or another non-standard endpoint, for example https://flink.us-east-2.aws.private.confluent.cloud. Set either endpoint or cloud_provider and cloud_region, not both.

global_api_key, global_api_secret

A global API key. Works for Flink and is required for Tableflow.

flink_api_key, flink_api_secret

A Flink API key, which is scoped to one environment and region. Supply one complete key pair, either global or Flink. If you supply both, the adapter uses the global key.

compute_pool_id

Optional. The default compute pool for every model. If you omit it, Flink runs statements in the environment’s default compute pool for the region. Models can override it; see compute_pool_id.

execution_mode

Optional. The default execution mode for statements that don’t set their own. Default: streaming_query. Materializations set their own mode, so you rarely need to change this.

statement_name_prefix

Optional. The prefix for generated statement names. Default: dbt-. See Statement names.

statement_label

Optional. A label applied to every statement that the adapter submits. Default: dbt-confluent.

threads

The number of models that dbt builds in parallel.

Store secrets in environment variables and read them with env_var(), as in the preceding example. To confirm that the profile works, run dbt debug, which also reports the configured compute pool.

Materializations

A materialization determines what the adapter creates for a model and which Flink SQL statements it submits.

Materialization

Use it for

Execution

materialized_table

Continuously maintained results. Use it for most streaming pipelines, because changes evolve in place.

Streaming. Flink maintains the table.

streaming_table

Continuously maintained results with an explicit table and a separate INSERT statement, or adopting a pipeline that you deployed outside dbt.

Streaming. A long-running INSERT INTO ... SELECT.

streaming_source

A table with a declared schema, backed by the faker or confluent Flink SQL connector.

Depends on the connector. A faker table generates data for each statement that reads it. A confluent table stores data in a Kafka topic.

table

A one-time, point-in-time result, like a batch warehouse table.

Snapshot. The query runs once and completes.

view

Reusable query logic that stores no data. This is the default.

Runs inside each statement that reads it.

ephemeral

Query logic that dbt inlines as a common table expression (CTE).

Runs inside each statement that reads it.

The following dbt materializations aren’t supported:

Materialization

Alternative

materialized_view

Use materialized_table for continuously maintained results, or table for a one-time result.

incremental

dbt’s batch-incremental model doesn’t map to continuous processing. Use materialized_table.

snapshot

dbt snapshots require MERGE and UPDATE operations that Flink SQL doesn’t support. See Limitations.

Set a materialization in a model’s config() block or for a whole directory in dbt_project.yml:

models:
  my_project:
    staging:
      +materialized: view
    marts:
      +materialized: materialized_table

materialized_table

A materialized_table model submits a CREATE OR ALTER MATERIALIZED TABLE statement with the model’s SELECT query. Flink then maintains the table continuously. For how materialized tables work, see Materialized Tables in Confluent Cloud for Apache Flink.

{{ config(
    materialized='materialized_table',
    distributed_by={'columns': ['customer_id'], 'buckets': 4},
    start_mode='RESUME_OR_FROM_BEGINNING',
) }}

SELECT
    customer_id,
    COUNT(*) AS order_count,
    SUM(price) AS lifetime_value
FROM {{ ref('orders') }}
GROUP BY customer_id

materialized_table is declarative. Every run submits the same statement, whether the table doesn’t exist yet, exists unchanged, or exists with a different definition. Flink treats an unchanged definition as a no-op and evolves a changed definition in place. This model doesn’t use schema drift detection.

Supported configuration: with, distributed_by, start_mode, statement_properties, tableflow, statement_name, compute_pool_id, and ignore_unsupported_config. You can also enforce a dbt model contract. See Contracts and primary keys. materialized_table also rejects the open source Flink options freshness_interval, refresh_mode, and partition_by, which Confluent Cloud doesn’t support. Setting any of them fails the run.

Start mode

start_mode controls where the query starts reading when the table is created, and where it restarts after the table evolves. Default: RESUME_OR_FROM_BEGINNING.

Value

Starts reading from

FROM_BEGINNING

The beginning of each source topic.

FROM_NOW

The current offset.

FROM_TIMESTAMP(TIMESTAMP '<yyyy-MM-dd HH:mm:ss>')

The given timestamp.

FROM_NOW(INTERVAL '<n>' <unit>)

<n> <unit> before now.

RESUME_OR_FROM_BEGINNING

Saved offsets if they exist, otherwise the beginning.

RESUME_OR_FROM_NOW

Saved offsets if they exist, otherwise the current offset.

RESUME_OR_FROM_TIMESTAMP(TIMESTAMP '<yyyy-MM-dd HH:mm:ss>')

Saved offsets if they exist, otherwise the given timestamp.

RESUME_OR_FROM_NOW(INTERVAL '<n>' <unit>)

Saved offsets if they exist, otherwise <n> <unit> before now.

Evolution and state

When you change a materialized_table model and run dbt run, Flink evolves the table in place. During an evolution, the query’s internal state, such as aggregations, joins, and windows, is discarded, and processing restarts according to start_mode:

  • Stateless queries, such as projections and filters, evolve without reprocessing or duplicates under a RESUME_* start mode, which includes the default.

  • Stateful queries recalculate their results from empty state instead of adjusting the previous results. Under a RESUME_* start mode, the new results cover only records that arrive after the change, so aggregates can look undercounted. To recompute from the beginning, run the model with --full-refresh.

For more information, see How to evolve materialized tables.

Some changes can’t evolve in place. Dropping a column that isn’t nullable, or changing distributed_by, fails the run. To apply these changes, run with --full-refresh.

--full-refresh drops and recreates the materialized table. This deletes the backing Kafka topic, all of its records, and its Schema Registry schema versions, and the new topic has a new identity even though its name is the same. Consumers of the topic, including other Flink SQL statements, can break. Coordinate full refreshes with downstream consumers.

Contracts and primary keys

To declare columns and a primary key, enforce a dbt model contract and add a model-level primary_key constraint in the model’s YAML file:

models:
  - name: customer_order_stats
    config:
      contract:
        enforced: true
    constraints:
      - type: primary_key
        columns: [customer_id]
        expression: "NOT ENFORCED"
    columns:
      - name: customer_id
        data_type: bigint
      - name: order_count
        data_type: bigint
      - name: lifetime_value
        data_type: decimal(38, 2)

Flink SQL requires primary keys to be NOT ENFORCED, so always set expression: "NOT ENFORCED" on the constraint. The adapter renders the columns and a PRIMARY KEY (...) NOT ENFORCED clause in the CREATE OR ALTER MATERIALIZED TABLE statement. Without an enforced contract, Flink infers the columns from the SELECT query.

Switch to or from materialized_table

The adapter doesn’t convert an existing table or view to a materialized table, or convert a materialized table in place to another materialization, even though Flink SQL can adopt an existing table as a materialized table. To switch, run the model with --full-refresh, which drops the existing table and its data. Without --full-refresh, the result depends on the direction of the switch:

  • To materialized_table from table, view, streaming_table, or streaming_source: the run fails with guidance.

  • From materialized_table to table, streaming_table, or streaming_source: schema drift detection fails the run with guidance. With on_schema_drift='ignore', the result depends on the new materialization:

    • table or streaming_source: drift detection doesn’t run, so the model is skipped and the materialized table keeps running.

    • streaming_table: the run still fails with guidance. The adapter never submits an INSERT statement into the materialized table.

  • From materialized_table to view: the run drops the materialized table, including its topic and data, and creates the view. Check a model’s downstream consumers before you make this change.

streaming_table

A streaming_table model submits two statements: a CREATE TABLE statement, and a long-running INSERT INTO ... SELECT statement that populates the table. Use it to adopt a pipeline that you deployed outside dbt, or when you need a separate INSERT statement. For other streaming pipelines, use materialized_table.

{{ config(
    materialized='streaming_table',
    distributed_by={'columns': ['customer_id'], 'buckets': 4},
    with={'changelog.mode': 'append'},
) }}

SELECT customer_id, order_id, price
FROM {{ ref('orders') }}
WHERE price > 0

Supported configuration: with, distributed_by, on_schema_drift, statement_properties, tableflow, statement_name, compute_pool_id, and ignore_unsupported_config.

Re-run behavior

When the table already exists and you run without --full-refresh:

  1. The adapter runs schema drift detection. If the model no longer matches the table, the run fails.

  2. The adapter checks the INSERT statement. If the statement is missing, or is in a terminal phase (COMPLETED, STOPPED, or FAILED), the adapter resubmits only the INSERT statement, under the same name. The table and its data are kept.

  3. A statement that is RUNNING, DEGRADED, or in transition (PENDING, STOPPING, or DELETING) is left alone.

Because a healthy statement is left alone, a change to a model’s query logic, statement_properties, or compute_pool_id takes effect only when you run with --full-refresh or the adapter resubmits a stopped statement.

Adopt an existing pipeline

You can bring a pipeline that you deployed outside dbt under dbt management without recreating it. Point a streaming_table model at the existing table with dbt’s alias configuration, and at the existing INSERT statement with statement_name:

{{ config(
    materialized='streaming_table',
    alias='orders_enriched',
    statement_name='orders-enriched-insert',
    with={'changelog.mode': 'append'},
) }}

SELECT order_id, price FROM {{ ref('orders') }}

On the next dbt run without --full-refresh, the adapter looks up the table and statement by name. It adopts a healthy statement as-is, resubmits a missing or terminal one under the same name, and never drops the table.

  • The model’s schema must match the existing table. Under the default on_schema_drift='fail', any mismatch fails the run.

  • The statement name must match exactly after the adapter normalizes it. For the normalization rules, see Statement names.

streaming_source

A streaming_source model creates a table with a Flink SQL connector. The model body is the table’s column list, not a SELECT query. The connector configuration is required and sets the connector option of the table’s CREATE TABLE statement. These are Flink SQL table connectors, not Kafka Connect connectors.

The adapter supports two connectors:

Connector

Creates

faker

A table that generates sample data. For options, see Generate Custom Sample Data with Confluent Cloud for Apache Flink.

confluent

A table backed by a Kafka topic, with the schema that the model declares. Use it when other applications produce to the topic and you want dbt to own the table’s schema. Running this model with --full-refresh deletes the topic and every record that other applications produced to it.

The following model creates a faker table:

{{ config(
    materialized='streaming_source',
    connector='faker',
    with={
        'rows-per-second': '1',
        'fields.order_id.expression': "#{number.numberBetween '1','1000000'}",
    },
) }}

order_id BIGINT,
price DECIMAL(10, 2),
order_time TIMESTAMP(3),
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND

The following model creates a confluent table:

{{ config(
    materialized='streaming_source',
    connector='confluent',
) }}

order_id BIGINT,
status STRING

The adapter merges connector into the table’s WITH options and escapes single quotes in keys and values, so write quotes as-is.

Supported configuration: connector, with, distributed_by, on_schema_drift, tableflow, statement_name, compute_pool_id, and ignore_unsupported_config. distributed_by and tableflow apply only to confluent tables, because a faker table has no topic.

A Faker table doesn’t store data. Flink generates rows separately for each statement that reads the table, so two models that read the same Faker table get different data. To give several models the same data, or to let Flink reprocess it later, read the Faker table from one materialized_table model, which stores the rows in a topic, and have the other models read from that model.

When the table already exists, the adapter runs schema drift detection and skips the model. Drift detection includes connector and the with options, so to change them, run the model with --full-refresh, which drops and recreates the table.

table

A table model submits a CREATE TABLE ... AS SELECT statement in snapshot mode. The query reads a point-in-time view of its sources, writes the result, and completes. Use it for batch-style results and as a first model when you learn the adapter.

{{ config(materialized='table') }}

SELECT customer_id, COUNT(*) AS order_count
FROM {{ ref('orders') }}
GROUP BY customer_id

Supported configuration: distributed_by, on_schema_drift, tableflow, statement_name, compute_pool_id, and ignore_unsupported_config.

When the table already exists, the adapter runs schema drift detection and skips the model. To recompute the result, run with --full-refresh, which drops the table and its data and runs the query again.

view

A view model submits a CREATE VIEW statement. A view stores no data. Flink inlines the view’s query into each statement that reads from it. On every run, the adapter drops and recreates the view.

Supported configuration: statement_name, compute_pool_id, and ignore_unsupported_config.

ephemeral

An ephemeral model behaves as in any other dbt adapter. dbt inlines the model’s SQL as a CTE in every model that references it, and the adapter creates nothing. The adapter doesn’t validate configuration keys on ephemeral models. If several models reference the same ephemeral model, each statement reads the underlying source separately.

Model configuration

The following table shows which adapter-specific configuration keys each materialization accepts. If you set a key on a materialization that doesn’t accept it, the run fails with an error instead of ignoring the key, unless you list the key in ignore_unsupported_config. For more information, see Configuration validation.

Configuration key

materialized_table

streaming_table

streaming_source

table

view

with

Yes

Yes

Yes

No

No

distributed_by

Yes

Yes

Yes

Yes

No

tableflow

Yes

Yes

Yes

Yes

No

start_mode

Yes

No

No

No

No

statement_properties

Yes

Yes

No

No

No

on_schema_drift

No

Yes

Yes

Yes

No

connector

No

No

Required

No

No

statement_name, compute_pool_id, ignore_unsupported_config

Yes

Yes

Yes

Yes

Yes

with

A dictionary of table options that the adapter renders in the table’s WITH clause, for example {'changelog.mode': 'append'} or {'value.format': 'avro-registry'}. For available options, see CREATE TABLE Statement in Confluent Cloud for Apache Flink.

distributed_by

Controls how rows are distributed across the partitions of the backing Kafka topic. The adapter renders a DISTRIBUTED BY HASH(...) INTO <n> BUCKETS clause.

{{ config(distributed_by={'columns': ['order_id'], 'buckets': 4}) }}

SELECT order_id, customer_id, price FROM {{ ref('orders') }}

Field

Description

columns

Required. A non-empty list of column names to hash.

buckets

Optional. A positive integer. If you omit it, Confluent Cloud chooses the number of buckets.

The distribution columns must be the first columns of the table, in the same order as columns. In a SELECT model, list them first. In a streaming_source model, declare them first. Otherwise, Flink rejects the statement. You can’t change distributed_by on an existing table without --full-refresh.

statement_properties

A dictionary of statement properties, the same properties that you set with the SET statement, for example {'sql.tables.scan.idle-timeout': '30 s'}. Values can be strings, integers, or Booleans.

{{ config(
    materialized='materialized_table',
    statement_properties={'sql.tables.scan.idle-timeout': '30 s'},
) }}

For materialized_table, the next run applies changed properties to the materialized table, even when the query is unchanged. Treat this like any other change to the model: the query restarts, which can reset a stateful model’s results. See Evolution and state. For streaming_table, they apply only to the INSERT statement and take effect only when the adapter resubmits a missing or terminal statement, or when you run with --full-refresh. To set table options, use with instead.

You can’t set sql.current-catalog, sql.current-database, or sql.snapshot.mode, because the adapter manages them.

compute_pool_id

Runs a model’s statements on a specific compute pool instead of the profile default, for example to isolate a heavy model:

{{ config(compute_pool_id='lfcp-abc123') }}

The pool must be in the same environment and region as the profile, and your API key must have access to it. To use a different pool in each environment, read the ID from an environment variable:

{{ config(compute_pool_id=env_var('FLINK_COMPUTE_POOL')) }}

How a changed pool takes effect depends on the materialization:

  • view: the next run uses the new pool.

  • materialized_table: changing only compute_pool_id doesn’t move an existing materialized table to the new pool. Its query keeps running on the pool where the table was created. To move it and keep its data, update the materialized table with the REST API. Or run the model with --full-refresh, which drops the table and its data.

  • streaming_table: a running INSERT statement doesn’t move. The new pool is used when the adapter resubmits a missing or stopped statement, or when you run with --full-refresh.

  • streaming_source and table: the new pool is used the next time you run the model with --full-refresh.

tableflow

Enables Tableflow on the Kafka topic behind a model, which materializes the topic as an Apache Iceberg™ table, a Delta Lake table, or both. You can query the result from analytics engines or from Flink with a snapshot query.

{{ config(
    materialized='materialized_table',
    tableflow={
        'table_formats': ['ICEBERG'],
        'storage': {'kind': 'Managed'},
        'config': {
            'retention_ms': 604800000,
            'error_handling': {'mode': 'LOG', 'target': 'error_log'},
        },
    },
) }}

Field

Description

table_formats

Required. 'ICEBERG', 'DELTA', or a list with one or both.

storage

Required. Where Tableflow writes the tables. Set kind to one of the following:

  • {'kind': 'Managed'} for Confluent-managed storage.

  • {'kind': 'ByobAws', 'bucket_name': '...', 'provider_integration_id': '...'} for your own Amazon S3 bucket.

  • {'kind': 'AzureDataLakeStorageGen2', 'storage_account_name': '...', 'container_name': '...', 'provider_integration_id': '...'} for your own Azure Data Lake Storage Gen2 container.

  • {'kind': 'GoogleCloudStorage', 'bucket_name': '...', 'provider_integration_id': '...'} for your own Google Cloud Storage bucket.

For more information, see Storage with Tableflow in Confluent Cloud.

config

Optional. Topic-level Tableflow settings:

  • retention_ms: How long to keep table snapshots, in milliseconds.

  • data_retention_ms: How long to keep table data, in milliseconds.

  • error_handling: What to do when a record can’t be materialized. {'mode': 'SUSPEND'} suspends Tableflow, which is the default. {'mode': 'SKIP'} skips the record. {'mode': 'LOG', 'target': '<topic>'} writes the record to a dead letter queue (DLQ) topic, error_log by default, and continues. For more information, see Error-handling mode.

The adapter reconciles Tableflow on every run:

  • If Tableflow isn’t enabled, the adapter enables it.

  • If you changed table_formats or config, the adapter updates Tableflow in place. If nothing changed, it sends no update.

  • If you changed storage, the adapter disables Tableflow and enables it again with the new storage.

  • If Tableflow is in a FAILED state and there is nothing to update, the run reports a warning.

The adapter changes only the settings that the model sets. Removing retention_ms, data_retention_ms, or error_handling from the configuration doesn’t reset them, and removing the tableflow configuration from a model doesn’t disable Tableflow. To disable Tableflow or reset a setting, use the Cloud Console.

When you run with --full-refresh, the adapter disables Tableflow before it drops the table only if the model still sets tableflow. If you removed the configuration, disable Tableflow in the Cloud Console before you run with --full-refresh.

Tableflow requirements and limitations:

  • Tableflow requires a global API key in the profile. A Flink API key can’t reach the Tableflow APIs. A global API key works for Flink too, so you don’t need both keys.

  • For Tableflow permissions, see Grant Role-Based Access for Tableflow in Confluent Cloud.

  • A storage change turns Tableflow off for the topic until the adapter enables it again, and Confluent Cloud can reject enabling Tableflow on the same topic for a while after it’s disabled. Until the run succeeds, the Tableflow table stops updating and dbt run fails for that model. For how disabling works, see Disable Tableflow. To rebuild immediately instead, run with --full-refresh, which also deletes the topic’s data.

Configuration validation

If you set an adapter-specific configuration key on a materialization that doesn’t use it, for example statement_properties on a table model, the run fails with an error that names the key. The check covers only the adapter’s own keys, which are listed in Model configuration, and the open source options that materialized_table rejects. It never inspects other configuration keys, including keys that your own macros read.

If one of your own keys has the same name as an adapter key, exclude it from the check for that model:

{{ config(
    materialized='table',
    statement_properties={'my_custom_key': 'value'},
    ignore_unsupported_config=['statement_properties'],
) }}

Schema drift detection

table, streaming_table, and streaming_source models skip creation when their table already exists. Before skipping, the adapter checks whether the existing table still matches the model, and fails the run if it doesn’t. materialized_table models don’t use drift detection, because Flink reconciles their definition on every run.

The check compares:

  • Columns: names and data types. Added, removed, renamed, or retyped columns are drift. Column order doesn’t matter.

  • distributed_by: the column list, column order, and bucket count, if you set buckets. If you don’t set distributed_by, the adapter doesn’t check distribution.

  • WITH options for streaming_table and streaming_source: every option that you set must have the same value on the table. For streaming_source, this includes connector.

The run reports every mismatch in a single error. To apply the model’s current definition, run with --full-refresh, which drops and recreates the table and deletes its data.

Drift detection doesn’t catch everything:

  • Removed WITH options aren’t detected, because connectors add their own options that the adapter can’t tell apart from yours.

  • Query logic changes that keep the same column names and types aren’t detected. For example, changing ROUND(price, 2) to ROUND(price, 4) isn’t drift. Run with --full-refresh to apply the change.

To turn off drift detection for a model, set on_schema_drift='ignore'. The adapter then skips the model whenever its table exists. The default is 'fail'. For streaming_table, the adapter still resubmits a missing or terminal INSERT statement, and before it does, it still checks that the columns match, because Flink would reject the statement otherwise.

{{ config(materialized='streaming_table', on_schema_drift='ignore') }}

To check the schema, the adapter creates a temporary table named __dbt_tmp_schema_check_<model> and drops it when the model finishes, even if the run fails.

Statement names

The adapter gives each statement a predictable name: <statement_name_prefix><project>-<model>, for example dbt-my-project-stg-orders-a1b2c3. The six-character hash at the end appears whenever the name contains characters that the adapter replaces, which includes the underscores in most dbt project and model names. The prefix defaults to dbt-. Other statements use suffixes:

  • A streaming_table model’s CREATE TABLE statement has a -ddl suffix.

  • A materialized_table model’s statement has a unique suffix for each run, because the statement completes as soon as Flink accepts the definition. If the statement fails, it’s kept for debugging.

To choose a name, set the statement_name configuration. The adapter uses it in place of the generated name, without statement_name_prefix. It still adds the suffixes in the preceding list, so a materialized_table statement always gets a per-run suffix. For an exact name, as when you adopt an existing pipeline, use streaming_table.

Statement names can contain only lowercase letters, digits, and hyphens, must start and end with a letter or digit, and can be up to 100 characters long. The adapter replaces other characters, including underscores, with hyphens and appends a six-character hash so that names stay unique. Names longer than 100 characters are shortened and get a hash.

When you run with --full-refresh, the adapter deletes a model’s statements before it drops and recreates the table.

Reference existing topics with sources

To read from a table that already exists in Flink, such as a topic that another application or team produces, declare it as a dbt source. A source creates nothing; it lets dbt resolve the table by name and show it in the lineage graph. This is different from streaming_source, which creates a new table.

# models/sources.yml
version: 2

sources:
  - name: marketplace
    database: examples
    tables:
      - name: orders
      - name: clicks

A source’s database is the Flink SQL catalog, which is a Confluent Cloud environment or the read-only examples catalog. Its schema is the Flink SQL database, which is a Kafka cluster name. If you omit schema, dbt uses the source name, so the preceding example resolves to examples.``marketplace``. For your own topics, set database to your environment ID and schema to your cluster name. Reference the source with {{ source() }}:

SELECT
    `order_id`,
    `customer_id`,
    `price`,
    `$rowtime` AS order_time
FROM {{ source('marketplace', 'orders') }}

Flink SQL quotes identifiers with backticks. Use backticks for names with special characters, such as the $rowtime system column.

Testing

The adapter supports dbt unit tests and data tests. For a worked example, see Step 6: Test the pipeline.

  • For each unit test, the adapter creates temporary tables with CREATE TABLE ... LIKE for the model’s inputs, inserts the mock rows, runs the model, and compares the output with the expected rows. It drops the temporary tables afterward. The model and its inputs must already exist in the target, so deploy with dbt run before you run unit tests. The adapter finds the tested model’s table by the model’s name, so unit tests don’t support models whose table name differs from the model name, such as models that set alias or projects that override generate_alias_name.

  • Data tests run as queries against the deployed models. A test passes when its query returns no rows. A data test checks every row that is in the table when the test runs, including tables that are still receiving data. Rows that arrive later are checked the next time you run dbt test.

  • Don’t use --store-failures or the store_failures configuration. By default, dbt stores failures in a dbt_test__audit schema, which fails the run unless a Kafka cluster with that name exists. If you point the failures at an existing cluster, a failing test can report PASS, because the adapter counts the stored failures before they are written.

In streaming mode, COUNT(*) over an empty result returns no rows instead of one row with the value 0. The adapter accounts for this in test assertions.

Manage deployments

Dependencies and selective deployment

{{ ref() }} and {{ source() }} define your pipeline’s dependency graph. dbt deploys models in dependency order, so upstream models exist before the models that read from them. Use dbt’s node selection syntax to list or deploy part of the graph:

# List stg_orders and everything downstream of it
dbt ls --select stg_orders+

# Deploy fct_revenue and everything upstream of it
dbt run --select +fct_revenue

# Deploy only models that changed since the last deployment
dbt run --select state:modified --defer --state <path-to-previous-artifacts>

For streaming_table, table, and streaming_source models, state:modified selects a changed model, but the adapter doesn’t apply query changes to an existing table. Rebuild those models with --full-refresh.

Before you rebuild a model with --full-refresh, list its downstream models. Rebuilding a table gives its topic a new identity, so statements that read from it can also need a rebuild. To explore the lineage graph, run dbt docs generate and dbt docs serve.

For schema compatibility and statement evolution beyond what the adapter manages, see Schema and Statement Evolution with Confluent Cloud for Apache Flink and Carry-over Offsets in Confluent Cloud for Apache Flink.

CI/CD

Run dbt from your CI/CD system to test and deploy on every merge. Unit tests need the tested models and their inputs to exist, so deploy to a CI environment first, test there, and then deploy to production. The following GitHub Actions workflow uses the ci and prod targets from the profiles.yml file in Multiple environments. It expects that file at the root of your repository, which is what --profiles-dir . points to. The file reads every credential with env_var(), so it contains no secrets and is safe to commit.

on:
  push:
    branches:
      - main

jobs:
  dbt_deploy:
    name: "Deploy Flink SQL with dbt"
    runs-on: ubuntu-latest
    steps:
      - name: Checkout
        uses: actions/checkout@v4

      - name: Set up Python
        uses: actions/setup-python@v5
        with:
          python-version: '3.12'

      - name: Install dependencies
        run: pip install "dbt-confluent==<adapter-version>"

      - name: Deploy to CI
        run: dbt run --target ci --profiles-dir .
        env:
          CI_CONFLUENT_GLOBAL_API_KEY: ${{ secrets.CI_CONFLUENT_GLOBAL_API_KEY }}
          CI_CONFLUENT_GLOBAL_API_SECRET: ${{ secrets.CI_CONFLUENT_GLOBAL_API_SECRET }}

      - name: Test in CI
        run: dbt test --target ci --profiles-dir .
        env:
          CI_CONFLUENT_GLOBAL_API_KEY: ${{ secrets.CI_CONFLUENT_GLOBAL_API_KEY }}
          CI_CONFLUENT_GLOBAL_API_SECRET: ${{ secrets.CI_CONFLUENT_GLOBAL_API_SECRET }}

      - name: Deploy to production
        run: dbt run --target prod --profiles-dir .
        env:
          PROD_CONFLUENT_GLOBAL_API_KEY: ${{ secrets.PROD_CONFLUENT_GLOBAL_API_KEY }}
          PROD_CONFLUENT_GLOBAL_API_SECRET: ${{ secrets.PROD_CONFLUENT_GLOBAL_API_SECRET }}

Each step gets only the credentials for its own target, so the CI steps can’t read the production key. Store the API keys and secrets for both environments as GitHub Actions secrets, and pin the adapter version so that every run uses the same adapter. To require approval before production deployments, run the production step in a GitHub environment with required reviewers.

Multiple environments

To deploy to more than one environment, for example a CI environment and production, define a target for each in profiles.yml and choose one with --target:

my_project:
  target: ci
  outputs:
    ci:
      type: confluent
      cloud_provider: aws
      cloud_region: us-east-2
      organization_id: <your-organization-id>
      environment_id: <ci-environment-id>
      compute_pool_id: <ci-compute-pool-id>
      dbname: <ci-kafka-cluster-name>
      global_api_key: "{{ env_var('CI_CONFLUENT_GLOBAL_API_KEY') }}"
      global_api_secret: "{{ env_var('CI_CONFLUENT_GLOBAL_API_SECRET') }}"
      threads: 1
    prod:
      type: confluent
      cloud_provider: aws
      cloud_region: us-east-2
      organization_id: <your-organization-id>
      environment_id: <prod-environment-id>
      compute_pool_id: <prod-compute-pool-id>
      dbname: <prod-kafka-cluster-name>
      global_api_key: "{{ env_var('PROD_CONFLUENT_GLOBAL_API_KEY') }}"
      global_api_secret: "{{ env_var('PROD_CONFLUENT_GLOBAL_API_SECRET') }}"
      threads: 4
dbt run --target prod

Troubleshooting

dbt run fails on the example models

The example models that dbt init creates in models/example/ select an untyped NULL (select null as id), which Flink SQL rejects. Delete the models/example/ directory and the example entry under models: in dbt_project.yml.

macro ‘dbt_macro__statement’ takes no keyword argument ‘hidden’

dbt Core 1.12 is installed, which isn’t compatible with the adapter. Install dbt Core 1.11:

pip install "dbt-core~=1.11.0"

Adapter versions 0.3.2 and later prevent this by pinning dbt Core to 1.11.x. See Version compatibility.

403 “You don’t have the required permissions” without a compute pool

If you omit compute_pool_id, Flink runs statements in the default compute pool, and it returns a 403 error when your principal can’t use it. This happens when default compute pools are disabled for your organization, or when your FlinkDeveloper role binding is scoped to a specific compute pool. Set compute_pool_id in the profile to a pool that your principal can use. For more information, see Compute Pools in Confluent Cloud for Apache Flink and Grant Role-Based Access in Confluent Cloud for Apache Flink.

Connection fails with endpoint and cloud region both set

endpoint replaces cloud_provider and cloud_region. If you set endpoint and either of the others, the connection fails. Remove cloud_provider and cloud_region when you use endpoint, for example for private networking. Adapter versions earlier than 0.3.1 report every connection failure as a generic confluent_sql connection error. Upgrade to see the underlying cause.

Tableflow fails with an authentication error

Tableflow requires a global API key. If the profile has only a Flink API key, replace it with a global API key, which also works for Flink.

Certificate errors behind a TLS inspection proxy

If your network inspects TLS traffic with a proxy, such as Zscaler, pip install or dbt debug can fail with CERTIFICATE_VERIFY_FAILED. On Python 3.13, the error can include Basic Constraints of CA cert not marked critical.

These errors have two causes:

  • The adapter verifies server certificates against the certifi CA bundle, not your operating system’s truststore. It doesn’t trust your proxy’s certificate authority (CA) unless you set the SSL_CERT_FILE or SSL_CERT_DIR environment variable, even when your browser works.

  • Python 3.13 and later check certificates strictly by default and reject CA certificates that don’t follow RFC 5280. Some proxy root CAs fail this check.

To fix the errors without turning off certificate verification:

  1. Get your proxy’s CA certificates in PEM format from your IT team. Include the proxy’s intermediate CA certificate. Python 3.13 accepts a trusted intermediate certificate without checking the root above it, which avoids the strict-mode failure.

  2. Create a bundle that contains the public CAs and your proxy’s CAs:

    mkdir -p ~/.config/certs
    cat "$(python3 -m pip._vendor.certifi)" proxy-ca.pem > ~/.config/certs/ca-bundle.pem
    
  3. Point Python at the bundle. Set the variable in your shell profile, because you need it both when you install the adapter and when you run dbt:

    export SSL_CERT_FILE=~/.config/certs/ca-bundle.pem
    

    For pip, also pass the bundle with --cert or set PIP_CERT to the same path.

Limitations

  • No dbt snapshots: Flink SQL doesn’t support the MERGE and UPDATE operations that dbt’s snapshot resource type needs to track slowly changing dimensions. This is unrelated to Flink SQL snapshot queries, which the adapter uses for table models and unit tests.

  • No incremental models: Use materialized_table instead.

  • No schema management: The adapter can’t create, rename, or drop Kafka clusters. See How dbt concepts map to Flink.

  • No renames: Flink SQL doesn’t support renaming tables. Renaming a model makes dbt create a new table under the new name. The old table, its topic, and any statement that writes to it keep running, so remove them yourself. Drop the old table with DROP MATERIALIZED TABLE or DROP TABLE. For a streaming_table model, also delete its old INSERT statement.

  • Not transactional: If a dbt run fails partway through, the models that already ran stay deployed.

  • No stored test failures: --store-failures and the store_failures configuration don’t work reliably. See Testing.