<a id="flink-dbt-reference"></a>

# dbt Adapter Reference for Confluent Cloud for Apache Flink

The `dbt-confluent` adapter runs [dbt](https://www.getdbt.com/) 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](deploy-flink-dbt.md#flink-deploy-dbt).

<a id="flink-dbt-reference-compatibility"></a>

## 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`,<br/>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<br/>compatible. Version 0.3.0 also allows Python 3.14, which isn’t<br/>supported. Pin dbt Core yourself with `pip install "dbt-core~=1.11.0"`,<br/>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`.

<a id="flink-dbt-reference-concepts"></a>

## 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`<br/>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.

<a id="flink-dbt-reference-profile"></a>

## Profile configuration

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

```yaml
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<br/>`aws` and `us-east-2`. Required unless you set `endpoint`.                                                                                                                                                                                      |
| `endpoint`                            | A full Flink endpoint URL, for a<br/>[private network](../concepts/flink-private-networking.md#flink-sql-private-networking) or another<br/>non-standard endpoint, for example<br/>`https://flink.us-east-2.aws.private.confluent.cloud`. Set either<br/>`endpoint` or `cloud_provider` and `cloud_region`, not both. |
| `global_api_key`, `global_api_secret` | A [global API key](../../security/authenticate/workload-identities/service-accounts/api-keys/overview.md#cloud-global-api-keys). Works for Flink and<br/>is required for Tableflow.                                                                                                                                   |
| `flink_api_key`, `flink_api_secret`   | A [Flink API key](generate-api-key-for-flink.md#flink-generate-api-key), which is scoped to<br/>one environment and region. Supply one complete key pair, either<br/>global or Flink. If you supply both, the adapter uses the global key.                                                                            |
| `compute_pool_id`                     | Optional. The default [compute pool](../concepts/compute-pools.md#flink-sql-compute-pools)<br/>for every model. If you omit it, Flink runs statements in the<br/>environment’s default compute pool for the region. Models can<br/>override it; see [compute_pool_id](#flink-dbt-reference-compute-pool).             |
| `execution_mode`                      | Optional. The default execution mode for statements that don’t set<br/>their own. Default: `streaming_query`. Materializations set their<br/>own mode, so you rarely need to change this.                                                                                                                             |
| `statement_name_prefix`               | Optional. The prefix for generated statement names. Default:<br/>`dbt-`. See [Statement names](#flink-dbt-reference-statement-names).                                                                                                                                                                                 |
| `statement_label`                     | Optional. A label applied to every statement that the adapter<br/>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.

<a id="flink-dbt-reference-materializations"></a>

## 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<br/>pipelines, because changes evolve in place.                                              | Streaming. Flink maintains the table.                                                                                                                |
| `streaming_table`    | Continuously maintained results with an explicit table and a<br/>separate `INSERT` statement, or adopting a pipeline that you<br/>deployed outside dbt. | Streaming. A long-running `INSERT INTO ... SELECT`.                                                                                                  |
| `streaming_source`   | A table with a declared schema, backed by the `faker` or<br/>`confluent` Flink SQL connector.                                                           | Depends on the connector. A `faker` table generates data for each<br/>statement that reads it. A `confluent` table stores data in a Kafka<br/>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<br/>`table` for a one-time result.                                            |
| `incremental`       | dbt’s batch-incremental model doesn’t map to continuous processing.<br/>Use `materialized_table`.                                              |
| `snapshot`          | dbt snapshots require `MERGE` and `UPDATE` operations that<br/>Flink SQL doesn’t support. See [Limitations](#flink-dbt-reference-limitations). |

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

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

<a id="flink-dbt-reference-materialized-table"></a>

### materialized_table

A `materialized_table` model submits a
[CREATE OR ALTER MATERIALIZED TABLE](../reference/statements/create-or-alter-materialized-table.md#flink-sql-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](../concepts/materialized-tables.md#flink-sql-materialized-tables).

```sql
{{ 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](#flink-dbt-reference-schema-drift).

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](#flink-dbt-reference-contracts).
`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.

<a id="flink-dbt-reference-start-mode"></a>

#### 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. |

<a id="flink-dbt-reference-evolution"></a>

#### 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](../concepts/materialized-tables.md#flink-sql-materialized-tables-evolving).

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.

<a id="flink-dbt-reference-contracts"></a>

#### Contracts and primary keys

To declare columns and a primary key, enforce a
[dbt model contract](https://docs.getdbt.com/reference/resource-configs/contract)
and add a model-level `primary_key` constraint in the model’s YAML file:

```yaml
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](../reference/statements/create-or-alter-materialized-table.md#flink-sql-adopt-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.

<a id="flink-dbt-reference-streaming-table"></a>

### 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`.

```sql
{{ 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](#flink-dbt-reference-schema-drift).
   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.

<a id="flink-dbt-reference-adopt"></a>

#### 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`:

```sql
{{ 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](#flink-dbt-reference-statement-names).

<a id="flink-dbt-reference-streaming-source"></a>

### 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](../reference/statements/create-table.md#flink-sql-create-table-with-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<br/>[Generate Custom Sample Data with Confluent Cloud for Apache Flink](../how-to-guides/custom-sample-data.md#flink-sql-custom-sample-data).                                                                                                            |
| `confluent` | A table backed by a Kafka topic, with the schema that the model<br/>declares. Use it when other applications produce to the topic and you<br/>want dbt to own the table’s schema. Running this model with<br/>`--full-refresh` deletes the topic and every record that other<br/>applications produced to it. |

The following model creates a `faker` table:

```sql
{{ 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:

```sql
{{ 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](#flink-dbt-reference-schema-drift) 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.

<a id="flink-dbt-reference-table"></a>

### table

A `table` model submits a `CREATE TABLE ... AS SELECT` statement in
[snapshot mode](../concepts/snapshot-queries.md#flink-sql-snapshot-queries). 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.

```sql
{{ 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](#flink-dbt-reference-schema-drift) and skips
the model. To recompute the result, run with `--full-refresh`, which drops
the table and its data and runs the query again.

<a id="flink-dbt-reference-view"></a>

### 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.

<a id="flink-dbt-reference-model-config"></a>

## 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](#flink-dbt-reference-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`,<br/>`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](../reference/statements/create-table.md#flink-sql-create-table).

<a id="flink-dbt-reference-distributed-by"></a>

### 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.

```sql
{{ 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<br/>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`.

<a id="flink-dbt-reference-statement-properties"></a>

### statement_properties

A dictionary of statement properties, the same properties that you set with
the [SET statement](../reference/statements/set.md#flink-sql-set-statement), for example
`{'sql.tables.scan.idle-timeout': '30 s'}`. Values can be strings,
integers, or Booleans.

```sql
{{ 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](#flink-dbt-reference-evolution). 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.

<a id="flink-dbt-reference-compute-pool"></a>

### 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:

```sql
{{ 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:

```sql
{{ 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](flink-rest-api.md#flink-rest-api-update-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`.

<a id="flink-dbt-reference-tableflow"></a>

### tableflow

Enables [Tableflow](../../topics/tableflow/overview.md#cloud-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](../../topics/tableflow/overview.md#cloud-tableflow) or from Flink with a snapshot
query.

```sql
{{ 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<br/>the following:<br/><br/>- `{'kind': 'Managed'}` for Confluent-managed storage.<br/>- `{'kind': 'ByobAws', 'bucket_name': '...', 'provider_integration_id': '...'}`<br/>  for your own Amazon S3 bucket.<br/>- `{'kind': 'AzureDataLakeStorageGen2', 'storage_account_name': '...', 'container_name': '...', 'provider_integration_id': '...'}`<br/>  for your own Azure Data Lake Storage Gen2 container.<br/>- `{'kind': 'GoogleCloudStorage', 'bucket_name': '...', 'provider_integration_id': '...'}`<br/>  for your own Google Cloud Storage bucket.<br/><br/>For more information, see [Storage with Tableflow in Confluent Cloud](../../topics/tableflow/concepts/tableflow-storage.md#tableflow-storage). |
| `config`        | Optional. Topic-level Tableflow settings:<br/><br/>- `retention_ms`: How long to keep table snapshots, in<br/>  milliseconds.<br/>- `data_retention_ms`: How long to keep table data, in<br/>  milliseconds.<br/>- `error_handling`: What to do when a record can’t be<br/>  materialized. `{'mode': 'SUSPEND'}` suspends Tableflow, which is<br/>  the default. `{'mode': 'SKIP'}` skips the record.<br/>  `{'mode': 'LOG', 'target': '<topic>'}` writes the record to a<br/>  dead letter queue (DLQ) topic, `error_log` by default, and<br/>  continues. For more information, see<br/>  [Error-handling mode](../../topics/tableflow/operate/configure-tableflow.md#tableflow-configure-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](../../topics/tableflow/operate/tableflow-rbac.md#tableflow-rbac).
- 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](../../topics/tableflow/concepts/tableflow-storage.md#tableflow-storage-disable). To rebuild
  immediately instead, run with `--full-refresh`, which also deletes the
  topic’s data.

<a id="flink-dbt-reference-validation"></a>

### 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](#flink-dbt-reference-model-config),
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:

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

<a id="flink-dbt-reference-schema-drift"></a>

## 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.

```sql
{{ 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.

<a id="flink-dbt-reference-statement-names"></a>

## 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](#flink-dbt-reference-adopt), 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.

<a id="flink-deploy-dbt-sources"></a>

## 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.

```yaml
# 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() }}`:

```sql
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.

<a id="flink-dbt-reference-testing"></a>

## Testing

The adapter supports dbt unit tests and data tests. For a worked example, see
[Step 6: Test the pipeline](deploy-flink-dbt.md#flink-deploy-dbt-test).

- 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.

<a id="flink-dbt-reference-deploy"></a>

## Manage deployments

<a id="flink-deploy-dbt-dependencies"></a>

### 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](https://docs.getdbt.com/reference/node-selection/syntax)
to list or deploy part of the graph:

```bash
# 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](../concepts/schema-statement-evolution.md#flink-sql-schema-and-statement-evolution) and
[Carry-over Offsets in Confluent Cloud for Apache Flink](carry-over-offsets.md#flink-sql-carry-over-offsets).

<a id="flink-deploy-dbt-cicd"></a>

<a id="flink-dbt-reference-cicd"></a>

### 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](#flink-deploy-dbt-multi-env). 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.

```yaml
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](https://docs.github.com/en/actions/security-for-github-actions/security-guides/using-secrets-in-github-actions),
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](https://docs.github.com/en/actions/deployment/targeting-different-environments/using-environments-for-deployment)
with required reviewers.

<a id="flink-deploy-dbt-multi-env"></a>

### 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`:

```yaml
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
```

```bash
dbt run --target prod
```

<a id="flink-dbt-reference-troubleshooting"></a>

## 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:

```bash
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](#flink-dbt-reference-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](../concepts/compute-pools.md#flink-sql-compute-pools) and
[Grant Role-Based Access in Confluent Cloud for Apache Flink](flink-rbac.md#flink-rbac).

### 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](../concepts/flink-private-networking.md#flink-sql-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.

<a id="flink-dbt-reference-troubleshooting-tls"></a>

### 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:
   ```bash
   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:
   ```bash
   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.

<a id="flink-dbt-reference-limitations"></a>

## Limitations

- **No dbt snapshots**: Flink SQL doesn’t support the `MERGE` and
  `UPDATE` operations that dbt’s
  [snapshot](https://docs.getdbt.com/docs/build/snapshots) resource
  type needs to track slowly changing dimensions. This is unrelated to
  Flink SQL [snapshot queries](../concepts/snapshot-queries.md#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](#flink-dbt-reference-concepts).
- **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](#flink-dbt-reference-testing).

## Related content

- [Build a Streaming Pipeline with dbt and Confluent Cloud for Apache Flink](deploy-flink-dbt.md#flink-deploy-dbt)
- [Materialized Tables in Confluent Cloud for Apache Flink](../concepts/materialized-tables.md#flink-sql-materialized-tables)
- [Flink SQL Development Lifecycle in Confluent Cloud for Apache Flink](development-lifecycle.md#flink-development-lifecycle)
- [Move SQL Statements to Production in Confluent Cloud for Apache Flink](best-practices.md#flink-sql-best-practices-for-statements)
- [Schema and Statement Evolution with Confluent Cloud for Apache Flink](../concepts/schema-statement-evolution.md#flink-sql-schema-and-statement-evolution)
- [Grant Role-Based Access in Confluent Cloud for Apache Flink](flink-rbac.md#flink-rbac)
- GitHub: [dbt-confluent repository](https://github.com/confluentinc/dbt-confluent)

#### NOTE
This website includes content developed at the [Apache Software Foundation](https://www.apache.org/)
under the terms of the [Apache License v2](https://www.apache.org/licenses/LICENSE-2.0.html).
