dbt en Snowflake
Cómo se construye la canalización Medallion de origen: fuentes, vistas de staging, modelos MERGE incrementales, instantáneas y pruebas.
Este documento explica cómo se usa dbt (data build tool) en el laboratorio de migración NYC Taxi desde Snowflake: qué hace, por qué existe cada pieza y cómo razonar sobre ella.
Qué hace dbt (y qué no hace)
dbt transforma datos que ya se encuentran en tu base de datos. No carga datos del exterior, mueve archivos ni administra infraestructura. Su trabajo es convertir tablas sin procesar en tablas limpias, probadas y listas para el análisis ejecutando el SQL que tú escribes.
Piensa en él como un sistema de compilación para SQL. Cada archivo .sql de models/
es un modelo que se convierte en tabla o vista de Snowflake. dbt se ocupa del código
repetitivo de CREATE OR REPLACE, resuelve dependencias y ejecuta tus pruebas.
Estructura del proyecto
dbt/nyc_taxi_dbt/
├── dbt_project.yml # Project config: name, folder layout, materialization defaults
├── profiles.yml.example # Connection config template (copy to ~/.dbt/profiles.yml)
├── packages.yml # Third-party dbt packages
│
├── macros/
│ ├── generate_schema_name.sql # Overrides dbt's default schema naming logic
│ └── generate_surrogate_key.sql # Wrapper for consistent surrogate key generation
│
└── models/
├── sources.yml # Declares RAW.TRIPS_RAW as an external source
│
├── staging/ # Layer 1: clean and rename raw columns
│ ├── schema.yml # Column-level tests for staging models
│ ├── stg_trips.sql
│ └── stg_taxi_zones.sql
│
├── intermediate/ # Layer 2: joins and enrichment (no physical table)
│ └── int_trips_enriched.sql
│
└── analytics/ # Layer 3: final tables consumed by dashboards
├── schema.yml
├── fact_trips.sql
├── agg_hourly_zone_trips.sql
├── dim_date.sql
├── dim_payment_type.sql
├── dim_taxi_zones.sql
└── dim_vendor.sqlLas tres capas (arquitectura Medallion)
Capa 1: staging (models/staging/)
Objetivo: tomar los datos sin procesar tal como llegaron y hacerlos utilizables.
Estos modelos se crean en el esquema STAGING como vistas (sin coste de
almacenamiento: se ejecutan al consultar). Cada modelo de staging cumple una función:
| Modelo | Origen | Qué hace |
|---|---|---|
stg_trips | RAW.TRIPS_RAW | Cambia los nombres de columnas a snake_case, añade duration_minutes y aplana la columna VARIANT TRIP_METADATA en columnas tipadas |
stg_taxi_zones | ANALYTICS.DIM_TAXI_ZONES | Limpieza ligera, añade protecciones COALESCE y proporciona un nodo de linaje dbt |
El trabajo más importante es aplanar la columna VARIANT TRIP_METADATA. La sintaxis de
ruta con dos puntos de Snowflake extrae campos JSON anidados:
-- Snowflake: colon-path notation
TRIP_METADATA:driver.rating::FLOAT AS driver_rating,
TRIP_METADATA:app.surge_multiplier::FLOAT AS surge_multiplierEs uno de los retos de la migración: ClickHouse usa en su lugar
JSONExtractFloat(TRIP_METADATA, 'driver', 'rating').
Capa 2: intermedia (models/intermediate/)
Objetivo: realizar todas las uniones en un único lugar para no repetirlas.
int_trips_enriched une stg_trips con todas las dimensiones (zonas, tipos de pago,
proveedores y fechas) y genera una fila amplia y totalmente desnormalizada por viaje.
Se declara ephemeral, por lo que dbt inserta su SQL en el modelo que lo referencia:
no se crea ninguna tabla o vista física en Snowflake.
-- dbt_project.yml
intermediate:
+materialized: ephemeral # compiled inline, no CREATE TABLEUsa ephemeral cuando el resultado intermedio solo lo necesita un modelo posterior y no quieres pagar por almacenamiento ni por la sobrecarga de compilación de consultas.
Capa 3: analítica (models/analytics/)
Objetivo: tablas finales listas para los paneles.
Se crean en el esquema ANALYTICS. Hay dos tipos:
Tablas de dimensiones estáticas: pequeñas y recargadas por completo en cada dbt run:
| Modelo | Filas | Notas |
|---|---|---|
dim_date | ~7 670 | Secuencia de fechas 2009–2029, con trimestres fiscales y festivos federales de EE. UU. |
dim_payment_type | 6 | Paso directo desde datos semilla |
dim_vendor | 3 | Paso directo desde datos semilla |
dim_taxi_zones | 265 | Paso directo mediante stg_taxi_zones |
Tablas incrementales de hechos/agregados: grandes y actualizadas con MERGE en cada ejecución:
| Modelo | Filas | Notas |
|---|---|---|
fact_trips | 50 M | Una fila por viaje, totalmente desnormalizada |
agg_hourly_zone_trips | ~9 M | Recuentos horarios preagregados por zona |
Materializaciones
Una materialización controla qué crea dbt en Snowflake para cada modelo.
| Materialización | Objeto de Snowflake | Cuándo usarla |
|---|---|---|
view | CREATE VIEW | Barata; siempre refleja los datos más recientes; usada en staging |
table | CREATE TABLE AS SELECT | Reconstrucción completa en cada ejecución; usada para dimensiones pequeñas |
incremental | MERGE INTO en una tabla existente | Tablas grandes; solo procesa filas nuevas |
ephemeral | Sin objeto; insertado como CTE | Lógica intermedia compartida por un modelo posterior |
Los dos modelos incrementales muestran estrategias distintas:
fact_trips: procesa los viajes nuevos desde la última ejecución:
{% if is_incremental() %}
WHERE pickup_at > (SELECT MAX(pickup_at) FROM {{ this }})
{% endif %}agg_hourly_zone_trips: vuelve a agregar una ventana móvil de 2 horas para capturar
datos que llegan tarde:
{% if is_incremental() %}
WHERE pickup_at >= DATEADD('hour', -2, CURRENT_TIMESTAMP())
{% endif %}En la primera ejecución (tabla vacía), is_incremental() devuelve false y se procesa
todo el conjunto. En las siguientes solo se procesan datos nuevos. Si cambia el esquema
y necesitas reconstruir desde cero, ejecuta:
dbt run --full-refreshLa estrategia MERGE (reto principal de la migración)
Cuando incremental_strategy = 'merge', dbt genera una sentencia MERGE INTO de
Snowflake:
MERGE INTO ANALYTICS.FACT_TRIPS AS target
USING (SELECT ...) AS source
ON target.trip_id = source.trip_id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...;Es uno de los retos de migración más importantes documentados en el laboratorio.
ClickHouse no tiene una sentencia MERGE. El equivalente consiste en usar el motor
de tabla ReplacingMergeTree y añadir FINAL a las consultas, o usar
CollapsingMergeTree para una semántica explícita de inserción/eliminación.
Nombres de esquemas: la macro generate_schema_name
El comportamiento predeterminado de dbt concatena el esquema de destino de
profiles.yml con el esquema personalizado de dbt_project.yml:
target schema = STAGING + custom schema = ANALYTICS → STAGING_ANALYTICS (wrong)Este proyecto lo sobrescribe con una macro personalizada en
macros/generate_schema_name.sql:
{% macro generate_schema_name(custom_schema_name, node) -%}
{%- if custom_schema_name is none -%}
{{ target.schema | upper }} -- no custom schema → use target schema
{%- else -%}
{{ custom_schema_name | upper }} -- custom schema → use it directly
{%- endif -%}
{%- endmacro %}Resultado: los modelos con +schema: ANALYTICS se crean en ANALYTICS, no en
STAGING_ANALYTICS.
Esta macro es necesaria siempre que haya varios esquemas en un proyecto dbt y no quieras anteponerles el nombre del esquema de destino.
Conexión y credenciales (profiles.yml)
dbt se conecta a Snowflake mediante un perfil definido en
~/.dbt/profiles.yml (nunca se incluye en git). El nombre del perfil en
dbt_project.yml debe coincidir:
# dbt_project.yml
profile: 'nyc_taxi'
# ~/.dbt/profiles.yml
nyc_taxi:
target: dev
outputs:
dev:
type: snowflake
account: "{{ env_var('SNOWFLAKE_ORG') }}-{{ env_var('SNOWFLAKE_ACCOUNT') }}"
role: DBT_ROLE
database: NYC_TAXI_DB
warehouse: TRANSFORM_WH
schema: STAGING # ← this is the "target schema" / default schema
threads: 4Puntos clave:
schema: STAGINGes el esquema predeterminado. Los modelos sin una excepción+schema:se crean aquí.role: DBT_ROLEes un rol de privilegios mínimos creado por Terraform únicamente con los permisos que necesita dbt.threads: 4controla cuántos modelos construye dbt en paralelo.- Las credenciales proceden de variables de entorno cargadas desde
.envantes de ejecutar la configuración.
Pruebas
Las pruebas dbt adoptan dos formas:
Pruebas de esquema (declaradas en schema.yml)
- name: trip_id
tests:
- not_null
- unique
- name: total_amount_usd
tests:
- dbt_expectations.expect_column_values_to_be_between:
min_value: 0
max_value: 1000not_null y unique están integradas. Las pruebas dbt_expectations proceden del
paquete calogica/dbt_expectations declarado en packages.yml.
Pruebas SQL personalizadas (tests/)
-- tests/assert_revenue_positive.sql
-- A passing test returns 0 rows
SELECT trip_id, total_amount_usd
FROM {{ ref('fact_trips') }}
WHERE total_amount_usd < 0Las pruebas personalizadas son simplemente consultas SQL. dbt las ejecuta y falla si devuelven alguna fila.
Ejecuta todas las pruebas con:
dbt testPaquetes de terceros (packages.yml)
packages:
- package: dbt-labs/dbt_utils
version: [">=1.0.0", "<2.0.0"]
- package: calogica/dbt_expectations
version: [">=0.10.0", "<1.0.0"]Instálalos antes del primer uso:
dbt depsdbt_utils proporciona el generador date_spine que usa dim_date.sql.
dbt_expectations aporta pruebas de intervalos y distribuciones que van más allá de
las integradas not_null/unique.
Grafo de dependencias
dbt construye los modelos en el orden correcto siguiendo automáticamente las llamadas
{{ ref() }}:
RAW.TRIPS_RAW (source — not managed by dbt)
└── stg_trips (view)
└── int_trips_enriched (ephemeral)
├── fact_trips (incremental table)
└── agg_hourly_zone_trips (incremental table)
ANALYTICS.DIM_TAXI_ZONES (seeded by SQL script)
└── stg_taxi_zones (view)
├── int_trips_enriched
└── dim_taxi_zones (table)
dbt_utils.date_spine
└── dim_date (table){{ ref('stg_trips') }} permite que un modelo declare su dependencia de otro.
{{ source('raw', 'TRIPS_RAW') }} declara una dependencia de una tabla externa
(definida en sources.yml).
Comandos frecuentes
| Comando | Qué hace |
|---|---|
dbt deps | Instala los paquetes de packages.yml |
dbt run | Construye todos los modelos (incrementalmente cuando es posible) |
dbt run --full-refresh | Reconstruye desde cero todos los modelos incrementales |
dbt run -s fact_trips | Construye solo fact_trips y sus dependencias |
dbt test | Ejecuta todas las pruebas de esquema y personalizadas |
dbt build | dbt run + dbt test juntos |
dbt compile | Genera SQL sin ejecutarlo (útil para depurar) |
dbt docs generate && dbt docs serve | Genera el grafo de linaje y permite explorarlo en un navegador |
En este proyecto, setup.sh activa automáticamente dbt run --full-refresh si
fact_trips está vacía (primera ejecución o tras desmontar el entorno).
Cómo encaja dbt en la configuración completa
terraform apply → creates warehouses, database, schemas, roles
scripts/01_create_tables.sql → creates raw tables, seeds dimension data
scripts/02_seed_data.sql → loads 50M synthetic trip rows
dbt deps && dbt build → transforms raw data into analytics-ready tables
scripts/03_create_streams_tasks.sql → creates CDC stream and scheduled taskdbt ocupa la parte central de la canalización. No puede ejecutarse hasta que existan
las tablas sin procesar y contengan datos. El script setup.sh se encarga de este orden.
Ejemplo resuelto: un plan completo
Un plan de migración cumplimentado para la carga de trabajo de taxis de NYC, con el que comparar el tuyo una vez escrito.
dbt en ClickHouse
Cómo configurar dbt-clickhouse: la estrategia incremental delete_insert, modelos ReplacingMergeTree y vistas materializadas actualizables.