Snowflake MigrationClickHouse Workshops

04 Reconstruction du pipeline dbt

Reconstruisez le pipeline Medallion dans ClickHouse avec dbt-clickhouse — modèles incrémentiels delete_insert, ReplacingMergeTree et vues matérialisées actualisables — puis créez le dictionnaire des zones.

Point de départ

Module 03 terminé : le service ClickHouse Cloud est actif et accessible, et setup.sh a écrit .clickhouse_state sur le disque avec CLICKHOUSE_HOST et CLICKHOUSE_PORT. Toutes les tables cibles et vues de préparation existent — default.trips_raw, les deux vues (stg_trips, stg_taxi_zones), les six tables analytics (fact_trips, agg_hourly_zone_trips, dim_taxi_zones, dim_payment_type, dim_vendor, dim_date) et la vue matérialisée actualisable analytics.mv_live_trip_feed — puisque le dbt run du module 03 a déjà créé les sept objets. Toutes les tables analytics restent vides sauf mv_live_trip_feed, qui contient la ligne d’instantané de cette création. En dehors de cela, seule default.trips_raw contient des données : environ 50 millions de lignes. Le producteur Snowflake fonctionne encore ; ClickHouse accuse donc un retard à peu près égal à la fenêtre de migration. Prévoyez environ 30 minutes.

Pourquoi

Le module 03 a prouvé que ClickHouse pouvait contenir 50 millions de lignes, mais pas que le pipeline pouvait y fonctionner — vues de préparation, table de faits incrémentielle, rechargement des dimensions et tests qui détectent un modèle cassé avant qu’un partenaire ne le voie. C’est ce que ce module reconstruit : les mêmes modèles Medallion que dans le module 01, exprimés avec dbt-clickhouse plutôt que dbt-snowflake et exécutés sur les tables déjà créées au module 03.

La logique des modèles ne change pas : stg_trips convertit toujours les types et extrait le JSON, int_trips_enriched joint toujours les dimensions et fact_trips aboutit toujours à une ligne par course. C’est la couche de matérialisation sous-jacente qui change : pas de MERGE INTO, pas de Task Snowflake et pas de cluster_by. Ce module montre que la migration n’est pas un simple déversement ponctuel : le pipeline exécuté quotidiennement par l’équipe d’un partenaire, selon le même planning et sous le contrôle du même dbt test, continue de fonctionner après le changement de warehouse.

Concepts — fonctionnement interne

L’ensemble des modèles diffère du pipeline Snowflake de quatre manières. La référence de configuration complète se trouve dans dbt sur ClickHouse ; voici la version courte nécessaire avant d’exécuter dbt run à l’étape 1. Pour le pipeline source remplacé par ces modèles, consultez dbt sur Snowflake.

1. delete_insert remplace MERGE. ClickHouse n’a pas d’instruction MERGE INTO. Là où le pipeline Snowflake employait incremental_strategy: merge pour effectuer l’upsert de fact_trips et agg_hourly_zone_trips, les modèles ClickHouse utilisent incremental_strategy: delete_insert : dbt supprime les lignes correspondant à la unique_key du lot entrant, puis insère ce lot. Pour fact_trips, la unique_key est trip_id, et le filtre incrémentiel prend updated_at comme marque plutôt que pickup_at. Une correction tarifaire réinsère le même trip_id avec le même pickup_at, mais un updated_at plus récent ; une marque basée sur pickup_at l’ignorerait silencieusement.

2. ReplacingMergeTree est la protection sous delete_insert, pas son substitut. Les deux modèles incrémentiels sont déclarés ReplacingMergeTree(updated_at). Si une exécution de delete_insert se termine normalement, la table possède déjà une ligne par clé et le moteur n’a rien à nettoyer. Si l’exécution s’interrompt à mi-parcours — panne après la suppression, avant l’insertion — les fusions en arrière-plan finissent par dédupliquer les lignes restantes en conservant celle dont updated_at est le plus élevé. Ne comptez jamais uniquement sur ReplacingMergeTree pour réaliser le travail prévu pour delete_insert : les fusions en arrière-plan sont asynchrones et peuvent prendre de quelques minutes à plusieurs heures sur une table de cette taille.

3. Les vues matérialisées actualisables remplacent la tâche planifiée. Le pipeline Snowflake utilisait une Task planifiée exécutant une procédure stockée afin de maintenir un agrégat glissant. Le projet dbt ClickHouse déclare à la place mv_live_trip_feed avec materialized = 'materialized_view' et engine = 'ReplacingMergeTree(refreshed_at)' — créée par dbt run comme vue matérialisée actualisable, au lieu des plus de 30 lignes de DDL Snowflake CREATE TASK nécessaires au même effet. Ce module n’active aucun intervalle d’actualisation ; il faudrait exécuter manuellement ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ..., ce que l’atelier n’automatise pas (le module 05 en explique la raison).

4. mv_live_trip_feed n’a aucun équivalent dans Snowflake. Elle ne traduit pas un modèle existant : il s’agit d’une capacité ajoutée par la migration. Une vue matérialisée ClickHouse standard se déclenche une fois par INSERT et ne voit que les lignes du lot ; elle ne peut donc pas calculer correctement un agrégat portant sur toute la durée, comme le nombre total de courses ou le tarif moyen. Une vue matérialisée ACTUALISABLE réexécute au contraire sa requête complète — ici, SELECT ... FROM {{ ref('fact_trips') }} — selon un planning, de sorte que chaque actualisation voit toute la table. La partie Snowflake de cet atelier n’offrait pas cette possibilité.

Les modèles effectivement créés par dbt et le moment où chacun reçoit des données :

ModèleCoucheMatérialisationRempliNotes
stg_tripsstagingVueÉtape 1 (chaque exécution)Convertit les types et utilise JSONExtract* pour trip_metadata
stg_taxi_zonesstagingVueÉtape 1 (chaque exécution)Transmission de la dimension des zones
int_trips_enrichedstagingÉphémère— (incorporé comme CTE)Toutes les jointures de dimensions ; aucune table physique
fact_tripsanalyticsIncrémentielÉtape 1delete_insert avec la clé trip_id, marque sur updated_at ; ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id)
agg_hourly_zone_tripsanalyticsIncrémentielModule 05, après la basculeFenêtre glissante de 2 heures ; le filtre incrémentiel ne correspond qu’aux lignes du producteur actif — voir l’étape 1
dim_taxi_zonesanalyticsTableÉtape 1Rechargement complet à chaque exécution ; source du dictionnaire des zones à l’étape 2
dim_payment_typeanalyticsTableÉtape 1Rechargement complet à chaque exécution
dim_vendoranalyticsTableÉtape 1Rechargement complet à chaque exécution
dim_dateanalyticsTableÉtape 1Calendrier statique, 2009 à 2029 ; rechargement complet à chaque exécution
mv_live_trip_feedanalyticsVue matérialisée (actualisable)dbt run du module 03 ; intervalle d’actualisation jamais activéAucun équivalent Snowflake — voir le point 4

Étape 1 — Remplir la couche analytics

Activez l’environnement dbt-clickhouse créé au module 00, puis exécutez dbt run une deuxième fois. Le module 03 l’a déjà exécuté sur des tables vides pour créer les schémas ; cette fois, des données réelles se trouvent derrière — trips_raw contient maintenant 50 millions de lignes.

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .venv/bin/activate
source .env && source .clickhouse_state

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt run

dbt lit la connexion ClickHouse dans ~/.dbt/profiles.yml, à partir de dbt/nyc_taxi_dbt_ch/profiles.yml.example — le même profil que le module 03 a employé pour créer les schémas vides.

Résultat attendu : environ 8 à 12 minutes (50 millions de lignes traitées par les modèles incrémentiels).

agg_hourly_zone_trips sera vide après cette exécution — c’est attendu, pas un échec. Son filtre incrémentiel est WHERE pickup_at >= now() - INTERVAL 2 HOUR, qui ne correspond qu’aux lignes écrites par le producteur actif. Toutes les lignes migrées sont historiques ; aucune n’entre dans une fenêtre de 2 heures calculée à partir de maintenant. La table reste vide jusqu’à ce que la bascule du module 05 démarre le producteur ClickHouse. Ne perdez pas de temps à diagnostiquer un pipeline qui n’est pas cassé.

Exécutez ensuite la suite de tests :

dbt test

Résultat attendu : tous les tests réussissent.

Vérification :

SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million

SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0 (normal — populated after cutover in module 05)

SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265

Étape 2 — Créer le dictionnaire des zones

analytics.dim_taxi_zones contient maintenant les 265 zones NYC TLC. Créez analytics.taxi_zones_dict, un dictionnaire en mémoire alimenté par cette table, afin que les requêtes en aval retrouvent l’arrondissement d’une zone avec dictGet() plutôt qu’une JOIN.

Avantage d’un dictionnaire sur une jointure. Un dictionnaire est chargé une fois en mémoire et y reste disponible ; les recherches suivantes sont pratiquement gratuites. Une JOIN avec dim_taxi_zones relit et associe la dimension à chaque exécution. Pour une petite table de référence rarement modifiée — 265 lignes, entièrement rechargées à chaque dbt run — le choix est évident. Le benchmark du module 05 interroge directement taxi_zones_dict avec dictGet : cette étape est donc une dépendance obligatoire, pas un complément facultatif.

Chargez les informations de connexion, puis appliquez le DDL du dictionnaire avec clickhouse-client ou l’API HTTP, selon ce qui est disponible :

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .clickhouse_state

# Via clickhouse-client
clickhouse-client \
  --host "${CLICKHOUSE_HOST}" \
  --port 9440 \
  --user default \
  --password "${CLICKHOUSE_PASSWORD}" \
  --secure \
  --multiquery \
  < scripts/04_create_dictionary.sql

# Or via HTTP API
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/" \
  --user "default:${CLICKHOUSE_PASSWORD}" \
  --data-binary @scripts/04_create_dictionary.sql

Vérification :

-- Should return 'Manhattan' for zone 42
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));

-- Should show status = LOADED, element_count = 265
SELECT name, status, element_count
FROM system.dictionaries
WHERE name = 'taxi_zones_dict';

Comment vérifier que vous avez terminé

SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million
SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0

Le zéro est correct et ne signale pas un échec. Le filtre incrémentiel de agg_hourly_zone_trips (WHERE pickup_at >= now() - INTERVAL 2 HOUR) ne correspond qu’aux lignes écrites par le producteur actif ; toutes les lignes actuellement dans ClickHouse sont des données historiques déplacées par le script du module 03. Aucune n’a moins de 2 heures par rapport à now(). La table ne se remplit qu’après le démarrage du producteur ClickHouse lors de la bascule du module 05. D’ici là, tout graphique alimenté par cette table reste vide, ce qui est attendu et ne nécessite aucun dépannage.

SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt test

Résultat attendu : tous les tests réussissent.

-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));

État final

La couche analytics est remplie et testée : fact_trips contient environ 50 millions de lignes ; dim_taxi_zones, dim_payment_type, dim_vendor et dim_date sont entièrement chargées ; dbt test réussit de bout en bout ; et analytics.taxi_zones_dict est actif et renvoie les arrondissements avec dictGet(). agg_hourly_zone_trips reste vide — par conception, pas à cause d’un défaut — jusqu’à la bascule du module 05.

Les tableaux de bord, le benchmark ClickHouse contre Snowflake et la bascule relèvent du module 05, pas de celui-ci.

Le producteur Snowflake fonctionne toujours et l’écart entre Snowflake et ClickHouse reste ouvert. Ce module n’a touché ni au producteur ni au script de migration. Le module 05 refermera volontairement l’écart grâce au même processus contrôlé en deux passages annoncé au module 03. N’arrêtez pas le producteur maintenant.

Sur cette page

Suivre votre progression ?

Facultatif. Nous envoyons un lien par e-mail pour confirmer votre adresse ; la progression est enregistrée après son ouverture.

Utilisez votre adresse e-mail professionnelle, et non une adresse personnelle.

Le suivi de la progression exige aussi d’accepter les Conditions d’utilisation actuelles dans les Paramètres de confidentialité.

FR