Skip to main content

dbt Modeling on Snowflake: Staging, Marts and Star Schemas (Part 5)

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.

Part 5 of 10 in the series Build a modern data pipeline with Snowflake, dbt and Airflow. Previous: Part 4, Your first dbt project on Snowflake. Next: Part 6, Testing and documenting with dbt.

Code for this part: models/ (staging, intermediate, marts). Full project: tpch_analytics on GitHub.

Sources (RAW)orderslineitemcustomernationregionpartStaging (views)stg_tpch__ordersstg_tpch__line_itemsstg_tpch__customersstg_tpch__nationsstg_tpch__regionsstg_tpch__partsIntermediateint_order_items__enrichedMarts (tables)fct_order_itemsfct_ordersdim_customersdim_partsAggregateagg_monthly_revenue

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 renamed

The 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_key

Notice 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.

ModelTypeGrain
fct_order_itemsFactOne row per order line
fct_ordersFactOne row per order
dim_customersDimensionOne row per customer
dim_partsDimensionOne row per part
agg_monthly_revenueAggregateMonth × 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_key

The 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_key

Finally, 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, 3

Build it

The + before a model name selects it and everything upstream:

bashdbt run --select +agg_monthly_revenue

Results 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.65

And 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.58

Conventions worth adopting

ConventionExampleWhy
Prefix by layerstg_, int_, fct_, dim_, agg_Purpose is obvious from the name
Double underscore before the entitystg_tpch__ordersSeparates source system from table
Keys end in _keycustomer_keyConsistent joins
Booleans start with is_ / has_is_late_deliveryReadable filters
Import CTEs at the topwith orders as (select * from {{ ref(...) }})Dependencies visible at a glance
One model, one grainfct_orders = one row per orderPrevents 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

  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 (this post)
  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

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