Pipeline Stage
Create a new transform stage with idempotency guarantees, a schema contract, and tests.
Inputs
- Stage name (snake_case, e.g.,
customer_ltv) - Input tables/models (list of upstream stage names or raw tables)
- Output table name (usually matches stage name)
- Grain (the primary key or unique key, e.g.,
customer_id,(order_id, date)) - Layer (
staging,intermediate,marts)
Steps
Create the transform file
For dbt:
models/{layer}/{stage_name}.sql{{ config( materialized='table', unique_key='{grain}' ) }} select {grain}, -- TODO: add business logic current_timestamp as updated_at from {{ ref('{input_table}') }}For pandas ETL:
pipelines/transforms/{stage_name}/transform.pydef run(df: pd.DataFrame) -> pd.DataFrame: """Transform {input_table} → {stage_name}.""" # TODO: add business logic return dfDefine the schema contract Create
models/{layer}/schema/{stage_name}.yaml(dbt) orpipelines/transforms/{stage_name}/schema.py:- name: {stage_name} columns: - name: {grain} tests: - unique - not_nullEvery non-nullable column must have
not_nulltest; every unique key must haveuniquetest.Add idempotency logic
- For
materialized='table': dbt handles full replacement — no extra work. - For incremental models: use
is_incremental()filter onupdated_ator an event timestamp. - For pandas: the output must be deterministic given the same input; add a dedup step on
{grain}.
- For
Write tests
tests/transforms/test_{stage_name}.pyRequired tests:
- Input fixture → expected output shape (column names, row count)
- Idempotency: running twice produces identical output
- Null check: no nulls in required columns after transform
Register in the pipeline DAG Add the stage after its upstream dependencies:
# dags/pipeline.py {stage_name}_task = DbtRunOperator( task_id="{stage_name}", models="{stage_name}", ) {upstream_task} >> {stage_name}_taskRun locally
dbt run --select {stage_name} dbt test --select {stage_name} # or for pandas: python -m pytest tests/transforms/test_{stage_name}.py -v
Conventions
- Layer hierarchy:
raw→staging→intermediate→marts - Never skip a layer (e.g., don't read from
rawin amartsmodel) - All stages have at least one
unique+not_nulltest on the grain column - Incremental models use
updated_atas the watermark; add it to every model
Edge Cases
- Fan-out (multiple downstream consumers): Create the stage at the
intermediatelayer; let downstreammartsmodels reference it. - Slowly changing dimension (SCD): Use dbt's
snapshotmaterialization or addvalid_from/valid_tocolumns manually. - Cross-database join: Materialize both inputs to the same database first, then join; cross-database SQL is not portable.
- Very wide table (>200 columns): Split into a core model plus an extension model; document the split in the schema YAML.