Snowflake MigrationClickHouse Workshops

dbt no Snowflake

Como o pipeline Medallion de origem é criado: sources, views de preparação, modelos MERGE incrementais, snapshots e testes.

Este documento explica como o dbt (data build tool) é usado no Laboratório de Migração do Snowflake com NYC Taxi — o que ele faz, por que cada elemento existe e como raciocinar sobre ele.


O que o dbt faz (e não faz)

O dbt transforma dados que já estão no seu banco. Ele não carrega dados externos, move arquivos nem gerencia infraestrutura. Sua função é transformar tabelas brutas em tabelas limpas, testadas e prontas para analytics — executando o SQL que você escreve.

Pense nele como um sistema de compilação para SQL. Cada arquivo .sql em models/ é um modelo que se torna uma tabela ou view no Snowflake. O dbt cuida do código repetitivo de CREATE OR REPLACE, resolve dependências entre modelos e executa seus testes.


Estrutura do projeto

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

As três camadas (arquitetura Medallion)

Camada 1 — Preparação (models/staging/)

Finalidade: receber os dados brutos exatamente como chegaram e torná-los utilizáveis.

Esses modelos são criados no esquema STAGING como views (sem custo de armazenamento — elas são executadas no momento da consulta). Cada modelo de preparação realiza uma tarefa:

ModeloOrigemO que faz
stg_tripsRAW.TRIPS_RAWRenomeia colunas para snake_case, adiciona duration_minutes e desaninha a coluna VARIANT TRIP_METADATA em colunas tipadas
stg_taxi_zonesANALYTICS.DIM_TAXI_ZONESFaz uma limpeza leve, adiciona proteções com COALESCE e fornece um nó de linhagem do dbt

O trabalho mais importante aqui é desaninhar a coluna VARIANT TRIP_METADATA. A sintaxe de caminho com dois-pontos do Snowflake extrai os campos JSON aninhados:

-- Snowflake: colon-path notation
TRIP_METADATA:driver.rating::FLOAT     AS driver_rating,
TRIP_METADATA:app.surge_multiplier::FLOAT AS surge_multiplier

Esse é um dos desafios da migração — o ClickHouse usa JSONExtractFloat(TRIP_METADATA, 'driver', 'rating') em seu lugar.

Camada 2 — Intermediária (models/intermediate/)

Finalidade: realizar todas as junções em um único lugar para que não precisem ser repetidas.

int_trips_enriched associa stg_trips a todas as dimensões (zonas, tipos de pagamento, fornecedores e datas) e produz uma linha larga e totalmente desnormalizada por corrida. Ele é declarado como efêmero, o que significa que o dbt incorpora seu SQL em todos os modelos que o referenciam — nenhuma tabela ou view física é criada no Snowflake.

-- dbt_project.yml
intermediate:
  +materialized: ephemeral   # compiled inline, no CREATE TABLE

Use uma materialização efêmera quando o resultado intermediário for necessário para apenas um modelo posterior e você não quiser arcar com o armazenamento nem com a sobrecarga de compilação da consulta.

Camada 3 — Analytics (models/analytics/)

Finalidade: tabelas finais prontas para dashboards.

Elas são criadas no esquema ANALYTICS. Há dois tipos:

Tabelas de dimensão estáticas — pequenas e totalmente recarregadas em cada dbt run:

ModeloLinhasObservações
dim_date~7.670Estrutura de datas de 2009 a 2029, com trimestres fiscais e feriados federais dos EUA
dim_payment_type6Repasse dos dados de seed
dim_vendor3Repasse dos dados de seed
dim_taxi_zones265Repasse por meio de stg_taxi_zones

Tabelas de fatos/agregações incrementais — grandes, atualizadas com MERGE em cada execução:

ModeloLinhasObservações
fact_trips50 miUma linha por corrida, totalmente desnormalizada
agg_hourly_zone_trips~9 miContagens pré-agregadas por hora e por zona

Materializações

Uma materialização controla o que o dbt cria no Snowflake para determinado modelo.

MaterializaçãoObjeto criado no SnowflakeQuando usar
viewCREATE VIEWBarata; sempre reflete os dados mais recentes; usada para preparação
tableCREATE TABLE AS SELECTReconstrução completa em cada execução; usada para dimensões pequenas
incrementalMERGE INTO na tabela existenteTabelas grandes; processa apenas as novas linhas
ephemeral(nenhum objeto — incorporado como CTE)Lógica intermediária compartilhada por um modelo posterior

Os dois modelos incrementais demonstram estratégias diferentes:

fact_trips — processa as novas corridas desde a última execução:

{% if is_incremental() %}
  WHERE pickup_at > (SELECT MAX(pickup_at) FROM {{ this }})
{% endif %}

agg_hourly_zone_trips — agrega novamente uma janela contínua de 2 horas para capturar dados atrasados:

{% if is_incremental() %}
  WHERE pickup_at >= DATEADD('hour', -2, CURRENT_TIMESTAMP())
{% endif %}

Na primeira execução (tabela vazia), is_incremental() retorna false e o conjunto de dados inteiro é processado. Nas execuções seguintes, somente os novos dados são processados. Se o esquema mudar e você precisar reconstruir tudo, execute:

dbt run --full-refresh

A estratégia MERGE (principal desafio da migração)

Quando incremental_strategy = 'merge', o dbt gera uma instrução MERGE INTO do 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 ...;

Esse é um dos desafios de migração mais importantes documentados no laboratório. O ClickHouse não possui uma instrução MERGE. O equivalente no ClickHouse é usar o mecanismo de tabela ReplacingMergeTree e adicionar FINAL às consultas ou usar um CollapsingMergeTree para uma semântica explícita de inserção/exclusão.


Nomes de esquema: a macro generate_schema_name

O comportamento padrão do dbt concatena o esquema de destino de profiles.yml com o esquema personalizado de dbt_project.yml:

target schema = STAGING  +  custom schema = ANALYTICS  →  STAGING_ANALYTICS  (wrong)

Este projeto substitui esse comportamento com uma macro personalizada em 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: os modelos com +schema: ANALYTICS são criados em ANALYTICS, não em STAGING_ANALYTICS.

Essa macro é obrigatória sempre que há vários esquemas em um projeto dbt e você não quer adicionar o nome do esquema de destino como prefixo.


Conexão e credenciais (profiles.yml)

O dbt se conecta ao Snowflake usando um perfil definido em ~/.dbt/profiles.yml (que nunca deve ser incluído no Git). O nome do perfil em dbt_project.yml precisa corresponder:

# 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

Pontos principais:

  • schema: STAGING é o esquema padrão. Modelos sem uma substituição +schema: são criados nele.
  • role: DBT_ROLE é uma função de privilégio mínimo criada pelo Terraform somente com as permissões de que o dbt precisa.
  • threads: 4 controla quantos modelos o dbt cria em paralelo.
  • As credenciais vêm de variáveis de ambiente, carregadas de .env antes da execução da preparação.

Testes

Há dois tipos de testes do dbt:

Testes de esquema (declarados em 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 e unique são integrados. Os testes dbt_expectations vêm do pacote calogica/dbt_expectations declarado em packages.yml.

Testes SQL personalizados (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

Testes personalizados são apenas consultas SQL. O dbt as executa e falha caso alguma linha seja retornada.

Execute todos os testes com:

dbt test

Pacotes de terceiros (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"]

Instale-os antes do primeiro uso:

dbt deps

dbt_utils fornece o gerador date_spine usado em dim_date.sql. dbt_expectations oferece testes de intervalo e distribuição que vão além dos testes integrados not_null/unique.


Grafo de dependências

O dbt cria os modelos automaticamente na ordem correta seguindo as chamadas {{ 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') }} é como um modelo declara uma dependência de outro. {{ source('raw', 'TRIPS_RAW') }} declara uma dependência de uma tabela externa (definida em sources.yml).


Comandos comuns

ComandoO que faz
dbt depsInstala os pacotes de packages.yml
dbt runCria todos os modelos (incrementalmente quando possível)
dbt run --full-refreshReconstrói todos os modelos incrementais desde o início
dbt run -s fact_tripsCria somente fact_trips e suas dependências
dbt testExecuta todos os testes de esquema e personalizados
dbt buildExecuta dbt run + dbt test em conjunto
dbt compileGera SQL sem executá-lo (útil para depuração)
dbt docs generate && dbt docs serveCria o grafo de linhagem e permite visualizá-lo em um navegador

Neste projeto, dbt run --full-refresh é acionado automaticamente por setup.sh se fact_trips estiver vazia (na primeira execução ou após uma desmontagem).


Como o dbt se encaixa na preparação 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

O dbt ocupa o meio do pipeline. Ele não pode ser executado até que as tabelas brutas existam e contenham dados. O script setup.sh garante essa ordem.

Nesta página

PT