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.
Code for this part: airflow/dags/tpch_daily.py, airflow/requirements.txt. Full project: tpch_analytics on GitHub.
New to Airflow? Apache Airflow for data engineers covers DAGs, tasks, scheduling and retries.
Two ways to run dbt from Airflow
One task: dbt build | Cosmos: one task per model | |
|---|---|---|
| Setup | A BashOperator, five lines | A Python library, about twenty lines |
| When a model fails | The whole task fails; retry reruns everything | Only that model fails; retry reruns only that model |
| Visibility | One log containing every model | Each model and its tests in the Airflow grid |
| Partial reruns | Manual --select commands | Clear one task in the UI |
| Best for | Small projects, quick start | Projects 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.txtStep 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.ymlfrom 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'smanifest.jsoninstead of runningdbt lsevery time the scheduler parses the DAG. In testing, thedbt lsapproach took about 6 seconds per parse; reading the manifest is near-instant. Your deployment must rundbt deps && dbt parseto produce the manifest (part 9).TestBehavior.AFTER_EACH: each model's tests run right after it, and downstream models wait for them, the same protectiondbt buildgives, now visible per task.SourceRenderingBehavior.WITH_TESTS_OR_FRESHNESSadds 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=1stops 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
| Symptom | Likely cause |
|---|---|
| DAG does not appear, import error mentions manifest.json | The 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 Cosmos | SQLAlchemy 2.1 was installed; install with the constraints file and sqlalchemy<2.1 |
| Freshness task fails every morning | The 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
- The stack and what we will build
- Setting up Snowflake for dbt
- Loading raw data into Snowflake
- Your first dbt project on Snowflake
- Modeling with dbt: staging to marts
- Testing and documenting with dbt
- Incremental models and snapshots
- Orchestrating dbt with Airflow (this post)
- CI/CD for dbt with GitHub Actions
- 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
Post a Comment