Snowflake MigrationClickHouse Workshops

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.

Flux de migration : le producteur de courses Snowflake continue d’écrire dans TRIPS_RAW tandis qu’un script Python ponctuel transfère 50 millions de lignes par lots de 100 000 vers ClickHouse Cloud

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 :

Côté cible ClickHouse : trips_raw sur ReplacingMergeTree alimente les vues de préparation créées par dbt, les tables de faits et de dimensions, un agrégat horaire et un dictionnaire de zones

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éthodeFonctionnementPourquoi elle n’est pas utilisée ici
ClickPipes (source Snowflake)Connecteur natif ClickHouse Cloud — pas d’ETL, interface géréeSnowflake 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 S3COPY INTO @stage exporte du Parquet/CSV vers S3 ; le connecteur S3 de ClickPipes charge les données dans ClickHouseNé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-clientMê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 → ClickHouseUn stream CDC Snowflake alimente un topic Kafka ; le connecteur Kafka ClickPipes réalise l’ingestionPipeline 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 directementAucune 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. --resume permet 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 :

  1. Arrêtez le producteur Snowflake afin de figer le jeu de données.
  2. Exécutez python scripts/02_migrate_trips.py --resume : seules les lignes de l’écart sont transférées (quelques secondes, pas des minutes).
  3. 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écisionMise en œuvre dans cet atelierPourquoi
Moteur de trips_rawReplacingMergeTree(_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_tripsReplacingMergeTree(updated_at)Les courses peuvent être corrigées (ajustements tarifaires) ; updated_at est la colonne de version
Moteur de agg_hourly_zone_tripsReplacingMergeTree(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 dbtStraté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 = 8443

L’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_state

Vé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: true

nyc_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 ClickHouse

Exé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 point

Ré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: ReplacingMergeTree

Les 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.py

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

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