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

    --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}
```
      --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
  ```
      --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 (...)
```

Sécurité, Observabilité et Coût

    --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
 ```
  1. 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.
  2. 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.
  3. 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.
  4. 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 →

Parcourir Google →

Related guides

Accès tout-en-un

Un seul abonnement. Chaque examen.

Chaque plan débloque la recherche de réponses illimitée, les tests pratiques, les explications IA et la bibliothèque complète de ressources — en plus de 20 langues.

Mensuel
24.87
Just €0.83/day
Tout inclus :
  • Recherche de réponses illimitée
  • Tests pratiques illimités
  • Explications basées sur l'IA
  • Bibliothèque complète de ressources
  • Plus de 20 langues
  • Mises à jour de contenu hebdomadaires
  • Récompenses et parrainages
  • Support prioritaire
Commencer l'essai gratuit

Aucune carte de crédit requise*

Meilleur rapport qualité/prix
12 mois
179.87
Just €0.49/daySave 40%
Tout inclus :
  • Recherche de réponses illimitée
  • Tests pratiques illimités
  • Explications basées sur l'IA
  • Bibliothèque complète de ressources
  • Plus de 20 langues
  • Mises à jour de contenu hebdomadaires
  • Récompenses et parrainages
  • Support prioritaire
Commencer l'essai gratuit

Aucune carte de crédit requise*

✓ Plan gratuit inclus · ✓ Annulez à tout moment · ✓ Tous les plans débloquent le produit complet