Snowflake MigrationClickHouse Workshops

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

Les 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èleSourceRôle
stg_tripsRAW.TRIPS_RAWRenomme les colonnes en snake_case, ajoute duration_minutes et transforme la colonne VARIANT TRIP_METADATA en colonnes typées
stg_taxi_zonesANALYTICS.DIM_TAXI_ZONESEffectue 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_multiplier

C’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 TABLE

Utilisez 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èleLignesRemarques
dim_dateenviron 7 670Plage de dates de 2009 à 2029, avec trimestres fiscaux et jours fériés fédéraux des États-Unis
dim_payment_type6Transmission directe des données d’amorçage
dim_vendor3Transmission directe des données d’amorçage
dim_taxi_zones265Transmission directe via stg_taxi_zones

Tables de faits et d’agrégats incrémentielles — grandes tables mises à jour avec MERGE à chaque exécution :

ModèleLignesRemarques
fact_trips50 millionsUne ligne par trajet, entièrement dénormalisée
agg_hourly_zone_tripsenviron 9 millionsNombre 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érialisationObjet SnowflakeCas d’utilisation
viewCREATE VIEWPeu coûteuse, reflète toujours les dernières données ; utilisée pour le staging
tableCREATE TABLE AS SELECTReconstruction complète à chaque exécution ; utilisée pour les petites dimensions
incrementalMERGE INTO dans une table existanteGrandes tables ; ne traite que les nouvelles lignes
ephemeralAucun objet, intégré sous forme de CTELogique 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-refresh

Straté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: 4

Points essentiels :

  • schema: STAGING est le schéma par défaut. Les modèles qui ne remplacent pas +schema: aboutissent ici.
  • role: DBT_ROLE est un rôle aux privilèges minimaux, créé par Terraform et limité aux autorisations nécessaires à dbt.
  • threads: 4 contrôle le nombre de modèles construits en parallèle par dbt.
  • Les identifiants proviennent de variables d’environnement chargées depuis .env avant 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: 1000

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

Les 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 test

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

dbt_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

CommandeRôle
dbt depsInstaller les paquets de packages.yml
dbt runConstruire tous les modèles, en mode incrémentiel lorsque possible
dbt run --full-refreshReconstruire intégralement tous les modèles incrémentiels
dbt run -s fact_tripsConstruire uniquement fact_trips et ses dépendances
dbt testExécuter tous les tests de schéma et personnalisés
dbt buildExécuter ensemble dbt run et dbt test
dbt compileGénérer le SQL sans l’exécuter, utile pour le débogage
dbt docs generate && dbt docs serveConstruire 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 task

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

Sur cette page

FR