Skip to main content

Running dbt, Snowflake and Airflow in Production (Part 10)

The pipeline is built, tested in CI and scheduled. This last part is about the months after go-live: knowing within minutes when a run fails, knowing within a week what the pipeline costs and which models are slow, and having a short runbook for the mornings when the dashboard is wrong. Everything here uses tools already in the stack: Airflow, dbt's own output files and Snowflake's usage views.

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

Code for this part: scripts/dbt_run_report.py, snowflake/04_monitoring.sql. Full project: tpch_analytics on GitHub.

SignalsChecked byActionAirflow task statesfailed tasks after retriesdbt run_results.jsontest failures, slow modelsSource freshnessRAW stopped updatingSnowflake ACCOUNT_USAGEcredits, slow queries, loadsAlert in minutesfailure callback (part 8)Weekly reviewcost and performance SQLRunbookfind cause, fix, rerun taskTuneresize, incremental, cluster

1. Know when a run fails

Part 8 already alerts on any task that fails after its retries. dbt adds a second signal: after every run it writes target/run_results.json with the status and duration of every model and test. A short script turns that into a summary and a non-zero exit code when anything failed, useful as a final Airflow task, in CI, or after a manual run:

scripts/dbt_run_report.py"""Summarise a dbt run from target/run_results.json: failures and slowest nodes."""
import json
import sys
from pathlib import Path

path = Path(sys.argv[1] if len(sys.argv) > 1 else "target/run_results.json")
results = json.loads(path.read_text())["results"]

problems = [r for r in results if r["status"] not in ("success", "pass")]
slowest = sorted(results, key=lambda r: r["execution_time"], reverse=True)[:5]

print(f"{len(results)} nodes, {len(problems)} not successful")
for r in problems:
    print(f"  {r['status'].upper():6} {r['unique_id']}  {(r.get('message') or '')[:80]}")
print("Slowest nodes:")
for r in slowest:
    print(f"  {r['execution_time']:6.2f}s  {r['unique_id'].split('.')[2]}")

sys.exit(1 if any(r["status"] in ("error", "fail") for r in problems) else 0)

After a healthy run:

34 nodes, 0 not successful
Slowest nodes:
    0.22s  unique_stg_tpch__line_items_order_item_key
    0.17s  snap_customers
    0.13s  fct_orders
    0.11s  fct_order_items
    0.10s  not_null_stg_tpch__line_items_order_item_key

And after I deliberately tightened the reconciliation test from part 6 so it would fail:

6 nodes, 2 not successful
  FAIL   test.tpch_analytics.assert_order_total_matches_line_items  Got 66465 results, configured to fail if != 0
  SKIPPED model.tpch_analytics.agg_monthly_revenue
Slowest nodes:
    0.22s  fct_orders
    0.03s  assert_order_total_matches_line_items
    ...
exit=1

Timings are from the local DuckDB test build; on Snowflake expect seconds to minutes per model. The useful part is the ranking: the same few models usually dominate.

2. Know what it costs

Snowflake bills compute per second of warehouse time, so cost questions are warehouse questions. The SNOWFLAKE.ACCOUNT_USAGE views answer them; they lag real time by up to a few hours, which is fine for a weekly review. You need a role with access to the SNOWFLAKE database (an admin can grant IMPORTED PRIVILEGES on it). All queries below are in snowflake/04_monitoring.sql.

SQL · Snowflake-- Credits per warehouse, last 30 days
SELECT warehouse_name,
       ROUND(SUM(credits_used), 2) AS credits
FROM SNOWFLAKE.ACCOUNT_USAGE.WAREHOUSE_METERING_HISTORY
WHERE start_time >= DATEADD(DAY, -30, CURRENT_TIMESTAMP())
GROUP BY warehouse_name
ORDER BY credits DESC;

-- Credits per day for the dbt warehouse: spot the day something changed
SELECT DATE_TRUNC('DAY', start_time) AS day,
       ROUND(SUM(credits_used), 2)   AS credits
FROM SNOWFLAKE.ACCOUNT_USAGE.WAREHOUSE_METERING_HISTORY
WHERE warehouse_name = 'TRANSFORM_WH'
  AND start_time >= DATEADD(DAY, -30, CURRENT_TIMESTAMP())
GROUP BY day
ORDER BY day;

Because part 2 gave each workload its own warehouse, the first query already splits cost into loading, transformation and reporting. A sudden step in the daily series usually lines up with a merge: compare the date with your Git history.

3. Find the slow queries

Part 4 set +query_tag: dbt_tpch_analytics, so every query dbt sends is labelled. That makes dbt's queries easy to isolate in the query history:

SQL · Snowflake-- Slowest dbt queries this week (dbt sets QUERY_TAG via +query_tag in dbt_project.yml)
SELECT start_time,
       ROUND(total_elapsed_time / 1000, 1)          AS seconds,
       ROUND(bytes_scanned / POWER(1024, 3), 2)     AS gb_scanned,
       partitions_scanned,
       partitions_total,
       LEFT(query_text, 120)                        AS query_start
FROM SNOWFLAKE.ACCOUNT_USAGE.QUERY_HISTORY
WHERE query_tag = 'dbt_tpch_analytics'
  AND start_time >= DATEADD(DAY, -7, CURRENT_TIMESTAMP())
ORDER BY total_elapsed_time DESC
LIMIT 20;

The two partition columns are the most useful: when partitions_scanned is close to partitions_total on a big table, the query reads the whole table. That points at a missing incremental filter (part 7) or a filter Snowflake cannot use for pruning, for example a function wrapped around the date column.

4. Check the loads

SQL · Snowflake-- Loads that did not fully succeed in the last 7 days
SELECT table_name, file_name, status, row_count, row_parsed, first_error_message, last_load_time
FROM SNOWFLAKE.ACCOUNT_USAGE.COPY_HISTORY
WHERE last_load_time >= DATEADD(DAY, -7, CURRENT_TIMESTAMP())
  AND status <> 'Loaded'
ORDER BY last_load_time DESC;

5. Cost levers, in the order to try them

LeverEffectCovered in
Auto-suspend at 60 secondsNo paying for idle warehousesPart 2
Incremental models for large, growing tablesProcess a day instead of yearsPart 7
Slim CIPull requests build a few models, not the projectPart 9
Views for staging, tables for martsNo storage or rebuild cost for thin layersPart 4
Right-size the warehouseDoubling size doubles credits per hour; only worth it if runtime roughly halvesThis part
Clustering keys on very large tablesBetter pruning; costs background maintenancePart 7
Resource monitorA hard ceiling on monthly spendPart 2

6. The runbook

When the morning alert arrives, these five situations cover most of them:

SymptomFirst checkUsual fix
Freshness task failedCOPY_HISTORY: did any files load? Did the source deliver?Chase the source; once files arrive, clear the load task in Airflow
Load task failedfirst_error_message in COPY_HISTORYFix or remove the bad file; rerun (COPY skips files already loaded)
A dbt test failedRun the compiled test SQL (in target/compiled/) to see the failing rowsFix the data or the model; if the rule was wrong, change the test, never just delete it
A model erroredThe task log: usually a SQL error, a privilege, or a changed source columnFix in a pull request (CI tests it), merge, clear the task
Credits jumpedDaily credits query, then the slow-query list for that dayAdd or repair an incremental filter, or revert the change

Two habits make every incident shorter: never edit production tables by hand (fix the code or the data and rerun), and write one line in the runbook each time something new breaks.

What we built

Over ten parts, the TPC-H data went from a sample database to a production pipeline:

  • A Snowflake account with separate warehouses, functional roles, a key-pair service user and a credit budget.
  • Files loaded idempotently with stages and COPY INTO, with audit columns on every row.
  • A dbt project with staging, intermediate and mart layers, a star schema and a revenue aggregate.
  • 20 data tests, source freshness checks and generated documentation.
  • An incremental fact table and a customer history snapshot.
  • A daily Airflow DAG where every model and test is its own task.
  • CI that builds only changed models, in isolated schemas, on every pull request.
  • Monitoring, cost queries and a runbook.

Where to go next: put a semantic layer on top of the marts so every tool shares one definition of revenue, or let people ask questions of them in plain English with an AI analytics assistant. For how all of this fits in a wider platform, see the modern data engineering architecture guide.

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
  9. CI/CD for dbt with GitHub Actions
  10. Running the pipeline in production (this post)

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