With dbt connected to Snowflake, the real work is deciding how to organise the transformations. The structure most dbt projects converge on has three layers: staging cleans each source table, intermediate combines them and applies business logic, and marts publish facts and dimensions people actually query. This part builds all three for the TPC-H data and ends with a star schema and a revenue table ready for a dashboard.
Code for this part: models/ (staging, intermediate, marts). Full project: tpch_analytics on GitHub.
Layer 1: staging, one model per source table
Staging models follow strict rules: one model per source table, rename and cast columns, decode codes, and no joins or aggregations. Everything downstream reads staging models, never the raw tables, so a change in a source system is fixed in one place. Part 4 built stg_tpch__orders; the line items model also creates a surrogate key, because TPC-H identifies a line by two columns (order and line number):
models/staging/tpch/stg_tpch__line_items.sqlwith source as (
select * from {{ source('tpch', 'lineitem') }}
),
renamed as (
select
{{ dbt_utils.generate_surrogate_key(['l_orderkey', 'l_linenumber']) }} as order_item_key,
l_orderkey as order_key,
l_linenumber as line_number,
l_partkey as part_key,
l_suppkey as supplier_key,
l_quantity as quantity,
l_extendedprice as extended_price,
l_discount as discount_rate,
l_tax as tax_rate,
l_returnflag as return_flag,
l_linestatus as line_status,
l_shipdate as ship_date,
l_commitdate as commit_date,
l_receiptdate as receipt_date,
l_shipmode as ship_mode,
_loaded_at
from source
)
select * from renamedThe smaller staging models are just a renamed select:
models/staging/tpch/stg_tpch__customers.sqlselect
c_custkey as customer_key,
c_name as customer_name,
c_address as customer_address,
c_nationkey as nation_key,
c_phone as phone_number,
c_acctbal as account_balance,
c_mktsegment as market_segment
from {{ source('tpch', 'customer') }}Layer 2: intermediate, where business logic lives
Revenue in TPC-H is not stored; it is calculated from price, discount and tax. Calculating it once, in one model, guarantees that every report uses the same definition. The intermediate model joins each line to its order and derives the amounts:
models/intermediate/int_order_items__enriched.sqlwith line_items as (
select * from {{ ref('stg_tpch__line_items') }}
),
orders as (
select * from {{ ref('stg_tpch__orders') }}
)
select
line_items.order_item_key,
line_items.order_key,
line_items.line_number,
orders.customer_key,
line_items.part_key,
line_items.supplier_key,
orders.order_date,
orders.order_status,
line_items.ship_date,
line_items.commit_date,
line_items.receipt_date,
line_items.ship_mode,
line_items.return_flag,
line_items.quantity,
line_items.extended_price as gross_amount,
line_items.extended_price * line_items.discount_rate as discount_amount,
line_items.extended_price * (1 - line_items.discount_rate) as net_amount,
line_items.extended_price * (1 - line_items.discount_rate)
* line_items.tax_rate as tax_amount,
line_items.extended_price * (1 - line_items.discount_rate)
* (1 + line_items.tax_rate) as total_amount,
line_items.receipt_date > line_items.commit_date as is_late_delivery,
line_items._loaded_at
from line_items
inner join orders
on line_items.order_key = orders.order_keyNotice ref() instead of table names. ref('stg_tpch__orders') resolves to the right schema for whoever is running dbt (your dev schema, a CI schema or production) and tells dbt that this model depends on that one. dbt uses those references to build the lineage graph above and to run models in the right order.
Layer 3: marts, a star schema
Marts are what BI tools and analysts query, so they are materialised as tables. We build a classic star schema: fact tables of events with numeric measures, surrounded by dimension tables that describe them.
| Model | Type | Grain |
|---|---|---|
| fct_order_items | Fact | One row per order line |
| fct_orders | Fact | One row per order |
| dim_customers | Dimension | One row per customer |
| dim_parts | Dimension | One row per part |
| agg_monthly_revenue | Aggregate | Month × region × market segment |
fct_order_items is simply the enriched intermediate model published as a table (part 7 turns it into an incremental model):
models/marts/core/fct_order_items.sql (for now)select * from {{ ref('int_order_items__enriched') }}fct_orders rolls the lines up to one row per order and keeps the header's own total for reconciliation (part 6 tests that the two agree):
models/marts/core/fct_orders.sqlwith order_items as (
select * from {{ ref('fct_order_items') }}
),
orders as (
select * from {{ ref('stg_tpch__orders') }}
),
order_totals as (
select
order_key,
count(*) as item_count,
sum(quantity) as total_quantity,
sum(gross_amount) as gross_amount,
sum(discount_amount) as discount_amount,
sum(net_amount) as net_amount,
sum(tax_amount) as tax_amount,
sum(total_amount) as total_amount,
sum(case when is_late_delivery then 1 else 0 end) as late_item_count
from order_items
group by order_key
)
select
orders.order_key,
orders.customer_key,
orders.order_date,
orders.order_status,
orders.order_priority,
orders.total_price,
order_totals.item_count,
order_totals.total_quantity,
order_totals.gross_amount,
order_totals.discount_amount,
order_totals.net_amount,
order_totals.tax_amount,
order_totals.total_amount,
order_totals.late_item_count
from orders
inner join order_totals
on orders.order_key = order_totals.order_keyThe customer dimension flattens the nation and region hierarchy into the dimension (no extra joins for report users) and adds order history:
models/marts/core/dim_customers.sqlwith customers as (
select * from {{ ref('stg_tpch__customers') }}
),
nations as (
select * from {{ ref('stg_tpch__nations') }}
),
regions as (
select * from {{ ref('stg_tpch__regions') }}
),
order_summary as (
select
customer_key,
min(order_date) as first_order_date,
max(order_date) as most_recent_order_date,
count(*) as lifetime_orders
from {{ ref('stg_tpch__orders') }}
group by customer_key
)
select
customers.customer_key,
customers.customer_name,
customers.market_segment,
customers.account_balance,
nations.nation_name,
regions.region_name,
order_summary.first_order_date,
order_summary.most_recent_order_date,
coalesce(order_summary.lifetime_orders, 0) as lifetime_orders
from customers
left join nations
on customers.nation_key = nations.nation_key
left join regions
on nations.region_key = regions.region_key
left join order_summary
on customers.customer_key = order_summary.customer_keyFinally, an aggregate for the revenue dashboard. Pre-aggregating keeps the dashboard fast and guarantees it uses the same revenue definition as everything else:
models/marts/finance/agg_monthly_revenue.sqlselect
date_trunc('month', orders.order_date) as order_month,
customers.region_name,
customers.market_segment,
count(*) as order_count,
sum(orders.net_amount) as net_revenue,
sum(orders.total_amount) as gross_revenue_incl_tax
from {{ ref('fct_orders') }} as orders
inner join {{ ref('dim_customers') }} as customers
on orders.customer_key = customers.customer_key
group by 1, 2, 3Build it
The + before a model name selects it and everything upstream:
bashdbt run --select +agg_monthly_revenueResults in this post come from the local DuckDB build of the project at a small scale factor (75,000 orders), used to test every snippet. On Snowflake's TPCH_SF1 data the shape is the same with about 20 times more rows.
The top five region and segment combinations for June 1998:
SQLselect order_month, region_name, market_segment, order_count, net_revenue
from dbt_aravind_marts.agg_monthly_revenue
where order_month = '1998-06-01'
order by net_revenue desc
limit 5;order_month region_name market_segment order_count net_revenue
1998-06 ASIA BUILDING 51 6887050.74
1998-06 MIDDLE EAST AUTOMOBILE 48 6739767.62
1998-06 AMERICA FURNITURE 42 6215626.95
1998-06 EUROPE MACHINERY 39 6042369.26
1998-06 ASIA MACHINERY 46 5950767.65And the first orders in fct_orders. Note that our calculated total_amount is a few cents off the source's total_price: the source rounds each line. Part 6 turns that observation into a test.
order_key order_status item_count net_amount total_amount total_price
1 open 6 176206.53 182897.05 182897.00
2 open 1 46143.40 48450.57 48450.57
3 fulfilled 6 217606.76 220604.61 220604.58Conventions worth adopting
| Convention | Example | Why |
|---|---|---|
| Prefix by layer | stg_, int_, fct_, dim_, agg_ | Purpose is obvious from the name |
| Double underscore before the entity | stg_tpch__orders | Separates source system from table |
| Keys end in _key | customer_key | Consistent joins |
| Booleans start with is_ / has_ | is_late_delivery | Readable filters |
| Import CTEs at the top | with orders as (select * from {{ ref(...) }}) | Dependencies visible at a glance |
| One model, one grain | fct_orders = one row per order | Prevents double counting |
Next
The models build, but nothing yet proves they are right. In part 6 we add tests for keys, relationships and business rules, check that raw data is fresh, and generate a documentation site with this lineage graph built in.
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 (this post)
- 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
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