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 |
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 |
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 |
|---|---|---|
|
Catalog |
Environment |
|
Database |
Kafka cluster |
Model |
Table or view, plus the statements that maintain it |
Kafka topic, except for views, |
These mappings determine two required profile fields:
environment_idis the Confluent Cloud environment ID, in the formenv-xxxxxx.dbnameis the name of the Kafka cluster, not the cluster ID (lkc-xxxxxx). A model-levelschemaconfiguration must also be a cluster name. The adapter uses a customschemavalue 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 |
|---|---|
|
Required. Always |
|
Required. Your Confluent Cloud organization ID. |
|
Required. The environment ID, for example |
|
Required. The name of the Kafka cluster that models are created in. |
|
The cloud provider and region of your Flink endpoint, for example
|
|
A full Flink endpoint URL, for a
private network or another
non-standard endpoint, for example
|
|
A global API key. Works for Flink and is required for Tableflow. |
|
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. |
|
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. |
|
Optional. The default execution mode for statements that don’t set
their own. Default: |
|
Optional. The prefix for generated statement names. Default:
|
|
Optional. A label applied to every statement that the adapter
submits. Default: |
|
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 |
|---|---|---|
|
Continuously maintained results. Use it for most streaming pipelines, because changes evolve in place. |
Streaming. Flink maintains the table. |
|
Continuously maintained results with an explicit table and a
separate |
Streaming. A long-running |
|
A table with a declared schema, backed by the |
Depends on the connector. A |
|
A one-time, point-in-time result, like a batch warehouse table. |
Snapshot. The query runs once and completes. |
|
Reusable query logic that stores no data. This is the default. |
Runs inside each statement that reads it. |
|
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 |
|---|---|
|
Use |
|
dbt’s batch-incremental model doesn’t map to continuous processing.
Use |
|
dbt snapshots require |
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 |
|---|---|
|
The beginning of each source topic. |
|
The current offset. |
|
The given timestamp. |
|
|
|
Saved offsets if they exist, otherwise the beginning. |
|
Saved offsets if they exist, otherwise the current offset. |
|
Saved offsets if they exist, otherwise the given timestamp. |
|
Saved offsets if they exist, otherwise |
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_tablefromtable,view,streaming_table, orstreaming_source: the run fails with guidance.From
materialized_tabletotable,streaming_table, orstreaming_source: schema drift detection fails the run with guidance. Withon_schema_drift='ignore', the result depends on the new materialization:tableorstreaming_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 anINSERTstatement into the materialized table.
From
materialized_tabletoview: 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:
The adapter runs schema drift detection. If the model no longer matches the table, the run fails.
The adapter checks the
INSERTstatement. If the statement is missing, or is in a terminal phase (COMPLETED,STOPPED, orFAILED), the adapter resubmits only theINSERTstatement, under the same name. The table and its data are kept.A statement that is
RUNNING,DEGRADED, or in transition (PENDING,STOPPING, orDELETING) 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 |
|---|---|
|
A table that generates sample data. For options, see Generate Custom Sample Data with Confluent Cloud for Apache Flink. |
|
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
|
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 |
|
|
|
|
|
|---|---|---|---|---|---|
|
Yes |
Yes |
Yes |
No |
No |
|
Yes |
Yes |
Yes |
Yes |
No |
|
Yes |
Yes |
Yes |
Yes |
No |
|
Yes |
No |
No |
No |
No |
|
Yes |
Yes |
No |
No |
No |
|
No |
Yes |
Yes |
Yes |
No |
|
No |
No |
Required |
No |
No |
|
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 |
|---|---|
|
Required. A non-empty list of column names to hash. |
|
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 onlycompute_pool_iddoesn’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 runningINSERTstatement 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_sourceandtable: 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 |
|---|---|
|
Required. |
|
Required. Where Tableflow writes the tables. Set
For more information, see Storage with Tableflow in Confluent Cloud. |
|
Optional. Topic-level Tableflow settings:
|
The adapter reconciles Tableflow on every run:
If Tableflow isn’t enabled, the adapter enables it.
If you changed
table_formatsorconfig, 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
FAILEDstate 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
storagechange 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 anddbt runfails 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 setdistributed_by, the adapter doesn’t check distribution.WITH options for
streaming_tableandstreaming_source: every option that you set must have the same value on the table. Forstreaming_source, this includesconnector.
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)toROUND(price, 4)isn’t drift. Run with--full-refreshto 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_tablemodel’sCREATE TABLEstatement has a-ddlsuffix.A
materialized_tablemodel’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 ... LIKEfor 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 withdbt runbefore 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 setaliasor projects that overridegenerate_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-failuresor thestore_failuresconfiguration. By default, dbt stores failures in adbt_test__auditschema, 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 reportPASS, 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.
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
certifiCA bundle, not your operating system’s truststore. It doesn’t trust your proxy’s certificate authority (CA) unless you set theSSL_CERT_FILEorSSL_CERT_DIRenvironment 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:
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.
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
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--certor setPIP_CERTto the same path.
Limitations
No dbt snapshots: Flink SQL doesn’t support the
MERGEandUPDATEoperations that dbt’s snapshot resource type needs to track slowly changing dimensions. This is unrelated to Flink SQL snapshot queries, which the adapter uses fortablemodels and unit tests.No incremental models: Use
materialized_tableinstead.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 TABLEorDROP TABLE. For astreaming_tablemodel, also delete its oldINSERTstatement.Not transactional: If a
dbt runfails partway through, the models that already ran stay deployed.No stored test failures:
--store-failuresand thestore_failuresconfiguration don’t work reliably. See Testing.