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.
| Mesure | Valeur |
|---|
| Nombre total de tables | 7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE) |
| Nombre total de vues | 2 (STG_TRIPS, STG_TAXI_ZONES) |
| Streams | 1 (TRIPS_CDC_STREAM sur TRIPS_RAW) |
| Tasks | 2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK) |
| Nombre total de lignes dans TRIPS_RAW | environ 50 000 000 |
| Plage de dates | Fenêtre glissante de 4 ans se terminant au moment de la préparation |
| Colonnes VARIANT | 1 (TRIPS_RAW.TRIP_METADATA) |
| Utilisations de QUALIFY détectées | 1, dans la requête Q3 |
| Utilisations de MERGE INTO détectées | 2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK) |
| Objet | Type | Schéma | Lignes | Niveau de complexité | Remarques |
|---|
trips_raw | Table | raw | environ 50 millions | B | RMT 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_trips | Vue dbt | staging | — | B | JSONExtract pour TRIP_METADATA ; les chemins JSON doivent être testés |
stg_taxi_zones | Vue dbt | staging | — | A | Transmission directe, triviale |
int_trips_enriched | Modèle dbt éphémère | staging | — | A | CTE ; les différences SQL sont traitées dans les modèles parents |
fact_trips | Modèle dbt incrémentiel | analytics | environ 50 millions | C | Moteur RMT ; delete_insert ; réécriture de QUALIFY ; FINAL obligatoire |
agg_hourly_zone_trips | Modèle dbt incrémentiel | analytics | environ 140 000 | B | RMT ; fenêtre glissante de recalcul de 2 heures ; tester soigneusement la limite de partition |
dim_taxi_zones | Table dbt | analytics | 265 | A | Référence statique ; rechargement complet ; triviale |
dim_payment_type | Table dbt | analytics | 6 | A | Référence statique ; triviale |
dim_vendor | Table dbt | analytics | 3 | A | Référence statique ; triviale |
taxi_zones_dict | Dictionnaire | analytics | 265 | B | Syntaxe propre à ClickHouse ; dictGet() au moment de la requête |
mv_hourly_revenue | Vue matérialisée actualisable | analytics | — | B | Syntaxe REFRESH EVERY ; vérifier le remplacement atomique |
TRIPS_CDC_STREAM / CDC_CONSUME_TASK | Stream + Task Snowflake | — | — | D | Aucun équivalent ClickHouse ; remplacés par la bascule directe du producteur dans la troisième partie |
| Table | Moteur | Colonne de version | Justification |
|---|
trips_raw | ReplacingMergeTree(_synced_at) | _synced_at | Le 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_trips | ReplacingMergeTree(updated_at) | updated_at | Les 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_trips | ReplacingMergeTree(updated_at) | updated_at | dbt 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_zones | MergeTree() | — | 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_type | MergeTree() | — | Même principe : rechargement complet. |
dim_vendor | MergeTree() | — | Même principe : rechargement complet. |
mv_hourly_revenue | MergeTree() | — | La vue matérialisée actualisable remplace atomiquement l’ensemble de ses résultats à chaque REFRESH. Aucun upsert. |
| Table | ORDER BY | Justification |
|---|
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. |
| Colonne | Type Snowflake | Type ClickHouse | Justification de la décision |
|---|
TRIP_METADATA | VARIANT | String | Conserve 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_AT | TIMESTAMP_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_ID | INTEGER | UInt16 | Les 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_ID | INTEGER | UInt8 | Les valeurs vont de 1 à 3. La limite maximale de UInt8, 255, convient, à raison d’un octet par ligne. |
DRIVER_RATING | FLOAT | Nullable(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_AT | TIMESTAMP_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. |
| 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::FLOAT | JSONExtractFloat(trip_metadata, 'driver', 'rating') |
TRIP_METADATA:surge_multiplier::FLOAT | JSONExtractFloat(trip_metadata, 'surge_multiplier') |
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10 | SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10 |
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATE | Mode incrémentiel delete_insert de dbt : DELETE des lignes dont les clés correspondent, puis INSERT de toutes les nouvelles lignes |
| Vague | Objets | Dépendances | Remarques |
|---|
| Vague 0 | Schéma de trips_raw, dim_taxi_zones, dim_payment_type, dim_vendor | Aucune | dbt 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 1 | Chargement en masse Python (scripts/02_migrate_trips.py) | Vague 0, car le schéma de trips_raw doit exister | 50 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 2 | stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_trips | Vague 1 terminée, trips_raw étant alimentée, et vague 0, les tables de dimensions existant | Exé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 3 | taxi_zones_dict, mv_live_trip_feed | Vague 2, avec dim_taxi_zones alimentée pour le dictionnaire et fact_trips pour la vue matérialisée | Dictionnaire 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 4 | Bascule du producteur (scripts/03_cutover.sh) | Vague 1 terminée et chargement en masse vérifié, plus vague 2 terminée et couche analytique construite | Arrê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. |
| Objet | Risque | Méthode de vérification |
|---|
fact_trips | Les 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_trips | La 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 producteur | Un 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. |
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.
| Critère | Seuil | Mesure |
|---|
| Parité du nombre de lignes | Correspondance ≥ 99,9 %, CH pouvant dépasser SF après la bascule | scripts/01_verify_migration.sh |
| Parité de la somme de contrôle | MD5 identique sur un échantillon de 10 000 lignes | scripts/02_validate_parity.sql |
| Taux de réussite des tests dbt | 100 % | dbt test dans dbt/nyc_taxi_dbt_ch |
| Parité des résultats de requêtes | Les 7 requêtes renvoient les mêmes résultats, dans la tolérance des nombres à virgule flottante | Comparaison manuelle dans le résultat de scripts/run_benchmark.sh |
| Modèle | Matérialisation | Justification |
|---|
stg_trips | view | Lit 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_zones | view | Même principe : nettoyage par transmission directe d’une table source |
int_trips_enriched | ephemeral | Logique 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_trips | incremental | Les 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_trips | incremental | Le recalcul glissant sur 2 heures est incrémentiel : il traite les lignes récentes, et non les 50 millions |
dim_taxi_zones | table | 265 zones statiques ; reconstruction complète avec échange atomique de table à chaque exécution dbt ; aucune mise à jour partielle |
dim_payment_type | table | 6 types statiques ; même raisonnement que pour dim_taxi_zones |
dim_vendor | table | 3 fournisseurs ; même raisonnement |
| Modèle | ENGINE | Colonne de version | Justification |
|---|
fact_trips | ReplacingMergeTree(updated_at) | updated_at | Les 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_trips | ReplacingMergeTree(updated_at) | updated_at | Le 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_zones | MergeTree() | — | 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_type | MergeTree() | — | Même principe que dim_taxi_zones. |
dim_vendor | MergeTree() | — | Même principe que dim_taxi_zones. |
| Modèle | unique_key | incremental_strategy | Filtre incrémentiel | Pourquoi ce filtre ? |
|---|
fact_trips | trip_id | delete_insert | WHERE 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_insert | WHERE pickup_at >= now() - INTERVAL 2 HOUR | La 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. |
| Modèle | FINAL dans la clause FROM ? | Justification |
|---|
stg_trips | Oui — FROM trips_raw FINAL | trips_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_enriched | Non | Lit stg_trips, qui est une vue, et non une table RMT ; FINAL ne s’applique pas aux vues. |
fact_trips | Non, dans le corps du modèle | delete_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.