Google PDE: Ingestion, intégration et migration de données — Guide d'étude
Fait partie du Google Professional Data Engineer — Guide d’étude. Entraînez-vous avec des réponses vérifiées dans le centre d’examens Google, ou passez des tests chronométrés sur ExamRoll.io.
Vue d’ensemble
L’ingestion, l’intégration et la migration de données dans Google Cloud englobent des patterns reproductibles, des services gérés et des contrôles opérationnels qui transforment des systèmes sources variés en ensembles de données fiables et interrogeables. Les conceptions efficaces séparent le transport de la transformation, découplent les producteurs des consommateurs, et privilégient des pipelines idempotents avec points de contrôle, disposant d’un lignage et d’une vérification clairs. Cette section couvre les patterns d’ingestion, les services Google Cloud pour le mouvement de données et le CDC, les contrôles de schéma et de qualité des données, la connectivité et l’intégration hybride, ainsi que les stratégies de basculement, en soulignant les compromis de conception et les modes de défaillance.
Patterns d’ingestion et charges de travail
- Ingestion par lots (Batch) : Extractions périodiques ou dépôts de fichiers à des intervalles définis. Idéal pour un coût prévisible et les remplissages de données historiques (backfills). Mode de défaillance : des lots volumineux et peu fréquents provoquent des pics de ressources, de longues fenêtres de rattrapage et le non-respect des SLA. Atténuation : dimensionner correctement les fenêtres de lots, partitionner (shard) par temps ou par clé, et utiliser le parallélisme.
- Chargement en masse (Bulk) : Chargements uniques ou à grande échelle (par ex., remplissage historique initial). Préférez les formats colonnaires ou auto-descriptifs (Parquet, Avro) et chargez directement vers un stockage analytique (BigQuery) ou une zone de transit (staging) dans Cloud Storage. Compromis : l’interrogation de tables externes évite les étapes de chargement mais déplace le coût vers l’analyse au moment de la requête.
- Chargement incrémental : Chargements périodiques de deltas via des horodatages ou des filigranes (high-water marks). Nécessite un dédoublonnage robuste et des opérations d’insertion/mise à jour (upserts) idempotentes. Mode de défaillance : désynchronisation des horloges ou enregistrements arrivant en retard. Utilisez des horodatages de commit côté serveur et le watermarking.
- Capture de données modifiées (CDC) : Réplication continue des insertions, mises à jour et suppressions depuis des bases de données opérationnelles. Idéal pour les analyses en quasi-temps réel et les migrations à faible temps d’arrêt. Compromis :
- Ordonnancement : La plupart des outils CDC préservent l’ordre au sein des transactions et généralement au sein d’un shard, mais ne garantissent pas un ordre global entre les shards. Utilisez les horodatages de commit de transaction et les clés primaires pour reconstituer la séquence.
- Sémantique de livraison : La sémantique « au moins une fois » (at-least-once) est typique ; construisez des récepteurs (sinks) idempotents ou dédoublonnez en utilisant des identifiants de changement uniques.
- Instantané (Snapshot) + CDC : Commencez avec un instantané cohérent, puis appliquez les changements à partir d’une séquence de journal précise pour atteindre la parité sans temps d’arrêt.
Sources relationnelles, SaaS, sur site (on-premises) et fichiers :
- Sources relationnelles : Utilisez le CDC natif ou des colonnes d’horodatage. Pour le chargement en masse, exportez en Avro/Parquet et stockez en zone intermédiaire dans Cloud Storage.
- Sources SaaS : Préférez les API des fournisseurs avec des jetons incrémentaux ; intégrez via des connecteurs gérés (par ex., dans Data Fusion). Limitez le débit (throttle) pour respecter les quotas et gérez la dérive de schéma (schema drift).
- Sources sur site (on-prem) : Choisissez entre un transfert basé sur un agent, VPN/Interconnect + Private Google Access, ou un amorçage hors ligne avec Transfer Appliance.
- Ingestion de fichiers : Pour de nombreux petits fichiers, regroupez-les (par ex., tar) pour réduire la surcharge des appels RPC. Utilisez
gsutil -mou des clients parallélisés ; composez ou transformez en fichiers colonnaires plus volumineux pour l’analyse.
Services Google Cloud pour l’ingestion, l’intégration et la migration
- Datastream (CDC sans serveur) : Capture les changements de MySQL, PostgreSQL et Oracle vers Cloud Storage, BigQuery (via des modèles) ou Pub/Sub. Il préserve les limites de transaction et les métadonnées de commit ; l’ordonnancement global n’est pas garanti. Appliquez un ordonnancement en aval par clé et par horodatage de commit. Attendez-vous à une livraison « au moins une fois » (at-least-once) ; concevez des consommateurs idempotents (par ex.,
MERGEdans BigQuery avec des identifiants de changement). - Database Migration Service (DMS) : Pour les migrations de bases de données avec un temps d’arrêt minimal en utilisant la réplication native. DMS crée un instantané cohérent, puis réplique continuellement les changements en utilisant GTID/LSN/SCN. Il est conçu spécifiquement pour le lift-and-shift, pas pour la transformation arbitraire. Pour l’analytique, complétez DMS avec Dataflow ou Data Fusion si nécessaire.
- Cloud Data Fusion : Un service d’intégration géré avec des connecteurs vers des systèmes relationnels, SaaS, de fichiers et de messagerie. Construisez des pipelines avec des étapes de transformation (jointures, agrégations, conversions de format, recettes Wrangler personnalisées) et capturez le lignage des données à travers les sources et les champs. Sur le plan opérationnel, il planifie, relance et émet des métriques. Utilisez Data Fusion pour l’ELT/ETL sans code ou avec peu de code (no/low-code) et pour centraliser la gestion des connecteurs.
- Storage Transfer Service (STS) : Transferts gérés et planifiés depuis AWS S3, Azure Blob, des sources sur site (à l’aide d’agents), SFTP et des listes d’URL vers Cloud Storage. Prend en charge les manifestes, la synchronisation incrémentale, le contrôle de la bande passante et l’intégrité vérifiée par somme de contrôle (checksum). Les modes de défaillance incluent l’inefficacité avec les petits fichiers et la limitation du débit des API (throttling) ; atténuez avec le regroupement par lots (batching) et une simultanéité (concurrency) ajustable.
- Transfer Appliance : Appliance hors ligne et chiffrée pour l’amorçage initial de plusieurs téraoctets à plusieurs pétaoctets lorsque la bande passante réseau est limitée ou que les données sont trop sensibles pour un transit prolongé. La chaîne de possession et le chiffrement sont intégrés. Après l’amorçage, poursuivez avec STS ou le CDC pour les deltas.
- Cloud Pub/Sub + Dataflow : Pub/Sub découple les producteurs et les consommateurs pour des modèles de streaming ou de micro-lots. Dataflow offre un traitement de flux/lots avec état (stateful), auto-adaptatif (autoscaled), avec points de contrôle (checkpointing) et filigranes (watermarking). Utilisez la BigQuery Storage Write API pour un streaming à faible latence avec des garanties « exactement une fois » (exactly-once) par flux par défaut ; sinon, appuyez-vous sur la sémantique de dédoublonnage de l’insertId.
Pour les migrations de Hadoop vers Dataproc, minimisez l’utilisation de Persistent Disk en stockant les données dans Cloud Storage avec le connecteur GCS et utilisez des clusters éphémères ou à autoscaling. Cela évite les coûts élevés de stockage par blocs tout en préservant une sémantique compatible HDFS pour le traitement.
Schéma, validation et qualité des données à la frontière
- Mappage de schéma et conversion de type : Standardisez tôt vers des schémas fortement typés. Avro ou Parquet préservent le schéma et évoluent proprement. Dans BigQuery, privilégiez les tables partitionnées et clusterisées pour réduire le coût des analyses (scans). Exemple : créer une table partitionnée pour une analyse quotidienne
undefined
- Gestion des enregistrements malformés : Acheminez les rejets vers une file d’attente de lettres mortes (dead-letter queue) (Pub/Sub) ou un bucket de quarantaine dans Cloud Storage. Utilisez les sorties secondaires (side outputs) dans Dataflow ou les collecteurs d’erreurs dans Data Fusion. Journalisez les erreurs d’analyse (parsing) avec des exemples de charges utiles (payloads) et les versions de schéma pour le triage.
- Validation : Effectuez des contrôles aux limites avant la persistance :
- Structurelle : conformité du schéma, champs obligatoires, types de données, domaines d’énumération (enum).
- Référentielle : existence de clés étrangères via des recherches dans des dimensions mises en cache.
- Plausibilité : plages pour les horodatages, barrières géographiques (geofences), montants non négatifs.
- Unicité : collisions de clés primaires ou de clés composites.
- Chargement idempotent : Utilisez des clés déterministes et des opérations d’upsert. Dans BigQuery, implémentez MERGE avec une clé de modification naturelle ou de substitution (surrogate). Exemple :
undefined
- Filigranes (Watermarking) et retard : Dans les pipelines de streaming, configurez les filigranes basés sur le temps de l’événement (event-time) et le retard autorisé (allowed lateness) pour équilibrer la complétude et la latence. Les données en retard sont acheminées vers des chemins correctifs ou déclenchent des rattrapages (backfills).
- Réconciliation : Suivez le nombre de lignes et les sommes de contrôle (checksums) par partition/fenêtre de la source à la destination (sink). Capturez les positions des journaux CDC (LSN/SCN) et les horodatages de commit ; stockez-les dans une table de contrôle pour prouver la continuité et identifier les lacunes.
Connectivité, fiabilité et opérations
Connectivité réseau et accès privé :
- Hybride : Utilisez Cloud VPN ou Dedicated/Partner Interconnect pour la connectivité privée. Activez Private Google Access ou Private Service Connect pour un accès privé aux API Google comme Cloud Storage.
- Sécurité : Utilisez des comptes de service pour l’identité des charges de travail (workload identity), le principe de moindre privilège IAM, VPC Service Controls pour la prévention de l’exfiltration de données, et CMEK si nécessaire.
- Débit : Augmentez le parallélisme côté client, mais c’est finalement la bande passante qui régit le débit. Pour les transferts massifs, préférez Transfer Appliance pour le chargement initial en masse, puis STS ou CDC pour les mises à jour incrémentielles.
Points de contrôle (Checkpoints) et contre-pression (backpressure) : Dataflow gère les points de contrôle et l’autoscaling ; concevez des récepteurs (sinks) capables d’absorber les pics de charge (mise en tampon vers Cloud Storage, écritures par lots vers BigQuery). Pour Pub/Sub, ajustez le contrôle de flux et les délais d’acquittement (ack deadlines) pour éviter les tempêtes de redistribution de messages.
Ordonnancement et cohérence avec le CDC :
- Datastream préserve l’ordre intra-transactionnel et émet des métadonnées de commit ; les consommateurs reconstruisent l’ordre par clé en utilisant les horodatages de commit. Attendez-vous à une sémantique “au moins une fois” (at-least-once) ; construisez l’idempotence.
- DMS assure la cohérence de la base de données lors du basculement de l’instantané (snapshot) à la réplication en utilisant les journaux natifs. Utilisez des réplicas en lecture ou des stratégies de double écriture pour un basculement progressif.
Stratégie de fichiers pour l’analytique : Pour un accès multi-moteur à grande échelle, stockez les données canoniques dans Cloud Storage et, lorsque c’est rentable, exposez des tables externes permanentes pour les requêtes ad hoc. Pour l’analytique de production, chargez les données dans des tables BigQuery partitionnées pour minimiser le coût d’analyse par requête.
Optimisation des petits fichiers : Regroupez les petits fichiers (par ex., ~1 000 par archive .tar) avant le transfert, puis décompressez-les dans le cloud. Utilisez gsutil en parallèle et des règles de cycle de vie pour hiérarchiser (tier) et faire expirer les artefacts de pré-production (staging).
Pièges opérationnels et mesures d’atténuation :
- Dérive de schéma depuis un SaaS : activez l’évolution de schéma dans Data Fusion et imposez la compatibilité. Alertez sur les changements non rétrocompatibles (breaking changes).
- Fuseau horaire et encodage : normalisez en UTC et UTF-8 à l’entrée (ingress).
- Lacunes dans le CDC : surveillez la rétention des journaux sources ; alertez lorsque le retard du réplica (replica lag) approche des limites de rétention.
- Quotas : insertion en streaming dans BigQuery, limites de taux d’API ; traitez par lots à l’approche des limites.
Basculement, Rattrapage et Vérification
- Planification du basculement :
- Big bang : gel court, basculement unique. Complexité opérationnelle la plus faible ; risque le plus élevé si un rollback est nécessaire.
- Progressif ou blue/green : exécution double avec écritures en miroir, déplacement progressif du trafic et lectures en mode shadow. Coût plus élevé ; rollback plus sûr.
- Rattrapage :
- Effectuer un chargement en masse initial (Transfer Appliance ou STS) en utilisant Avro/Parquet pour préserver le schéma. Partitionner et clusteriser pendant le chargement pour éviter de retravailler les données.
- Démarrer le CDC (Change Data Capture) à une position connue du journal en même temps que le snapshot pour capturer les deltas pendant le transfert en masse. Réconcilier à un watermark commun avant d’ouvrir à la production.
- Rollback :
- Maintenir l’ancien système en lecture seule pendant la vérification. Pour les scénarios de double écriture, contrôler les écritures via un feature flag pour revenir en arrière rapidement. Conserver un checkpoint cohérent pour rejouer ou annuler les changements du CDC si nécessaire.
- Vérification de la migration :
- Structurelle : les nombres de lignes et les checksums par partition correspondent ; le schéma et les contraintes sont équivalents.
- Temporelle : pas de lacunes entre la limite du snapshot et le basculement ; les positions du CDC sont continues.
- Parité métier : comparer les agrégats et les indicateurs de performance (KPI) sur des fenêtres de temps ; exécuter des requêtes de validation.
- Performance : valider le débit d’ingestion, la latence des requêtes et le coût par rapport aux budgets.
Scénario de problème pratique
Northstar Retail doit consolider un ensemble mondial de systèmes transactionnels Oracle et MySQL sur site, d’événements CRM SaaS et de dépôts CSV quotidiens dans Google Cloud pour alimenter des analyses en quasi-temps réel et du machine learning. Ils doivent également migrer un cluster Hadoop hérité sans encourir de dépenses élevées en stockage par blocs, et réaliser un basculement avec un temps d’arrêt nul ou faible.
- Établir une connectivité hybride sécurisée
- Utiliser Partner Interconnect pour la bande passante principale et Cloud VPN comme solution de secours. Activer Private Google Access pour que les charges de travail sur site puissent accéder à Cloud Storage et Pub/Sub de manière privée. Justification : Les chemins privés minimisent l’exposition en sortie (egress) et la latence, et Private Google Access évite les exigences d’adresses IP publiques tout en respectant la politique de sécurité.
- Injecter les données historiques efficacement
- Pour 800 To de données HDFS historiques, copier vers Cloud Storage en utilisant Transfer Appliance (masse initiale). Après l’injection initiale, exécuter Storage Transfer Service quotidiennement depuis l’export NFS sur site pour récupérer les changements jusqu’au basculement. Justification : Transfer Appliance évite une saturation prolongée du réseau ; STS fournit une synchronisation incrémentielle planifiée et vérifiée par checksum. Le stockage dans Cloud Storage avec le connecteur GCS permet le traitement par Dataproc sans nécessiter 50 To de Persistent Disk par nœud.
- Migrer les bases de données opérationnelles avec le CDC
- Utiliser DMS pour migrer MySQL et PostgreSQL avec un temps d’arrêt minimal. Pour le CDC d’Oracle vers l’analytique, utiliser Datastream pour le landing dans Cloud Storage, puis un modèle Dataflow fourni par Google pour charger dans BigQuery. Justification : DMS exploite la réplication native pour un snapshot fiable + une synchronisation continue ; Datastream fournit un CDC serverless avec des métadonnées de commit, tandis que le modèle Dataflow assure des écritures ordonnées et idempotentes dans BigQuery.
- Ingestion des flux SaaS et basés sur des fichiers
- Construire des pipelines Cloud Data Fusion utilisant des connecteurs SaaS pour les événements CRM avec des jetons incrémentiels, et un pipeline de fichiers pour ingérer les CSV quotidiens depuis un SFTP fournisseur via STS. Normaliser en Avro dans un bucket Cloud Storage organisé (curated), puis charger des tables BigQuery partitionnées. Justification : Cloud Data Fusion centralise les connecteurs, la transformation et le lignage (lineage). La standardisation sur Avro préserve le schéma et facilite son évolution ; les tables BigQuery partitionnées réduisent le coût des requêtes.
- Diffuser des événements en temps réel
- Publier les événements web et en magasin sur Pub/Sub. Traiter avec Dataflow pour l’analyse syntaxique, la validation, l’enrichissement et le watermarking ; écrire dans BigQuery via la Storage Write API et archiver les données brutes Avro dans Cloud Storage. Justification : Pub/Sub découple les producteurs et les consommateurs ; Dataflow fournit l’autoscaling, le traitement avec état (stateful), les checkpoints et la gestion des données en retard ; la double écriture assure à la fois des analyses à faible latence et une rétention durable des données brutes.
- Appliquer des contrôles de qualité des données et de schéma en périphérie
- Implémenter un registre de schémas et une validation dans Dataflow/Data Fusion. Acheminer les enregistrements mal formés vers un bucket de quarantaine GCS et un sujet de lettres mortes (dead-letter) Pub/Sub. Appliquer des vérifications de domaine (ex: codes de devise, horodatages UTC) et dédoublonner à l’aide de clés composites. Justification : Le rejet précoce et la mise en quarantaine empêchent la propagation de données incorrectes ; l’idempotence et le dédoublonnage protègent contre la livraison “au moins une fois” des sources CDC et de streaming.
- Optimiser le stockage et l’accès pour l’analytique
- Charger les jeux de données organisés (curated) dans des tables BigQuery partitionnées et clusterisées. Exposer les archives brutes en tant que tables externes permanentes pour une exploration peu fréquente. Pour les charges de travail OLTP qui restent transactionnelles, conserver Cloud SQL avec des réplicas en lecture. Justification : Le partitionnement et la clusterisation minimisent le coût de balayage (scan) ; les tables externes évitent des chargements inutiles pour un accès occasionnel ; Cloud SQL préserve les sémantiques ACID pour les applications transactionnelles.
- Planifier le basculement, le rattrapage et le rollback
- Exécuter un snapshot + CDC pour chaque RDBMS ; atteindre un point de réconciliation où les nombres de lignes et les checksums correspondent. Exécuter en mode blue/green avec des doubles écritures pendant 48 heures, en déplaçant progressivement les lectures vers BigQuery. Maintenir un feature flag pour annuler les écritures si des divergences sont détectées. Justification : Le mode blue/green réduit le risque ; la vérification à un watermark connu garantit l’exhaustivité ; les flags permettent un rollback rapide.
- Vérification et observabilité
- Créer des tables de contrôle capturant les LSN/SCN source, les horodatages de commit, les nombres de lignes et les checksums par partition. Surveiller le retard de Datastream, l’état de réplication de DMS, les watermarks de Dataflow, le backlog de Pub/Sub, l’état des tâches STS et les métriques d’insertion en streaming de BigQuery. Justification : Un lignage de bout en bout et des contrôles quantitatifs fournissent une preuve de conformité auditable et des alertes rapides sur les manques ou les retards.
En séparant les couches de landing, de curation et de service (serving) ; en utilisant Cloud Storage comme zone de transit (staging) et d’archivage durable et à faible coût ; en tirant parti de DMS/Datastream pour le CDC avec des consommateurs idempotents ; et en appliquant le schéma et la qualité à l’entrée (ingress), Northstar Retail réalise une ingestion sécurisée et évol
← Spark · Tous les domaines · Orchestration de workflows et automatisation de pipelines →
Entraînez-vous sur ces questions → · Tests chronométrés sur ExamRoll.io →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
Réussissez votre examen →