dbt sur Snowflake
Construction du pipeline Medallion source : sources, vues de staging, modèles MERGE incrémentiels, snapshots et tests.
Ce document explique comment dbt (data build tool) est utilisé dans le laboratoire de migration Snowflake NYC Taxi : son rôle, la raison d’être de chaque élément et la façon de les aborder.
Ce que fait dbt, et ce qu’il ne fait pas
dbt transforme des données déjà présentes dans votre base. Il ne charge pas de données externes, ne déplace aucun fichier et ne gère pas l’infrastructure. Son rôle est de convertir des tables brutes en tables propres, testées et prêtes pour l’analyse, en exécutant le SQL que vous écrivez.
Voyez-le comme un système de construction pour SQL. Chaque fichier .sql du répertoire
models/ est un modèle qui devient une table ou une vue dans Snowflake. dbt prend en
charge le code répétitif de CREATE OR REPLACE, résout les dépendances entre les
modèles et exécute vos tests.
Organisation du projet
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.sqlLes trois couches de l’architecture Medallion
Couche 1 — Staging (models/staging/)
Objectif : rendre exploitables les données brutes exactement telles qu’elles sont arrivées.
Ces modèles sont créés dans le schéma STAGING sous forme de vues, sans coût de
stockage puisqu’ils s’exécutent au moment de la requête. Chaque modèle de staging
remplit une fonction :
| Modèle | Source | Rôle |
|---|---|---|
stg_trips | RAW.TRIPS_RAW | Renomme les colonnes en snake_case, ajoute duration_minutes et transforme la colonne VARIANT TRIP_METADATA en colonnes typées |
stg_taxi_zones | ANALYTICS.DIM_TAXI_ZONES | Effectue un léger nettoyage, ajoute des protections COALESCE et fournit un nœud de lignage dbt |
L’opération la plus importante consiste à aplatir la colonne VARIANT TRIP_METADATA.
La syntaxe de chemin avec deux-points de Snowflake extrait les champs JSON imbriqués :
-- Snowflake: colon-path notation
TRIP_METADATA:driver.rating::FLOAT AS driver_rating,
TRIP_METADATA:app.surge_multiplier::FLOAT AS surge_multiplierC’est l’un des défis de la migration : ClickHouse utilise à la place
JSONExtractFloat(TRIP_METADATA, 'driver', 'rating').
Couche 2 — Intermédiaire (models/intermediate/)
Objectif : effectuer toutes les jointures au même endroit afin de ne pas les répéter.
int_trips_enriched joint stg_trips à chaque dimension, notamment les zones, les
types de paiement, les fournisseurs et les dates, et produit une ligne large et
entièrement dénormalisée par trajet. Il est déclaré éphémère : dbt insère son SQL
directement dans tout modèle qui le référence. Aucune table ni vue physique n’est créée
dans Snowflake.
-- dbt_project.yml
intermediate:
+materialized: ephemeral # compiled inline, no CREATE TABLEUtilisez ce mode éphémère lorsque le résultat intermédiaire n’est requis que par un seul modèle en aval et que vous souhaitez éviter les coûts de stockage ou de compilation de requête.
Couche 3 — Analytics (models/analytics/)
Objectif : créer les tables finales prêtes à alimenter les tableaux de bord.
Elles aboutissent dans le schéma ANALYTICS. Il en existe deux catégories :
Tables de dimensions statiques — petites tables entièrement rechargées à chaque
dbt run :
| Modèle | Lignes | Remarques |
|---|---|---|
dim_date | environ 7 670 | Plage de dates de 2009 à 2029, avec trimestres fiscaux et jours fériés fédéraux des États-Unis |
dim_payment_type | 6 | Transmission directe des données d’amorçage |
dim_vendor | 3 | Transmission directe des données d’amorçage |
dim_taxi_zones | 265 | Transmission directe via stg_taxi_zones |
Tables de faits et d’agrégats incrémentielles — grandes tables mises à jour avec MERGE à chaque exécution :
| Modèle | Lignes | Remarques |
|---|---|---|
fact_trips | 50 millions | Une ligne par trajet, entièrement dénormalisée |
agg_hourly_zone_trips | environ 9 millions | Nombre horaire préagrégé par zone |
Matérialisations
Une matérialisation détermine l’objet que dbt crée dans Snowflake pour un modèle.
| Matérialisation | Objet Snowflake | Cas d’utilisation |
|---|---|---|
view | CREATE VIEW | Peu coûteuse, reflète toujours les dernières données ; utilisée pour le staging |
table | CREATE TABLE AS SELECT | Reconstruction complète à chaque exécution ; utilisée pour les petites dimensions |
incremental | MERGE INTO dans une table existante | Grandes tables ; ne traite que les nouvelles lignes |
ephemeral | Aucun objet, intégré sous forme de CTE | Logique intermédiaire partagée par un seul modèle en aval |
Les deux modèles incrémentiels illustrent des stratégies différentes :
fact_trips — traite les nouveaux trajets depuis la dernière exécution :
{% if is_incremental() %}
WHERE pickup_at > (SELECT MAX(pickup_at) FROM {{ this }})
{% endif %}agg_hourly_zone_trips — recalcule une fenêtre glissante de deux heures afin de
prendre en compte les données arrivées en retard :
{% if is_incremental() %}
WHERE pickup_at >= DATEADD('hour', -2, CURRENT_TIMESTAMP())
{% endif %}Lors de la toute première exécution, lorsque la table est vide, is_incremental()
renvoie false et l’ensemble des données est traité. Aux exécutions suivantes, seules
les nouvelles données le sont. Si le schéma change et que vous devez tout reconstruire,
exécutez :
dbt run --full-refreshStratégie MERGE, un défi essentiel de la migration
Lorsque incremental_strategy = 'merge', dbt génère une instruction Snowflake
MERGE INTO :
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 ...;Il s’agit de l’un des principaux défis de migration présentés dans le laboratoire.
ClickHouse ne possède pas d’instruction MERGE. Son équivalent consiste à employer
le moteur de table ReplacingMergeTree et à ajouter FINAL aux requêtes, ou à utiliser
CollapsingMergeTree pour une sémantique explicite d’insertion et de suppression.
Nommage des schémas : macro generate_schema_name
Par défaut, dbt concatène le schéma cible de profiles.yml et le schéma
personnalisé de dbt_project.yml :
target schema = STAGING + custom schema = ANALYTICS → STAGING_ANALYTICS (wrong)Ce projet remplace ce comportement avec une macro personnalisée dans
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 %}Résultat : les modèles configurés avec +schema: ANALYTICS aboutissent dans
ANALYTICS, et non dans STAGING_ANALYTICS.
Cette macro est nécessaire dès qu’un même projet dbt possède plusieurs schémas et que vous ne souhaitez pas ajouter le nom du schéma cible en préfixe.
Connexion et identifiants (profiles.yml)
dbt se connecte à Snowflake grâce à un profil défini dans
~/.dbt/profiles.yml, qui n’est jamais ajouté à git. Le nom du profil dans
dbt_project.yml doit correspondre :
# 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: 4Points essentiels :
schema: STAGINGest le schéma par défaut. Les modèles qui ne remplacent pas+schema:aboutissent ici.role: DBT_ROLEest un rôle aux privilèges minimaux, créé par Terraform et limité aux autorisations nécessaires à dbt.threads: 4contrôle le nombre de modèles construits en parallèle par dbt.- Les identifiants proviennent de variables d’environnement chargées depuis
.envavant l’exécution de la préparation.
Tests
Les tests dbt prennent deux formes :
Tests de schéma, déclarés dans 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 et unique sont intégrés. Les tests dbt_expectations proviennent du
paquet calogica/dbt_expectations déclaré dans packages.yml.
Tests SQL personnalisés (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 < 0Les tests personnalisés sont de simples requêtes SQL. dbt les exécute et échoue si la requête renvoie la moindre ligne.
Exécutez tous les tests avec :
dbt testPaquets tiers (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"]Installez-les avant la première utilisation :
dbt depsdbt_utils fournit le générateur date_spine utilisé dans dim_date.sql.
dbt_expectations apporte des tests de plage et de distribution qui vont au-delà des
tests intégrés not_null et unique.
Graphe de dépendances
dbt construit automatiquement les modèles dans le bon ordre en suivant les appels
{{ 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') }} permet à un modèle de déclarer une dépendance envers un autre.
{{ source('raw', 'TRIPS_RAW') }} déclare une dépendance envers une table externe,
définie dans sources.yml.
Commandes courantes
| Commande | Rôle |
|---|---|
dbt deps | Installer les paquets de packages.yml |
dbt run | Construire tous les modèles, en mode incrémentiel lorsque possible |
dbt run --full-refresh | Reconstruire intégralement tous les modèles incrémentiels |
dbt run -s fact_trips | Construire uniquement fact_trips et ses dépendances |
dbt test | Exécuter tous les tests de schéma et personnalisés |
dbt build | Exécuter ensemble dbt run et dbt test |
dbt compile | Générer le SQL sans l’exécuter, utile pour le débogage |
dbt docs generate && dbt docs serve | Construire et parcourir le graphe de lignage dans un navigateur |
Dans ce projet, dbt run --full-refresh est déclenché automatiquement par setup.sh
si fact_trips est vide, lors de la première exécution ou après la suppression de
l’environnement.
Place de dbt dans la préparation complète
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 se situe au milieu du pipeline. Il ne peut s’exécuter avant la création et le
chargement des tables brutes. Le script setup.sh respecte cet ordre.
Exemple corrigé : un plan terminé
Un plan de migration entièrement rempli pour la charge de travail NYC Taxi, à comparer avec le vôtre une fois sa rédaction terminée.
dbt sur ClickHouse
Configurer dbt-clickhouse : stratégie incrémentielle delete_insert, modèles ReplacingMergeTree et vues matérialisées actualisables.