Skip to main content

Orchestrating dbt on Snowflake with Airflow and Cosmos (Part 8)

The pipeline works when you run it by hand. Production needs it to run every day, in the right order, with retries when Snowflake has a hiccup and an alert when something really breaks. That is Airflow's job. In this part we write a DAG that loads new files with COPY INTO and then runs the whole dbt project, with every model and its tests as separate Airflow tasks, using Astronomer's open-source Cosmos library.

Part 8 of 10 in the series Build a modern data pipeline with Snowflake, dbt and Airflow. Previous: Part 7, Incremental models and snapshots. Next: Part 9, CI/CD for dbt with GitHub Actions.

Code for this part: airflow/dags/tpch_daily.py, airflow/requirements.txt. Full project: tpch_analytics on GitHub.

DAG tpch_daily · 02:00 UTC daily · 23 tasksload_raw_ordersCOPY INTO transform (Cosmos DbtTaskGroup)orders sourcefreshnessstg_tpch__orders.runstg_tpch__orders.testlineitem sourcefreshnessstg_…line_items.runstg_…line_items.testfct_order_items.runfct_order_items.test… and so on for every model, snapshot and teston failureretry ×2 after 5 min,then alert

New to Airflow? Apache Airflow for data engineers covers DAGs, tasks, scheduling and retries.

Two ways to run dbt from Airflow

One task: dbt buildCosmos: one task per model
SetupA BashOperator, five linesA Python library, about twenty lines
When a model failsThe whole task fails; retry reruns everythingOnly that model fails; retry reruns only that model
VisibilityOne log containing every modelEach model and its tests in the Airflow grid
Partial rerunsManual --select commandsClear one task in the UI
Best forSmall projects, quick startProjects you operate every day

We use Cosmos. It reads the dbt project, turns every model, snapshot, test and source freshness check into an Airflow task, and wires the dependencies from ref() and source(), so the DAG always matches the dbt project without being edited by hand.

Step 1: install

Airflow has many dependencies, so install it with its official constraints file, then add Cosmos and the Snowflake provider. In testing, Airflow 3.1 failed to start with SQLAlchemy 2.1, so pin it below 2.1:

bashpip install "apache-airflow==3.1.8" \
  --constraint "https://raw.githubusercontent.com/apache/airflow/constraints-3.1.8/constraints-3.12.txt"
pip install "astronomer-cosmos==1.15.1" "apache-airflow-providers-snowflake==6.18.0" "sqlalchemy<2.1"

Install dbt in its own virtual environment so its dependencies never conflict with Airflow's. Cosmos only needs the path to the dbt executable:

bashpython -m venv /opt/airflow/dbt_venv
/opt/airflow/dbt_venv/bin/pip install -r /opt/airflow/dbt/tpch_analytics/requirements.txt

Step 2: one Snowflake connection

Both the load and dbt use the AIRFLOW_SVC user from part 2, with its private key. Create the connection in the Airflow UI or as an environment variable. The password field holds the key passphrase, not a Snowflake password:

bashexport AIRFLOW_CONN_SNOWFLAKE_DEFAULT='{
  "conn_type": "snowflake",
  "login": "AIRFLOW_SVC",
  "password": "<private key passphrase>",
  "schema": "TPCH",
  "extra": {
    "account": "<your_account_identifier>",
    "warehouse": "LOADING_WH",
    "database": "RAW",
    "role": "LOADER",
    "private_key_file": "/opt/airflow/keys/rsa_key.p8"
  }
}'

The connection's defaults (LOADER role, LOADING_WH) suit the load. For dbt, the DAG overrides role, warehouse and database, so transformations run as TRANSFORMER on TRANSFORM_WH.

Step 3: the DAG

airflow/dags/tpch_daily.py"""Daily TPC-H pipeline: load new raw data into Snowflake, then build and test with dbt."""
import os
from datetime import datetime, timedelta
from pathlib import Path

from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from airflow.sdk import dag
from cosmos import DbtTaskGroup, ExecutionConfig, ProfileConfig, ProjectConfig, RenderConfig
from cosmos.constants import LoadMode, SourceRenderingBehavior, TestBehavior
from cosmos.profiles import SnowflakeEncryptedPrivateKeyFilePemProfileMapping

DBT_PROJECT_DIR = Path(os.getenv("DBT_PROJECT_DIR", "/opt/airflow/dbt/tpch_analytics"))
DBT_EXECUTABLE = os.getenv("DBT_EXECUTABLE", "/opt/airflow/dbt_venv/bin/dbt")

profile_config = ProfileConfig(
    profile_name="tpch_analytics",
    target_name="prod",
    profile_mapping=SnowflakeEncryptedPrivateKeyFilePemProfileMapping(
        conn_id="snowflake_default",
        profile_args={
            "database": "ANALYTICS",
            "schema": "PROD",
            "warehouse": "TRANSFORM_WH",
            "role": "TRANSFORMER",
        },
    ),
)


def notify_failure(context):
    """Called when a task fails after its retries. Swap the print for Slack or email."""
    ti = context["task_instance"]
    print(f"FAILED: {ti.dag_id}.{ti.task_id} for {context['logical_date']}")


@dag(
    schedule="0 2 * * *",              # every day at 02:00 UTC
    start_date=datetime(2026, 1, 1),
    catchup=False,
    max_active_runs=1,
    dagrun_timeout=timedelta(hours=2),
    default_args={
        "owner": "data-eng",
        "retries": 2,
        "retry_delay": timedelta(minutes=5),
        "on_failure_callback": notify_failure,
    },
    tags=["tpch", "dbt", "snowflake"],
)
def tpch_daily():

    load_raw = SQLExecuteQueryOperator(
        task_id="load_raw_orders",
        conn_id="snowflake_default",
        sql="sql/load_raw_orders.sql",
        split_statements=True,
    )

    transform = DbtTaskGroup(
        group_id="transform",
        project_config=ProjectConfig(
            DBT_PROJECT_DIR,
            manifest_path=DBT_PROJECT_DIR / "target" / "manifest.json",
        ),
        profile_config=profile_config,
        execution_config=ExecutionConfig(dbt_executable_path=DBT_EXECUTABLE),
        render_config=RenderConfig(
            load_method=LoadMode.DBT_MANIFEST,
            test_behavior=TestBehavior.AFTER_EACH,
            source_rendering_behavior=SourceRenderingBehavior.WITH_TESTS_OR_FRESHNESS,
        ),
        operator_args={"install_deps": True},
    )

    load_raw >> transform


tpch_daily()

What each setting does

  • SQLExecuteQueryOperator runs the COPY statements from part 3 (sql/load_raw_orders.sql, next to the DAG). Because COPY skips files it has already loaded, a retry is always safe.
  • Profile mapping: Cosmos builds dbt's profiles.yml from the Airflow connection at run time, so the key never lives in the dbt project. target_name="prod" activates the clean schema names from part 4.
  • LoadMode.DBT_MANIFEST: Cosmos builds the task graph from dbt's manifest.json instead of running dbt ls every time the scheduler parses the DAG. In testing, the dbt ls approach took about 6 seconds per parse; reading the manifest is near-instant. Your deployment must run dbt deps && dbt parse to produce the manifest (part 9).
  • TestBehavior.AFTER_EACH: each model's tests run right after it, and downstream models wait for them, the same protection dbt build gives, now visible per task.
  • SourceRenderingBehavior.WITH_TESTS_OR_FRESHNESS adds a task per source with freshness configured (orders and lineitem, from part 6). If the load delivered nothing new, the freshness check fails before any model builds on stale data.
  • max_active_runs=1 stops two days' runs overlapping, which matters for the incremental model from part 7.

Step 4: check the DAG

Before deploying, confirm Airflow can import the file and see what Cosmos generated:

pythonfrom airflow.models.dagbag import DagBag

db = DagBag(dag_folder="airflow/dags", include_examples=False)
print("import errors:", db.import_errors)
dag = db.get_dag("tpch_daily")
print("tasks:", len(dag.tasks))
import errors: {}
tasks: 23
  transform.tpch_orders_source <- ['load_raw_orders']
  transform.tpch_lineitem_source <- ['load_raw_orders']
  transform.stg_tpch__orders.run <- ['transform.tpch_orders_source']
  transform.stg_tpch__orders.test <- ['transform.stg_tpch__orders.run']
  ...

23 tasks: the load, two source freshness checks, and a run task (plus a test task where tests exist) for each model and the snapshot. I also ran the whole task group end to end with airflow dags test on the local DuckDB version of the project: every task succeeded, freshness passed for both sources, and each model's tests ran before its dependants started.

Alerts that reach someone

The notify_failure callback fires only after a task has used up its retries, so a brief network error does not wake anyone. Replace the print with a real channel; the Slack provider's webhook notifier and Airflow's email settings are the common choices. Include the DAG, task and date in the message, and a link to the log.

Common problems

SymptomLikely cause
DAG does not appear, import error mentions manifest.jsonThe deployment did not run dbt parse, so target/manifest.json is missing
Tasks fail with "insufficient privileges"The dbt tasks are running as LOADER; check profile_args sets role TRANSFORMER
Airflow fails to start after installing CosmosSQLAlchemy 2.1 was installed; install with the constraints file and sqlalchemy<2.1
Freshness task fails every morningThe source really stopped sending files, or the load ran before the files arrived; check COPY_HISTORY (part 3)

Next

Airflow now runs production every morning. But how do changes get there safely? In part 9 we set up CI/CD: every pull request builds and tests only the models it changed, in its own Snowflake schema, before anyone merges.

The series

  1. The stack and what we will build
  2. Setting up Snowflake for dbt
  3. Loading raw data into Snowflake
  4. Your first dbt project on Snowflake
  5. Modeling with dbt: staging to marts
  6. Testing and documenting with dbt
  7. Incremental models and snapshots
  8. Orchestrating dbt with Airflow (this post)
  9. CI/CD for dbt with GitHub Actions
  10. Running the pipeline in production

All parts: Data Pipeline Series. Every file shown in this series is in the companion project tpch_analytics on GitHub; clone it to follow along, or run it locally on DuckDB without a Snowflake account.

Comments