Google PDE: Spark, Dataproc et traitement de données distribué — 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
Apache Spark sur Google Cloud Dataproc fournit une plateforme élastique et gérée pour le traitement distribué des données. Vous pouvez choisir entre des clusters Dataproc éphémères ou de longue durée et Dataproc Serverless pour Spark, en fonction des besoins de contrôle, de la variabilité de l’exécution et de la charge de gestion. Spark offre des abstractions résilientes (RDDs), des API relationnelles (DataFrames et Spark SQL), et un moteur d’exécution DAG tolérant aux pannes, optimisé pour l’ETL itératif et par lots à grande échelle. Sur Google Cloud, Cloud Storage remplace HDFS pour un stockage durable et à faible coût ; le connecteur BigQuery permet un déchargement analytique direct ; et Dataproc Metastore centralise la gestion des schémas. Les solutions efficaces alignent les cycles de vie du stockage et du calcul, ajustent Spark à la charge de travail, instrumentent l’observabilité et appliquent la sécurité avec le moindre privilège et l’isolation réseau.
Architecture Dataproc : Clusters, Serverless, Stockage et Metastore
- Types de clusters et rôles des nœuds
- Les nœuds primaires (maîtres) hébergent YARN, le HDFS NameNode (si utilisé), les interfaces utilisateur du pilote Spark ; le mode HA (haute disponibilité) utilise plusieurs nœuds primaires.
- Les nœuds de travail (workers) exécutent les exécuteurs et les HDFS DataNodes (si utilisés).
- Les workers secondaires/auxiliaires sont généralement préemptifs/spot pour une capacité élastique à moindre coût, sans les rôles HDFS.
- Les images regroupent les versions de l’OS et des composants (par exemple, 2.1-debian11, 2.2-ubuntu20) ; épinglez les versions d’image pour contrôler la compatibilité Spark/Hadoop et effectuer les mises à niveau de manière délibérée.
- Component Gateway publie les interfaces utilisateur (Spark History Server, YARN RM) de manière sécurisée via HTTPS.
- Dataproc Serverless pour Spark
- Pas de provisionnement de cluster, autoscaling automatique et facturation à la seconde pour les exécuteurs et les pilotes. Idéal pour les tâches sporadiques ou en rafale, ou pour minimiser la charge opérationnelle.
- Inconvénients : moins de réglages de bas niveau que les clusters ; la latence de démarrage des tâches peut être plus élevée que sur des clusters préchauffés ; utilisez les métriques et les journaux d’événements serverless pour le dépannage.
- Autoscaling
- Les politiques d’autoscaling de cluster ajoutent/suppriment des workers en fonction des métriques YARN/Spark et des délais de récupération (cooldowns), en ajustant séparément les groupes de workers primaires et secondaires.
- L’autoscaling Serverless est géré par le service ; concevez des traitements parallélisables par partition et évitez les goulots d’étranglement sérialisés pour une mise à l’échelle optimale.
- Stockage et connecteurs
- Préférez Google Cloud Storage (GCS) comme système de référence (system-of-record) ; il découple le calcul du stockage, réduit le coût des disques persistants et survit aux cycles de vie des clusters.
- Le connecteur GCS (gs://) s’intègre avec Hadoop/Spark. Les écritures dans les stockages d’objets utilisent des protocoles de commit ; configurez l’algorithme FileOutputCommitter v2 pour réduire la surcharge liée au renommage et accélérer les commits de tâches sur GCS :
--conf mapreduce.fileoutputcommitter.algorithm.version=2
```
- Utilisez Parquet/ORC avec l'élagage de colonnes (column pruning) et la délégation de prédicats (predicate pushdown). Gérez les petits fichiers par compaction pour viser 128–512 Mio par fichier pour une lecture efficace.
- Metastore Hive
- Centralisez les schémas et les métadonnées des tables dans Dataproc Metastore (un Metastore Apache Hive géré) ou un metastore adossé à Cloud SQL pour partager les catalogues entre les clusters.
- Utilisez des tables externes pointant vers GCS pour la durabilité ; partitionnez par date/heure pour limiter le coût des analyses (scans).
- Tâches, initialisation et workflows
- Soumettez des tâches spark, pyspark, spark-sql ou hadoop. Les actions d'initialisation installent des bibliothèques ou des agents supplémentaires lors de la création du cluster (par exemple, des connecteurs, des bibliothèques Python).
- Les modèles de workflow (Workflow templates) paramètrent des pipelines à plusieurs étapes ; ils peuvent créer des clusters éphémères par workflow, puis les détruire. Cela améliore l'isolation et réduit les coûts d'inactivité.
- Les clusters éphémères sont recommandés pour l'ETL par lots ; les données et le metastore résident en dehors du cluster (GCS, Dataproc Metastore, BigQuery).
- Intégration avec BigQuery
- Le connecteur Spark BigQuery lit/écrit directement dans BigQuery ; envisagez la BigQuery Storage Read API pour le débit et la Write API pour des insertions en streaming à plus faible latence et avec sémantique exactly-once.
- Pour la maintenance des tables, effectuez des opérations MERGE ou des remplacements de partitions en aval dans BigQuery pour finaliser les chargements de manière atomique.
### Modèle Spark, réglage des performances et fiabilité
- API et exécution
- RDDs : bas niveau, immuables, typés en Scala/Java ; vous contrôlez le partitionnement et la persistance.
- DataFrames/Datasets : relationnels, optimisés par Catalyst ; préférez-les pour l'ETL en raison de l'optimisation des requêtes et de la génération de code.
- Les transformations sont paresseuses (map, filter, join) ; les actions déclenchent l'exécution (count, collect, save). Spark construit un DAG d'étapes (stages) séparées par des shuffles ; les tâches s'exécutent par partition.
- Partitionnement et shuffle
- Partitionnement en entrée : suffisamment de partitions pour utiliser tous les cœurs ; commencez avec 2 à 4 fois le nombre total de cœurs d'exécuteur. Contrôlez via `spark.default.parallelism` (pour les RDDs) et les options du lecteur (pour les DataFrames).
- Partitions de shuffle : la valeur par défaut de 200 est souvent sous ou sur-provisionnée. Réglez :
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
- Ciblez environ 100–256 Mio par partition après les transformations larges (wide transforms) ; une taille trop petite entraîne une surcharge du planificateur (scheduler) et une taille trop grande risque un OOM de l’exécuteur.
- Le shuffle est le coût dominant pour les jointures (joins), les
groupByet lesorderBy. Assurez-vous que la mémoire et le disque de l’exécuteur sont adéquats ; envisagez des SSD locaux pour les shuffles intensifs sur les clusters. - Asymétrie (skew) et stratégie de jointure
- Détectez l’asymétrie (temps d’exécution des tâches à longue traîne, grandes tailles de partition). Atténuations :
- Diffusez (broadcast) les petites tables pour éviter les shuffles :
- Détectez l’asymétrie (temps d’exécution des tâches à longue traîne, grandes tailles de partition). Atténuations :
--conf spark.sql.autoBroadcastJoinThreshold=64m
```
- Salez (salt) les clés pour les partitions surchargées (hot partitions) ; appliquez une pré-agrégation côté map ; filtrez tôt.
- Activez l'Adaptive Query Execution (AQE) pour fusionner les partitions post-shuffle et gérer les jointures asymétriques (skewed joins) :
--conf spark.sql.adaptive.enabled=true
```
- Mise en cache, checkpointing et lignage
- Mettez en cache avec parcimonie les DataFrames intermédiaires fréquemment réutilisés (hot) ; préférez
MEMORY_AND_DISKpour éviter les OOM. - Faites des checkpoints (checkpointing) des lignages longs vers GCS ou HDFS pour limiter le recalcul en cas d’échec.
- Mettez en cache avec parcimonie les DataFrames intermédiaires fréquemment réutilisés (hot) ; préférez
- Exécuteurs et allocation dynamique
- Dimensionnez correctement les exécuteurs pour équilibrer le parallélisme et la surcharge du GC :
- Cœurs par exécuteur : 2 à 5 pour des tâches équilibrées en I/O/CPU ; moins de cœurs réduit les pauses du GC.
- Surcharge mémoire (memory overhead) : définissez
spark.yarn.executor.memoryOverheadpour les shuffles larges. - Activez l’allocation dynamique avec un service de shuffle externe sur les clusters pour adapter les exécuteurs à la charge de travail :
- Dimensionnez correctement les exécuteurs pour équilibrer le parallélisme et la surcharge du GC :
--conf spark.dynamicAllocation.enabled=true
--conf spark.shuffle.service.enabled=true
--conf spark.dynamicAllocation.minExecutors=0
--conf spark.dynamicAllocation.maxExecutors=200
```
- Patrons de tolérance aux pannes pour l'ETL par lots
- Écritures idempotentes : écrivez dans un chemin temporaire/de staging, puis promouvez atomiquement avec un commit au niveau du répertoire ; pour BigQuery, écrivez dans une table de staging et utilisez `MERGE` :
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```
- Traitement incrémental : utilisez le filtrage basé sur les watermarks sur les partitions
ingestion_date; maintenez un manifeste des fichiers traités dans GCS pour éviter le retraitement. - Gestion des lettres mortes (dead-letter) : en cas d’erreurs d’analyse/validation, dérivez les enregistrements incorrects vers un chemin/une table de quarantaine avec des diagnostics. Pour une application stricte du schéma et des DLQ intégrées, envisagez Dataflow ; avec Spark, implémentez un
try/catchpar enregistrement et un récepteur (sink) séparé.
Sécurité, Observabilité et Coût
- Identité et accès
- Exécutez les clusters et les jobs avec des comptes de service dédiés en appliquant le principe de moindre privilège IAM. N’accordez que les rôles requis, par exemple :
- roles/dataproc.worker aux comptes de service des instances
- roles/storage.objectViewer ou objectAdmin pour les chemins d’E/S GCS
- roles/bigquery.dataEditor sur les ensembles de données cibles
- Pour Dataproc Serverless, utilisez des comptes de service par job pour limiter la portée des accès.
- Exécutez les clusters et les jobs avec des comptes de service dédiés en appliquant le principe de moindre privilège IAM. N’accordez que les rôles requis, par exemple :
- Isolation réseau et chiffrement
- Utilisez des clusters à IP privée dans un sous-réseau VPC, restreignez les interfaces utilisateur du maître par pare-feu et activez l’Accès privé à Google pour GCS/BigQuery sans sortie publique.
- Placez les clusters dans des projets VPC partagé pour un contrôle centralisé. Activez optionnellement Kerberos sur Dataproc pour l’authentification au sein du cluster.
- Chiffrez au repos avec CMEK : configurez CMEK sur les buckets GCS, les disques persistants, Dataproc Metastore et BigQuery ; utilisez TLS en transit par défaut.
- Journalisation, historique et métriques
- Activez les journaux d’événements Spark vers GCS et déployez le History Server :
--conf spark.eventLog.enabled=true
--conf spark.eventLog.dir=gs://bucket/spark-events/
```
- Dataproc envoie les journaux du pilote et de YARN vers Cloud Logging ; exportez-les vers des récepteurs pour la rétention/l'analyse forensique.
- Surveillez avec les métriques Cloud Monitoring : conteneurs YARN en attente, CPU, mémoire, santé HDFS (si utilisé), débit GCS. Déclenchez des alertes en cas de nouvelles tentatives prolongées d'étapes, de perte d'exécuteur et de pics d'exécution spéculative.
- Analyse des défaillances : les causes courantes incluent les traînards induits par le déséquilibre (skew), les OOM (Out of Memory) de l'exécuteur pendant le shuffle, les échecs de commit du stockage objet et la perte de nœuds préemptifs/spot. Augmentez le nombre de tentatives judicieusement ; des tentatives excessives peuvent amplifier le coût et les délais.
- Optimisation des coûts
- Utilisez des clusters éphémères ou Dataproc Serverless pour éviter les coûts d'inactivité ; conservez les données dans GCS pour minimiser l'utilisation de disques persistants.
- Ajoutez des nœuds de calcul secondaires préemptifs/spot pour absorber les pics de demande ; concevez pour le recalcul, car les tâches sur les nœuds perdus sont relancées. Ne placez pas les maîtres sur des nœuds préemptifs.
- Dimensionnez correctement les types de machine et utilisez l'autoscaling pour réduire la capacité lorsque les files d'attente sont vides. Privilégiez Parquet/ORC avec l'élagage de partitions pour réduire le coût de scan et l'utilisation du CPU.
- Évitez les petits fichiers en compactant les sorties ; moins de fichiers, mais plus volumineux, réduisent la surcharge de métadonnées et le temps d'exécution des jobs.
- Pour les jobs courts et périodiques (par exemple, un ETL Spark hebdomadaire de 30 minutes), les nœuds de calcul préemptifs ou le mode serverless offrent souvent le meilleur profil de coût.
#### Scénario de problème pratique
Acme Retail migre un cluster Hadoop sur site (on-prem) de 30 nœuds exécutant des ETL Spark et Hive nocturnes qui alimentent des analyses en aval. Ils souhaitent réutiliser les jobs existants avec des modifications minimales, éviter de gérer des clusters à plein temps, persister les données au-delà de la durée de vie des clusters et réduire les coûts de stockage.
Approche :
1) Placer les données et métadonnées dans des services gérés
- Stocker toutes les données brutes et traitées dans Cloud Storage en utilisant Parquet avec partitionnement (par exemple, dt=YYYY-MM-DD).
- Justification : GCS est durable, peu coûteux et découple le calcul du stockage, permettant ainsi aux clusters éphémères et aux jobs serverless de s'exécuter sans disques persistants. Le format Parquet partitionné permet la délégation de prédicat (predicate pushdown) et des scans efficaces.
2) Centraliser le catalogue avec Dataproc Metastore
- Migrer le métastore Hive vers Dataproc Metastore. Créer des tables Hive externes référençant les chemins GCS et conserver la logique de schéma/partitionnement existante.
- Justification : Un métastore géré permet à plusieurs clusters éphémères et jobs serverless de partager les définitions de tables sans avoir à exécuter une instance MySQL/PostgreSQL en haute disponibilité.
3) Utiliser des clusters Dataproc éphémères pour l'ETL par lots et des modèles de workflow pour l'orchestration
- Définir un modèle de workflow qui crée un cluster avec l'image requise (par exemple, 2.1-debian11), exécute des jobs Spark (spark-sql et pyspark) et supprime le cluster à la fin. Ajouter des actions d'initialisation pour installer toute bibliothèque personnalisée.
- Justification : Les clusters éphémères éliminent les coûts d'inactivité et isolent les dépendances des jobs. Les modèles de workflow assurent la répétabilité et la paramétrisation (dates, chemins d'entrée).
4) Activer l'autoscaling et les nœuds de calcul préemptifs
- Attacher une politique d'autoscaling avec un petit groupe de nœuds de calcul principaux (core) et un pool plus large de nœuds de calcul secondaires préemptifs ; ajuster les délais de récupération (cooldowns) pour réduire la taille rapidement après l'exécution.
- Justification : Les nœuds de calcul principaux maintiennent la stabilité du cluster ; les nœuds de calcul préemptifs absorbent les shuffles et les transformations larges à moindre coût. Les nouvelles tentatives de Spark/YARN gèrent les tâches perdues lors de la préemption.
5) Intégrer avec BigQuery via le connecteur Spark BigQuery
- Pour les chargements de dimensions/faits, écrire les résultats Spark dans des tables de pré-production (staging) BigQuery, puis exécuter des instructions MERGE pour mettre à jour les cibles de manière atomique. Lorsque l'écrasement direct est sûr, écrire dans des tables partitionnées en utilisant le mode de remplacement de partition (partition overwrite).
- Justification : BigQuery sert les analyses et la BI à grande échelle ; la combinaison pré-production+MERGE produit des opérations de type upsert transactionnel à partir de Spark en batch, réduisant l'incohérence en aval.
6) Optimiser Spark pour la performance et la fiabilité
- Définir les partitions de shuffle par rapport aux cœurs des exécuteurs et activer AQE :
--conf spark.sql.shuffle.partitions=600
--conf spark.sql.adaptive.enabled=true
```
- Utiliser des jointures de diffusion (broadcast joins) pour les petites dimensions et créer des points de contrôle (checkpoint) des longs lignages vers GCS pour la stabilité.
- Justification : Un partitionnement approprié réduit le déséquilibre (skew) et la surcharge du planificateur ; AQE s’adapte aux profils de données à l’exécution ; les points de contrôle limitent le recalcul après une défaillance.
Renforcer la sécurité et le réseau
- Exécuter les clusters avec des comptes de service dédiés n’accordant que les rôles nécessaires pour les chemins GCS, le métastore et les ensembles de données BigQuery. Créer des clusters à IP privée dans un sous-réseau restreint avec l’Accès privé à Google et limiter l’accès à l’interface utilisateur via des règles de pare-feu.
- Justification : Le moindre privilège et l’isolation réseau réduisent la surface d’attaque ; une sortie privée pour le plan de contrôle évite l’exposition publique.
Instrumenter la journalisation, l’historique et les alertes
- Activer les journaux d’événements Spark vers GCS et déployer le History Server ; acheminer les journaux du pilote/YARN vers Cloud Logging avec rétention. Ajouter des alertes Monitoring pour les conteneurs en attente prolongée, les échecs de tâches répétés ou une durée de job excessive.
- Justification : Les journaux centralisés facilitent l’analyse des causes profondes ; les alertes proactives détectent tôt les déséquilibres (skew), les OOM ou les E/S dégradées.
Moderniser sélectivement avec Dataproc Serverless pour les charges de travail ad hoc et les pics élastiques
- Déplacer les charges de travail Spark SQL sporadiques ou exploratoires vers Dataproc Serverless ; conserver les pipelines nocturnes sur des clusters éphémères jusqu’à leur validation complète en mode serverless.
- Justification : Le mode serverless supprime les opérations de cluster et s’adapte automatiquement, ce qui est idéal pour les charges imprévisibles ; les workflows existants continuent avec un minimum de changement de code.
Valider les committers de stockage objet et la gestion des petits fichiers
- Définir l’algorithme FileOutputCommitter v2 et compacter les sorties à 256–512 Mio par fichier via repartition/coalesce avant les écritures.
- Justification : Les stockages objet n’ont pas de renommage atomique ; les committers optimisés réduisent la surcharge de copie/renommage. Le compactage atténue le problème des petits fichiers pour la performance et le coût.
Cette conception réutilise les jobs Spark et Hive existants avec un remaniement minimal, assure la durabilité des données dans GCS, centralise les schémas, contient le rayon d’impact en matière de sécurité, fournit une observabilité robuste et optimise les coûts grâce à des clusters éphémères, l’autoscaling, la capacité préemptive et une utilisation ciblée de l’exécution serverless.
← Messagerie · Tous les domaines · Ingestion →
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 →