Snowflake MigrationClickHouse Workshops

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.

Voici le corrigé complet des cinq fiches (1, 2, 3, 4, 5) appliquées à la charge de travail NYC Taxi. Utilisez-le pour :

  • vérifier vos réponses après avoir terminé chaque section ;
  • comprendre le raisonnement derrière les décisions mises en œuvre dans la troisième partie ;
  • comparer vos choix avec le tableau d’alignement des décisions de la troisième partie s’ils diffèrent.

Ceci est le corrigé : ne le remplissez pas comme s’il s’agissait de votre plan. Remplissez plutôt migration-plan.md.


Liste de contrôle d’achèvement

  • Choix des moteurs : terminé
  • Conception des clés de tri : terminée
  • Traduction du schéma : terminée
  • Plan des vagues de migration : terminé
  • Conception des modèles dbt : terminée

Section 1 : synthèse du profil

MesureValeur
Nombre total de tables7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE)
Nombre total de vues2 (STG_TRIPS, STG_TAXI_ZONES)
Streams1 (TRIPS_CDC_STREAM sur TRIPS_RAW)
Tasks2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK)
Nombre total de lignes dans TRIPS_RAWenviron 50 000 000
Plage de datesFenêtre glissante de 4 ans se terminant au moment de la préparation
Colonnes VARIANT1 (TRIPS_RAW.TRIP_METADATA)
Utilisations de QUALIFY détectées1, dans la requête Q3
Utilisations de MERGE INTO détectées2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK)

Section 2 : inventaire des objets

ObjetTypeSchémaLignesNiveau de complexitéRemarques
trips_rawTablerawenviron 50 millionsBRMT avec la colonne de version _synced_at ; le chevauchement du chargement en masse et du CDC impose une déduplication ; stg_trips doit employer FINAL
stg_tripsVue dbtstaging—BJSONExtract pour TRIP_METADATA ; les chemins JSON doivent être testés
stg_taxi_zonesVue dbtstaging—ATransmission directe, triviale
int_trips_enrichedModèle dbt éphémèrestaging—ACTE ; les différences SQL sont traitées dans les modèles parents
fact_tripsModèle dbt incrémentielanalyticsenviron 50 millionsCMoteur RMT ; delete_insert ; réécriture de QUALIFY ; FINAL obligatoire
agg_hourly_zone_tripsModèle dbt incrémentielanalyticsenviron 140 000BRMT ; fenêtre glissante de recalcul de 2 heures ; tester soigneusement la limite de partition
dim_taxi_zonesTable dbtanalytics265ARéférence statique ; rechargement complet ; triviale
dim_payment_typeTable dbtanalytics6ARéférence statique ; triviale
dim_vendorTable dbtanalytics3ARéférence statique ; triviale
taxi_zones_dictDictionnaireanalytics265BSyntaxe propre à ClickHouse ; dictGet() au moment de la requête
mv_hourly_revenueVue matérialisée actualisableanalytics—BSyntaxe REFRESH EVERY ; vérifier le remplacement atomique
TRIPS_CDC_STREAM / CDC_CONSUME_TASKStream + Task Snowflake——DAucun équivalent ClickHouse ; remplacés par la bascule directe du producteur dans la troisième partie

Section 3 : décisions relatives aux moteurs

TableMoteurColonne de versionJustification
trips_rawReplacingMergeTree(_synced_at)_synced_atLe script de migration Python (scripts/02_migrate_trips.py) peut relancer un lot et réinsérer le même trip_id. Après la bascule, le producteur en direct peut également réessayer lors d’une erreur temporaire. _synced_at DateTime DEFAULT now() est défini automatiquement lors de l’INSERT ; une nouvelle tentative ultérieure porte donc un horodatage supérieur et RMT conserve l’écriture la plus récente. stg_trips interroge avec FINAL afin d’imposer la déduplication avant l’exécution de tout modèle en aval.
fact_tripsReplacingMergeTree(updated_at)updated_atLes trajets peuvent être corrigés, par exemple pour ajuster un tarif ou changer un statut. Le même trip_id est réinséré avec les nouvelles valeurs. updated_at augmente de façon monotone à chaque correction ; la valeur supérieure l’emporte pendant la déduplication RMT. Toujours interroger avec FINAL.
agg_hourly_zone_tripsReplacingMergeTree(updated_at)updated_atdbt recalcule les 2 dernières heures et les réinsère. Sans RMT, les anciens et nouveaux agrégats s’accumulent et sont comptés deux fois. updated_at, défini sur now() à chaque exécution dbt, garantit que les valeurs les plus récentes l’emportent.
dim_taxi_zonesMergeTree()—Rechargement complet par dbt, avec échange atomique de table. Aucun doublon ne peut s’accumuler et aucune déduplication n’est nécessaire.
dim_payment_typeMergeTree()—Même principe : rechargement complet.
dim_vendorMergeTree()—Même principe : rechargement complet.
mv_hourly_revenueMergeTree()—La vue matérialisée actualisable remplace atomiquement l’ensemble de ses résultats à chaque REFRESH. Aucun upsert.

Section 4 : conception des clés de tri

TableORDER BYJustification
trips_raw(pickup_at, trip_id)Les lectures par plage temporelle filtrent d’abord sur pickup_at. trip_id est la clé de déduplication RMT et doit figurer dans ORDER BY pour permettre à RMT d’identifier les doublons. pickup_at vient en premier car les lectures analytiques par plage dominent ; trip_id vient en dernier car sa cardinalité est forte et qu’il ne sert qu’à distinguer les valeurs uniques.
fact_trips(toStartOfMonth(pickup_at), pickup_at, trip_id)Les 7 requêtes analytiques filtrent sur pickup_at. Le préfixe mensuel regroupe les données d’un mois civil dans des blocs adjacents et permet d’ignorer grossièrement des blocs pour les agrégations mensuelles sans ajouter PARTITION BY. trip_id vient en dernier pour garantir l’unicité RMT sans nuire à l’élimination des blocs.
agg_hourly_zone_trips(hour_bucket, zone_id)Q6, comme toutes les requêtes d’agrégation, filtre sur hour_bucket et zone_id. hour_bucket comporte environ 35 000 valeurs distinctes et zone_id en comporte 265. hour_bucket vient en premier car la lecture par plage temporelle est le principal mode d’accès ; zone_id vient ensuite pour le filtrage secondaire.
dim_taxi_zones(location_id)Avec 265 lignes, la table tient dans un granule et ORDER BY n’influe pas sur les performances. Utiliser location_id, la clé de jointure, est conventionnel et facilite la lecture.

Section 5 : remarques sur la traduction du schéma

ColonneType SnowflakeType ClickHouseJustification de la décision
TRIP_METADATAVARIANTStringConserve exactement le JSON brut. JSONExtract* gère des chemins arbitraires au moment de la requête. Map(String,String) perd les structures imbriquées et Tuple exige un schéma fixe. String est le choix sûr pour un JSON arbitraire.
PICKUP_DATETIME / PICKUP_ATTIMESTAMP_NTZ(9)DateTime64(3, 'UTC')Une précision à la milliseconde suffit aux horodatages des trajets ; la nanoseconde (9) est excessive. 'UTC' rend le fuseau horaire explicite et évite les surprises liées à l’heure d’été dans les agrégations par plage temporelle.
PICKUP_LOCATION_IDINTEGERUInt16Les valeurs vont de 1 à 265. La limite maximale de UInt8, 255, est insuffisante ; celle de UInt16, 65535, convient. Deux octets au lieu de quatre pour Int32 économisent environ 95 Mo non compressés par colonne sur 50 millions de lignes.
VENDOR_IDINTEGERUInt8Les valeurs vont de 1 à 3. La limite maximale de UInt8, 255, convient, à raison d’un octet par ligne.
DRIVER_RATINGFLOATNullable(Float32)Souvent NULL, car tous les trajets ne comportent pas de note. Nullable conserve la bonne sémantique des valeurs nulles. Float32 suffit à une plage de 1,0 à 5,0 ; Float64 gaspillerait du stockage sans précision utile.
UPDATED_ATTIMESTAMP_NTZ(9)DateTime64(3, 'UTC')Colonne de version de ReplacingMergeTree. Elle doit utiliser DateTime64 plutôt que DateTime : deux corrections pendant la même seconde seraient indéterminées avec une précision à la seconde. La milliseconde garantit le bon ordre de déduplication.

Traductions de fonctions requises

Expression SnowflakeÉquivalent ClickHouse
DATE_TRUNC('hour', pickup_at)toStartOfHour(pickup_at)
DATEADD('day', -7, CURRENT_DATE)today() - 7
DATEDIFF('minute', pickup_at, dropoff_at)dateDiff('minute', pickup_at, dropoff_at)
TRIP_METADATA:driver.rating::FLOATJSONExtractFloat(trip_metadata, 'driver', 'rating')
TRIP_METADATA:surge_multiplier::FLOATJSONExtractFloat(trip_metadata, 'surge_multiplier')
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATEMode incrémentiel delete_insert de dbt : DELETE des lignes dont les clés correspondent, puis INSERT de toutes les nouvelles lignes

Section 6 : vagues de migration

VagueObjetsDépendancesRemarques
Vague 0Schéma de trips_raw, dim_taxi_zones, dim_payment_type, dim_vendorAucunedbt crée des tables vides. Les tables de dimensions sont immédiatement alimentées à partir des données de référence statiques, sans dépendre des trajets. Exécution : dbt run --select trips_raw dim_*
Vague 1Chargement en masse Python (scripts/02_migrate_trips.py)Vague 0, car le schéma de trips_raw doit exister50 millions de lignes provenant de TRIPS_RAW dans Snowflake. Reprise possible avec --resume. Vérifier le nombre de lignes avec scripts/01_verify_migration.sh. Environ 40 à 50 minutes.
Vague 2stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_tripsVague 1 terminée, trips_raw étant alimentée, et vague 0, les tables de dimensions existantExécution complète de dbt run. stg_trips lit trips_raw ; int_trips_enriched effectue les jointures avec les dimensions ; fact_trips et agg_hourly_zone_trips sont construits au-dessus.
Vague 3taxi_zones_dict, mv_live_trip_feedVague 2, avec dim_taxi_zones alimentée pour le dictionnaire et fact_trips pour la vue matérialiséeDictionnaire créé avec scripts/04_create_dictionary.sql. Vue matérialisée actualisable créée par le modèle dbt ; l’activation ultérieure de sa fréquence d’actualisation constitue une étape manuelle ALTER TABLE ... MODIFY REFRESH, et non une action exécutée automatiquement par dbt.
Vague 4Bascule du producteur (scripts/03_cutover.sh)Vague 1 terminée et chargement en masse vérifié, plus vague 2 terminée et couche analytique construiteArrêter le producteur Snowflake, démarrer le producteur ClickHouse qui écrit directement dans ClickHouse Cloud, puis exécuter dbt run afin d’alimenter agg_hourly_zone_trips avec les données en direct.

Registre des risques, objets de niveau C ou D

ObjetRisqueMéthode de vérification
fact_tripsLes requêtes sans FINAL surestiment les résultats pendant le retard de fusion. La plage de partitions de delete_insert doit suivre le préfixe ORDER BY afin de ne pas supprimer des partitions hors cible.SELECT COUNT(*) FINAL correspond à Snowflake, au décalage CDC près. Exécuter dbt test. Comparer les résultats de Q3 entre les systèmes. Rechercher des trip_id en double avec SELECT trip_id, count() FROM fact_trips GROUP BY trip_id HAVING count() > 1 LIMIT 10.
agg_hourly_zone_tripsLa fenêtre glissante de recalcul de 2 heures doit délimiter correctement la plage supprimée. Trop large, elle supprime d’anciens agrégats ; trop étroite, elle conserve des agrégats obsolètes.Contrôler quelques tuples (hour_bucket, zone_id) par rapport à Snowflake. Vérifier que le total trip_count de toutes les zones correspond à Snowflake AGG_HOURLY_ZONE_TRIPS sur la même période.
Bascule du producteurUn script de migration interrompu laisse un écart de lignes ; le relancer avec --resume pour le combler. Après la bascule, une nouvelle tentative du producteur peut réinsérer des trajets déjà présents dans ClickHouse.scripts/01_verify_migration.sh vérifie la parité du nombre de lignes entre Snowflake et ClickHouse. ReplacingMergeTree(_synced_at) rend les insertions en double idempotentes.

Section 7 : écarts de dialecte connus

  • QUALIFY — affecte Q3 (queries/q03_top_trips_qualify.sql)
  • Chemin VARIANT avec deux-points — affecte Q4 et Q5, accès au JSON TRIP_METADATA
  • LATERAL FLATTEN — absent de cette charge ; VARIANT est lu par chemin avec deux-points et non avec FLATTEN
  • MERGE INTO — affecte les modèles dbt incrémentiels fact_trips et agg_hourly_zone_trips
  • Streams Snowflake → bascule du producteur, les écritures en direct allant directement dans ClickHouse après la bascule
  • Différences entre les fonctions de date — affectent Q1 avec DATE_TRUNC, Q3 avec DATEADD et Q4 avec DATEDIFF

Section 8 : stratégie de migration

Déplacement des données : script de migration Python (scripts/02_migrate_trips.py)

Pourquoi un script Python plutôt qu’un relais par stockage d’objets ou ClickPipes ?

  • remoteSecure() sert au transfert de ClickHouse vers ClickHouse et ne s’applique pas ici.
  • Un relais par stockage d’objets, de Snowflake vers S3 puis vers la fonction de table S3 de ClickHouse, fonctionnerait, mais ajouterait de la complexité : provisionnement d’un compartiment S3, rôles IAM et COPY INTO Snowflake, soit un surcoût inutile pour un laboratoire.
  • ClickPipes ne prend pas Snowflake en charge comme source. Ses sources compatibles sont Kafka, S3, Kinesis, CDC PostgreSQL et CDC MySQL.
  • Le script Python emploie snowflake-connector-python et clickhouse-connect, déjà installés pour le laboratoire. Il affiche la progression en temps réel, accepte --resume après une interruption et son code est entièrement consultable.

Stratégie incrémentielle dbt : delete_insert

Pourquoi delete_insert plutôt que les stratégies append ou merge ?

  • append insère de nouvelles lignes sans toucher aux anciennes. Pour fact_trips, dont les lignes peuvent être mises à jour, cela crée des doublons et serait incorrect.
  • merge, si disponible, serait le plus proche du MERGE INTO de Snowflake, mais la stratégie merge de dbt-clickhouse présente des limites avec ReplacingMergeTree et n’est pas recommandée.
  • delete_insert supprime les lignes de la plage de clés du lot entrant, puis insère toutes les nouvelles lignes. Cette opération est idempotente, car sa réexécution produit le même résultat, gère les insertions et les mises à jour et fonctionne correctement avec ReplacingMergeTree. C’est la recommandation standard de la communauté dbt-clickhouse pour les modèles d’upsert.

Section 9 : critères de bascule

CritèreSeuilMesure
Parité du nombre de lignesCorrespondance ≥ 99,9 %, CH pouvant dépasser SF après la basculescripts/01_verify_migration.sh
Parité de la somme de contrôleMD5 identique sur un échantillon de 10 000 lignesscripts/02_validate_parity.sql
Taux de réussite des tests dbt100 %dbt test dans dbt/nyc_taxi_dbt_ch
Parité des résultats de requêtesLes 7 requêtes renvoient les mêmes résultats, dans la tolérance des nombres à virgule flottanteComparaison manuelle dans le résultat de scripts/run_benchmark.sh


Section 10 : conception des modèles dbt

Choix de la matérialisation

ModèleMatérialisationJustification
stg_tripsviewLit et nettoie trips_raw ; aucune mise à jour du modèle ; aucun coût de stockage ; reflète toujours l’état courant de la source
stg_taxi_zonesviewMême principe : nettoyage par transmission directe d’une table source
int_trips_enrichedephemeralLogique de jointure pure, utilisée uniquement par fact_trips ; son insertion sous forme de CTE évite une table physique redondante ; aucun modèle ne l’interroge directement
fact_tripsincrementalLes trajets peuvent être corrigés après coup ; seules les lignes nouvelles et modifiées doivent être traitées à chaque exécution
agg_hourly_zone_tripsincrementalLe recalcul glissant sur 2 heures est incrémentiel : il traite les lignes récentes, et non les 50 millions
dim_taxi_zonestable265 zones statiques ; reconstruction complète avec échange atomique de table à chaque exécution dbt ; aucune mise à jour partielle
dim_payment_typetable6 types statiques ; même raisonnement que pour dim_taxi_zones
dim_vendortable3 fournisseurs ; même raisonnement

Configuration des moteurs

ModèleENGINEColonne de versionJustification
fact_tripsReplacingMergeTree(updated_at)updated_atLes trajets peuvent être corrigés ; updated_at, défini sur now() à chaque insertion, permet à la version la plus récente de l’emporter lors de la déduplication RMT en arrière-plan. delete_insert assure principalement l’exactitude et RMT sert de filet de sécurité.
agg_hourly_zone_tripsReplacingMergeTree(updated_at)updated_atLe recalcul glissant réinsère des agrégats pour les mêmes paires (hour_bucket, zone_id) ; RMT garantit la suppression des agrégats obsolètes pendant la fusion en arrière-plan.
dim_taxi_zonesMergeTree()—Le rechargement complet par dbt entraîne un échange atomique de table à chaque exécution ; aucun doublon ne peut s’accumuler et aucune déduplication n’est nécessaire.
dim_payment_typeMergeTree()—Même principe que dim_taxi_zones.
dim_vendorMergeTree()—Même principe que dim_taxi_zones.

Stratégie incrémentielle

Modèleunique_keyincremental_strategyFiltre incrémentielPourquoi ce filtre ?
fact_tripstrip_iddelete_insertWHERE updated_at > (SELECT max(updated_at) FROM {{ this }})Le repère supérieur sur updated_at capture les nouveaux trajets comme les trajets corrigés, car un ajustement de tarif réinsère le même trip_id avec le même pickup_at, mais un updated_at plus récent. Un repère sur pickup_at omettrait silencieusement les corrections.
agg_hourly_zone_trips[hour_bucket, zone_id]delete_insertWHERE pickup_at >= now() - INTERVAL 2 HOURLa fenêtre glissante de 2 heures force la réagrégation des heures en limite afin de toujours corriger le décompte d’une heure partielle. Un repère supérieur max(pickup_at) sous-estimerait définitivement l’heure en limite.

Placement de FINAL

ModèleFINAL dans la clause FROM ?Justification
stg_tripsOui — FROM trips_raw FINALtrips_raw est une ReplacingMergeTree ; elle peut contenir plusieurs lignes avec le même trip_id après de nouvelles tentatives du script de migration ou du producteur après la bascule. stg_trips constitue l’unique point d’application : la déduplication s’effectue ici pour que chaque modèle en aval, int_trips_enriched, fact_trips et agg_hourly_zone_trips, reçoive des données propres.
int_trips_enrichedNonLit stg_trips, qui est une vue, et non une table RMT ; FINAL ne s’applique pas aux vues.
fact_tripsNon, dans le corps du modèledelete_insert maintient fact_trips propre à la fin de chaque exécution réussie. Ajouter FINAL dans le modèle l’appliquerait inutilement à la sous-requête is_incremental() qui lit max(updated_at) depuis {{ this }}. Les tableaux de bord et les tests dbt emploient FINAL à l’extérieur lorsqu’ils interrogent directement fact_trips.

Ceci est l’exemple terminé. Votre fichier migration-plan.md doit reprendre les décisions essentielles présentées ici, ou expliquer explicitement pourquoi vos choix diffèrent.

Sur cette page

FR