03 Provisionnement et migration
Provisionnez ClickHouse Cloud avec Terraform, créez les tables cibles à partir de votre plan et transférez 50 millions de lignes grâce à un script de migration Python reprenable.
Point de départ
Module 02 terminé : migration-plan.md est rempli et toutes les cases de sa liste de
contrôle sont cochées, tandis que le producteur Snowflake fonctionne toujours. Le
setup.sh de ce module recherche ce fichier et avertit s’il est absent ou incomplet,
mais ne bloque jamais. Rien ne vous empêche de poursuivre sans le plan ; seule votre
compréhension des deux modules suivants en souffrira. Prévoyez environ 60 minutes au
total, dont 40 à 50 minutes de transfert autonome que vous pouvez laisser tourner en
arrière-plan. C’est également ici que commencent les dépenses de l’essai ClickHouse
Cloud : le provisionnement du service et ce module consomment environ 1 à 2 $ de crédits
d’essai (l’ensemble de l’atelier coûte à peu près 2 à 4 $).
Pourquoi
Ce module transforme le plan en réalité. Toutes les décisions consignées dans
migration-plan.md au module 02 — moteur MergeTree de chaque table, clé ORDER BY
déduite des requêtes réelles et traduction des constructions propres à Snowflake — sont
directement inscrites dans le DDL des tables, pas redéduites à partir de zéro. ClickHouse
n’offre aucun index que vous pourriez ajouter après coup : si une clé ORDER BY se
révèle incorrecte une fois 50 millions de lignes chargées, la solution est un
rechargement complet, pas un rapide ALTER.
C’est aussi pourquoi le contrôle souple est important, même s’il ne peut pas vous
arrêter. Sans plan terminé, le module réussira tout de même mécaniquement : dbt run
créera fact_trips comme ReplacingMergeTree et le script déplacera 50 millions de
lignes. Mais vous ignorerez pourquoi ce moteur a été choisi plutôt qu’un simple
MergeTree, pourquoi la clé de tri prend cette forme et comment justifier les gains
d’environ 6 à 9 fois du benchmark que le module 04 montrera plus tard. Le tableau
d’alignement ci-dessous rattache chaque choix mis en œuvre ici à la question de fiche à
laquelle il répond, afin que vous puissiez vérifier votre plan avant tout provisionnement.
Concepts — fonctionnement interne
Architecture cible. Le producteur continue d’envoyer de nouvelles courses à Snowflake pendant qu’un script Python ponctuel transfère les 50 millions de lignes existantes vers ClickHouse. Les deux systèmes fonctionnent en parallèle pendant toute la migration ; la bascule n’a pas encore eu lieu.
Côté ClickHouse, trips_raw est la table d’arrivée dans laquelle écrit le script. dbt
construit ensuite au-dessus les vues de préparation et le reste de la couche analytics.
Ce module crée le schéma, mais ne remplit encore rien au-delà de trips_raw :
Légende des couleurs du schéma :
- Vert — sources de données (producteur de courses avant et après la bascule)
- Bleu — tables Snowflake
- Orange — modèles et pipeline dbt
- Rouge — tables et vues matérialisées ClickHouse
- Cyan — tableaux de bord Apache Superset
- Flèches en pointillés — flux après la bascule
Pourquoi un script Python plutôt qu’un connecteur natif. Il existe plusieurs moyens de déplacer les données de Snowflake vers ClickHouse. Cet atelier utilise un script Python par lots. Voici pourquoi, en comparaison des solutions possibles :
| Méthode | Fonctionnement | Pourquoi elle n’est pas utilisée ici |
|---|---|---|
| ClickPipes (source Snowflake) | Connecteur natif ClickHouse Cloud — pas d’ETL, interface gérée | Snowflake n’est pas une source ClickPipes prise en charge. ClickPipes accepte Kafka, S3, Kinesis, le CDC PostgreSQL et MySQL, ainsi que le stockage objet. |
| Export vers S3 → ClickPipes S3 | COPY INTO @stage exporte du Parquet/CSV vers S3 ; le connecteur S3 de ClickPipes charge les données dans ClickHouse | Nécessite un bucket S3, un rôle IAM, un stage Snowflake et un compte AWS. Ajoute environ 3 étapes de configuration avant tout déplacement. Viable en production, mais trop d’infrastructure pour un atelier. |
Export vers S3 → clickhouse-client | Même export S3, chargé avec INSERT INTO ... SELECT FROM s3(...) | Mêmes prérequis S3. Le partenaire doit aussi gérer manuellement le découpage des fichiers et la reprise. |
| Snowflake → Kafka → ClickHouse | Un stream CDC Snowflake alimente un topic Kafka ; le connecteur Kafka ClickPipes réalise l’ingestion | Pipeline de streaming complet, adapté aux besoins de latence inférieure à une minute en production. Un cluster Kafka est trop lourd pour l’atelier. |
| Script Python (cet atelier) | snowflake-connector-python lit des lots de 100 000 lignes depuis le curseur ; clickhouse-connect les insère directement | Aucune infrastructure supplémentaire en dehors des paquets déjà nécessaires. Reprenable avec --resume (marque max(pickup_at)). Progression en temps réel. Environ 40 à 50 min pour 50 millions de lignes à près de 20 000 lignes/s, acceptable pour un exercice ponctuel. |
Pourquoi le script Python convient à cet atelier :
- Aucun compte AWS requis. Les approches S3 imposent la création d’un bucket, de politiques IAM et d’un stage externe Snowflake : trois étapes sans rapport avec ClickHouse.
- Solution autonome. Les deux paquets (
snowflake-connector-python,clickhouse-connect) sont installés dans le même environnement virtuel que dbt. Aucun service ni identifiant supplémentaire. - Reprise possible.
--resumepermet d’interrompre et de redémarrer le script sans risque.ReplacingMergeTree(_synced_at)assure la déduplication automatique des insertions répétées lors d’une nouvelle tentative. - Transparence. Les partenaires peuvent lire le script, comprendre la correspondance des colonnes et l’adapter à leur schéma, ce qui est plus formateur qu’un assistant graphique.
Gestion de l’écart de migration. Le producteur Snowflake reste actif pendant les 40 à 50 minutes d’exécution. Toutes les courses écrites dans Snowflake pendant ce temps manquent dans ClickHouse. L’atelier referme cet écart en deux passages lors de la bascule, présentée au module 05 :
- Arrêtez le producteur Snowflake afin de figer le jeu de données.
- Exécutez
python scripts/02_migrate_trips.py --resume: seules les lignes de l’écart sont transférées (quelques secondes, pas des minutes). - Démarrez le producteur ClickHouse.
La même déduplication ReplacingMergeTree(_synced_at) qui protège les nouvelles
tentatives de migration traite aussi ce cas : si une ligne se chevauche entre ce module
et le passage ultérieur avec --resume, la valeur _synced_at la plus récente l’emporte.
Quand choisir S3 en production. Si le jeu dépasse 500 millions de lignes, ou si le coût du warehouse Snowflake qui analyse toute la table devient significatif, préférez un export S3 : Snowflake exporte du Parquet compressé en parallèle (beaucoup plus vite qu’un curseur unique) et ClickHouse charge également depuis S3 en parallèle. Le script Python convient bien à l’échelle de l’atelier.
Alignement des décisions. Le tableau suivant reprend les décisions des fiches 1
(moteurs), 2 (clés de tri) et 3 (traduction du schéma) dans migration-plan.md, puis les
compare à ce que construit réellement l’atelier. Vérifiez votre plan avant de
provisionner quoi que ce soit :
| Décision | Mise en œuvre dans cet atelier | Pourquoi |
|---|---|---|
Moteur de trips_raw | ReplacingMergeTree(_synced_at) | Le script Python de migration effectue des INSERT par lots qui peuvent être retentés après une interruption. _synced_at DateTime DEFAULT now() est défini à chaque INSERT ; une ligne répétée arrive donc plus tard avec une valeur _synced_at plus élevée, l’emporte lors de la déduplication RMT et rend les reprises idempotentes. Les nouvelles tentatives du producteur après la bascule sont sûres pour la même raison. stg_trips interroge avec FINAL pour garantir une ligne par course. |
Moteur de fact_trips | ReplacingMergeTree(updated_at) | Les courses peuvent être corrigées (ajustements tarifaires) ; updated_at est la colonne de version |
Moteur de agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | Recalcul sur fenêtre glissante = modèle d’upsert |
Moteur des tables dim_* | MergeTree() | Rechargement complet à chaque exécution dbt ; aucun upsert |
ORDER BY de fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Les requêtes Q1 à Q7 filtrent toutes sur pickup_at ; trip_id assure l’unicité au dernier niveau |
ORDER BY de agg_hourly_zone_trips | (hour_bucket, zone_id) | Les deux colonnes figurent dans toutes les requêtes d’agrégation |
| VARIANT → | String + JSONExtract* | Préserve le JSON brut ; extraction lors de la requête |
| QUALIFY → | Sous-requête autour de ROW_NUMBER() | ClickHouse possède une clause QUALIFY native depuis la version 24.5, mais la forme en sous-requête reste enseignée car elle est portable vers les versions antérieures et les moteurs SQL sans QUALIFY |
| MERGE INTO → | Modèle incrémentiel delete_insert dans dbt | Stratégie d’upsert idiomatique de dbt-clickhouse ; évite de réécrire toute la table |
Étape 1 — Provisionner le cluster ClickHouse
setup.sh réalise une seule opération : il exécute terraform apply et écrit les
informations de connexion dans .clickhouse_state. Il vérifie aussi l’existence de
migration-plan.md avant le provisionnement — voir la section Pourquoi —, mais ne fait
qu’avertir et ne bloque jamais.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
# Configure credentials
cp .env.example .env
vim .env
# Fill in: CLICKHOUSE_ORG_ID, CLICKHOUSE_TOKEN_KEY, CLICKHOUSE_TOKEN_SECRET, CLICKHOUSE_PASSWORD
# Provision
source .env && ./setup.sh.env est ignoré par Git ; ne le validez jamais dans le dépôt.
Résultat attendu : Terraform crée 2 ressources (service et liste d’accès IP) en 2 à 3 minutes environ :
Apply complete! Resources: 2 added, 0 changed, 0 destroyed.
Outputs:
clickhouse_host = "abc123xyz.us-east-1.aws.clickhouse.cloud"
clickhouse_port = 8443L’hôte et le port sont enregistrés dans .clickhouse_state. Chargez ce fichier dans
n’importe quel terminal pour obtenir la connexion :
source .clickhouse_stateVérification :
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1Étape 2 — Créer les tables vides
Commencez par créer manuellement trips_raw avec le bon moteur. Le script de migration
y charge les données à l’étape 3. La table doit déjà exister avec
ReplacingMergeTree, afin que la colonne de version soit configurée avant l’arrivée de
la première ligne.
-- Run in the ClickHouse SQL console (cloud.clickhouse.com -> SQL console)
CREATE TABLE IF NOT EXISTS default.trips_raw (
trip_id String,
vendor_id UInt8,
pickup_at DateTime64(3, 'UTC'),
dropoff_at DateTime64(3, 'UTC'),
passenger_count UInt8,
trip_distance_miles Float32,
pickup_location_id UInt16,
dropoff_location_id UInt16,
payment_type_id UInt8,
rate_code_id UInt8,
store_fwd_flag String,
fare_amount_usd Float32,
extra_amount_usd Float32,
mta_tax_usd Float32,
tip_amount_usd Float32,
tolls_amount_usd Float32,
total_amount_usd Float32,
ingested_at DateTime64(3, 'UTC'),
trip_metadata String,
_synced_at DateTime DEFAULT now()
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id);_synced_at est défini automatiquement à chaque INSERT. Si le script de migration est
interrompu puis relancé avec --resume, des lignes en double peuvent brièvement exister
pour un même trip_id. RMT conserve la plus récente (celle dont _synced_at est le plus
élevé). stg_trips interroge trips_raw FINAL afin de forcer la déduplication avant que
les données ne soient visibles par un modèle en aval.
Chargez ensuite les données de référence des zones. Ces données statiques (265 zones NYC
TLC) sont lues comme source par stg_taxi_zones dans dbt.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
clickhouse-client --host "${CLICKHOUSE_HOST}" --port 9440 --secure \
--user default --password "${CLICKHOUSE_PASSWORD}" \
--multiquery < scripts/00_seed_zones.sql(Vous pouvez également coller directement le contenu de
scripts/00_seed_zones.sql dans la console SQL ClickHouse.)
Configurez le profil dbt. Le fichier dbt_project.yml de ce projet déclare
profile: 'nyc_taxi_ch'. Sans profil correspondant dans ~/.dbt/profiles.yml,
dbt run échoue immédiatement avec Could not find profile named 'nyc_taxi_ch'. Le
modèle se trouve dans
workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch/profiles.yml.example.
Le module 01 a déjà écrit dans ~/.dbt/profiles.yml un profil nyc_taxi: pour
Snowflake, et la boucle d’actualisation de l’étape 4 continue de l’utiliser tant que le
producteur Snowflake reste actif. Ne remplacez pas ce fichier par le modèle ClickHouse :
écraser le fichier avec profiles.yml.example effacerait le profil nyc_taxi: et
interromprait la boucle du module 01. Ouvrez plutôt le modèle et ajoutez le bloc
nyc_taxi_ch: au fichier
~/.dbt/profiles.yml existant comme second profil de premier niveau, à côté de
nyc_taxi: :
nyc_taxi: # from module 01 — leave this one alone
target: dev
outputs:
dev:
type: snowflake
# ...
nyc_taxi_ch: # add this block
target: dev
outputs:
dev:
type: clickhouse
schema: nyc_taxi_ch
host: "{{ env_var('CLICKHOUSE_HOST') }}"
port: 8443
user: "{{ env_var('CLICKHOUSE_USER', 'default') }}"
password: "{{ env_var('CLICKHOUSE_PASSWORD') }}"
secure: truenyc_taxi_ch: lit CLICKHOUSE_HOST, CLICKHOUSE_USER et CLICKHOUSE_PASSWORD dans
l’environnement avec env_var(). .env et .clickhouse_state doivent donc être
chargés avant toute commande dbt de ce module ; le dbt run ci-dessous le fait déjà.
Comme le profil Snowflake, ~/.dbt/profiles.yml contient des identifiants et est ignoré
par Git. Ne le validez jamais. Cette fusion ajoute un deuxième jeu d’identifiants à un
fichier qui contenait déjà le premier.
Vérification :
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt debug
# Expected: "All checks passed!" — confirms dbt found the nyc_taxi_ch profile and
# connected to ClickHouseExécutez ensuite dbt run pour créer les tables analytics et les vues de préparation :
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt deps # install packages (first run only)
dbt run # creates analytics tables and staging views; all empty at this pointRésultat attendu : environ 8 modèles créés en moins de 2 minutes (toutes les tables sont vides).
Vérification :
# Check analytics tables were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+analytics" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: agg_hourly_zone_trips, dim_date, dim_payment_type, dim_vendor, dim_taxi_zones, fact_trips
# Check staging views were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+staging" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: stg_trips, stg_taxi_zones
# Check trips_raw exists with the correct engine
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+engine+FROM+system.tables+WHERE+database%3D%27default%27+AND+name%3D%27trips_raw%27" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: ReplacingMergeTreeLes six tables analytics et les deux vues de préparation existent désormais, mais
restent vides : dbt run n’a fait que créer leur schéma. La seule table destinée à
contenir des données après ce module est trips_raw, elle-même encore vide à ce stade ;
les données arrivent maintenant.
Étape 3 — Migrer les données
Chargez toutes les lignes de NYC_TAXI_DB.RAW.TRIPS_RAW dans Snowflake vers
default.trips_raw dans ClickHouse à l’aide du script Python de migration par lots.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
source .venv/bin/activate
python scripts/02_migrate_trips.pySortie attendue (environ 40 à 50 minutes pour 50 millions de lignes) :
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
NYC Taxi Migration: Snowflake -> ClickHouse
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Rows to migrate: 50,000,000
Batch size: 100,000
Rows inserted Elapsed ETA Rate
-------------------- ------------ ---------------------- ---------------
100,000 0m 07s 56m 14s remaining 13,945 rows/s
200,000 0m 14s 55m 28s remaining 14,021 rows/s
...Le producteur Snowflake continue d’écrire dans TRIPS_RAW pendant toute l’exécution ;
ClickHouse prend donc un retard approximativement égal à la durée du transfert. Cet
écart est attendu et sera traité au module 05, pas ici.
Si le script est interrompu, relancez-le avec --resume pour reprendre au dernier point :
python scripts/02_migrate_trips.py --resume--resume lit max(pickup_at) dans ClickHouse et ignore les lignes déjà chargées. Vous
pouvez donc interrompre puis redémarrer le script sans risque de créer un chargement
partiel irrécupérable.
Comment vérifier que vous avez terminé
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
# Row count in trips_raw
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+count()+FROM+default.trips_raw" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: approximately 50000000
# .clickhouse_state was written by setup.sh
ls -la .clickhouse_state
# Service is reachable
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1État final
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 — sept objets au total
dans le schéma analytics, créés par le dbt run de ce module. Toutes les tables
analytics sont encore vides, à l’exception de mv_live_trip_feed, qui contient déjà la
ligne d’instantané produite lorsque dbt run a créé la vue. En dehors de cela, seule default.trips_raw contient des données :
environ 50 millions de lignes transférées à l’étape 3.
Le producteur Snowflake fonctionne toujours. Il n’a pas été arrêté dans ce module et
ne doit pas l’être ici. Chaque course écrite dans TRIPS_RAW après le dernier lot du
script est absente de ClickHouse. ClickHouse accuse donc un retard d’environ la durée de
la fenêtre de migration (40 à 50 minutes, plus le temps de préparation de ce module).
L’écart est réel et continue de croître tant que le producteur reste actif. Ne le
refermez pas dans ce module. La bascule du module 05 le ferme volontairement en deux
passages contrôlés et en mesure la taille avant de l’éliminer. Arrêter le producteur ou
relancer le script maintenant supprimerait précisément le phénomène que le module 05 est
conçu pour démontrer.
Fiche 5 : conception des modèles dbt
Configurez la matérialisation, le moteur et la stratégie incrémentielle de chaque modèle dbt, avec une correction immédiate pour chaque réponse.
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.