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.
Code for this part: scripts/dbt_run_report.py, snowflake/04_monitoring.sql. Full project: tpch_analytics on GitHub.
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_keyAnd 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=1Timings 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
| Lever | Effect | Covered in |
|---|---|---|
| Auto-suspend at 60 seconds | No paying for idle warehouses | Part 2 |
| Incremental models for large, growing tables | Process a day instead of years | Part 7 |
| Slim CI | Pull requests build a few models, not the project | Part 9 |
| Views for staging, tables for marts | No storage or rebuild cost for thin layers | Part 4 |
| Right-size the warehouse | Doubling size doubles credits per hour; only worth it if runtime roughly halves | This part |
| Clustering keys on very large tables | Better pruning; costs background maintenance | Part 7 |
| Resource monitor | A hard ceiling on monthly spend | Part 2 |
6. The runbook
When the morning alert arrives, these five situations cover most of them:
| Symptom | First check | Usual fix |
|---|---|---|
| Freshness task failed | COPY_HISTORY: did any files load? Did the source deliver? | Chase the source; once files arrive, clear the load task in Airflow |
| Load task failed | first_error_message in COPY_HISTORY | Fix or remove the bad file; rerun (COPY skips files already loaded) |
| A dbt test failed | Run the compiled test SQL (in target/compiled/) to see the failing rows | Fix the data or the model; if the rule was wrong, change the test, never just delete it |
| A model errored | The task log: usually a SQL error, a privilege, or a changed source column | Fix in a pull request (CI tests it), merge, clear the task |
| Credits jumped | Daily credits query, then the slow-query list for that day | Add 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
- 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
- CI/CD for dbt with GitHub Actions
- 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
Post a Comment