Snowflake MigrationClickHouse Workshops

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

Las 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:

ModeloOrigenQué hace
stg_tripsRAW.TRIPS_RAWCambia los nombres de columnas a snake_case, añade duration_minutes y aplana la columna VARIANT TRIP_METADATA en columnas tipadas
stg_taxi_zonesANALYTICS.DIM_TAXI_ZONESLimpieza 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_multiplier

Es 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 TABLE

Usa 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:

ModeloFilasNotas
dim_date~7 670Secuencia de fechas 2009–2029, con trimestres fiscales y festivos federales de EE. UU.
dim_payment_type6Paso directo desde datos semilla
dim_vendor3Paso directo desde datos semilla
dim_taxi_zones265Paso directo mediante stg_taxi_zones

Tablas incrementales de hechos/agregados: grandes y actualizadas con MERGE en cada ejecución:

ModeloFilasNotas
fact_trips50 MUna fila por viaje, totalmente desnormalizada
agg_hourly_zone_trips~9 MRecuentos horarios preagregados por zona

Materializaciones

Una materialización controla qué crea dbt en Snowflake para cada modelo.

MaterializaciónObjeto de SnowflakeCuándo usarla
viewCREATE VIEWBarata; siempre refleja los datos más recientes; usada en staging
tableCREATE TABLE AS SELECTReconstrucción completa en cada ejecución; usada para dimensiones pequeñas
incrementalMERGE INTO en una tabla existenteTablas grandes; solo procesa filas nuevas
ephemeralSin objeto; insertado como CTELó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-refresh

La 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: 4

Puntos clave:

  • schema: STAGING es el esquema predeterminado. Los modelos sin una excepción +schema: se crean aquí.
  • role: DBT_ROLE es un rol de privilegios mínimos creado por Terraform únicamente con los permisos que necesita dbt.
  • threads: 4 controla cuántos modelos construye dbt en paralelo.
  • Las credenciales proceden de variables de entorno cargadas desde .env antes 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: 1000

not_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 < 0

Las pruebas personalizadas son simplemente consultas SQL. dbt las ejecuta y falla si devuelven alguna fila.

Ejecuta todas las pruebas con:

dbt test

Paquetes 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 deps

dbt_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

ComandoQué hace
dbt depsInstala los paquetes de packages.yml
dbt runConstruye todos los modelos (incrementalmente cuando es posible)
dbt run --full-refreshReconstruye desde cero todos los modelos incrementales
dbt run -s fact_tripsConstruye solo fact_trips y sus dependencias
dbt testEjecuta todas las pruebas de esquema y personalizadas
dbt builddbt run + dbt test juntos
dbt compileGenera SQL sin ejecutarlo (útil para depurar)
dbt docs generate && dbt docs serveGenera 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 task

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

En esta página

ES