Apache Spark est un moteur de calcul distribué pour le traitement et l'analyse de données à grande échelle. Sa puissance repose sur une séparation claire des responsabilités entre le driver, les executors et un gestionnaire de cluster, ainsi que sur un modèle d'exécution résilient qui transforme le code utilisateur en étapes de tâches parallèles. Cet article explique comment les composants de Spark s'articulent, comment les flux de contrôle et de données circulent dans un job, et comment exécuter un pilote pratique et à faible risque que vous pouvez vérifier de bout en bout.
Ce Que Vous Allez Apprendre
- Les principaux composants de Spark et leurs responsabilités opérationnelles
- Comment le flux de contrôle (planification et ordonnancement) et le flux de données (shuffles, broadcasts) fonctionnent en pratique
- Un chemin de configuration sûr pour exécuter un job petit mais réaliste
- Les étapes de vérification à l'aide de l'interface utilisateur Spark et de contrôles de données simples
- Les modes de défaillance courants et les actions de récupération précises
- Une liste de contrôle opérationnelle pratique que vous pouvez appliquer à chaque job
Ce guide s'adresse aux développeurs, consultants DevOps et équipes techniques de startups qui doivent déployer des jobs Spark de manière fiable, raisonner sur la performance et récupérer rapidement en cas de problème.
Composants Spark en Un Coup d'Œil
Le tableau suivant résume les composants d'exécution principaux et ce qu'il faut surveiller opérationnellement.
| Composant | Rôle | Préoccupations Opérationnelles |
|---|---|---|
| Driver | Orchestre l'application : crée la SparkSession, construit les plans logiques/physiques, planifie les étapes et les tâches, suit les métadonnées | Éviter les OOM du driver via un travail côté executor, limiter collect() sur les gros volumes, s'assurer que la mémoire et le CPU du driver sont suffisants |
| Executors | Exécutent les tâches, mettent en cache les données, effectuent les shuffles et les écritures, renvoient l'état au driver | Définir cœurs et mémoire par executor de manière appropriée, éviter les OOM d'executor, surveiller le GC et les spills sur disque |
| Gestionnaire de Cluster | Alloue les ressources à l'application Spark (Standalone, YARN, Kubernetes) | Capacité des files d'attente, limites d'application, préemption, santé des nœuds |
| Stockage Externe | Systèmes d'entrée et de sortie (HDFS, stockage objet, puits JDBC) | Débit, cohérence, protocole de commit de sortie, partitionnement et évolution de schéma |
Comment les Flux de Contrôle et de Données Circulent dans Spark
Flux de contrôle : Votre code (DataFrame, Dataset, SQL) définit un plan logique. L'optimiseur Catalyst de Spark construit un plan physique et le divise en étapes séparées par des frontières de shuffle. Le driver soumet les tâches aux executors via le gestionnaire de cluster. Les executors exécutent les tâches, rapportent la progression et produisent des fichiers intermédiaires pour les shuffles.
Flux de données : Les dépendances étroites (map, filter, withColumn) s'enchaînent au sein d'une partition ; les dépendances larges (groupBy, join, reduceByKey, repartition) déclenchent des shuffles où les données sont repartitionnées et échangées entre executors. Les variables broadcast envoient de petites données de référence à tous les executors. Les caches persistent les données en mémoire ou sur disque pour réutilisation.
Inventaire des Versions et de l'Environnement
Établir les versions et une topologie minimale mais réaliste évite les incompatibilités subtiles.
Prérequis (exemple construit pour un petit pilote) :
- OS : Linux x86_64, 4 à 8 vCPU par worker, 32 à 64 Go de RAM par nœud worker
- Java : OpenJDK 11 ou 17
- Python/Scala : Python 3.9+ ou Scala 2.12.x selon le choix d'API
- Spark : 3.3.x, 3.4.x ou 3.5.x
- Stockage : HDFS 3.2+ ou stockage objet compatible Hadoop (S3, GCS, ADLS), plus un petit puits de sortie (chemin HDFS ou table de test JDBC)
- Réseau : Les workers peuvent atteindre le stockage et le driver ; ports de l'UI Spark ouverts (4040+ pour l'application active, History Server souvent 18080) selon les besoins
Topologie de référence pour le pilote (exemple construit) :
- 1 hôte driver (8 vCPU, 16 Go RAM)
- 3 hôtes workers (chacun 8 vCPU, 64 Go RAM, espace SSD scratch pour le shuffle)
- Gestionnaire de cluster : Spark Standalone, YARN dans une file de développement, ou Kubernetes
- Stockage : HDFS avec facteur de réplication 2 ou 3 ; un bucket S3-compatible de développement convient s'il est disponible
Données et charge de travail pour le pilote :
- Taille d'entrée : 50 Go à 150 Go de données semi-structurées (ex. CSV ou Parquet)
- Transformations : 1 ou 2 opérations larges déclenchant un shuffle (ex. groupBy ou join) pour exercer les chemins réseau et disque
- Sortie : Parquet partitionné ou petite table JDBC pour une recherche dimensionnelle, selon vos contraintes
Notes sur les fonctionnalités par version :
- L'exécution de requêtes adaptative (AQE - Adaptive Query Execution) est disponible dans Spark 3.x et peut améliorer la gestion du décalage et la sélection des jointures. Commencez avec elle activée pour les charges de travail SQL/DataFrame et vérifiez avec l'UI.
- Les nouvelles implémentations de shuffle et l'allocation dynamique fonctionnent bien pour des pilotes modérés ; gardez-les conservatrices jusqu'à ce que vous mesuriez.
Chemin de Configuration Sûr
Votre objectif est d'exécuter un job facile à vérifier et difficile à casser. Les paramètres de base suivants supposent la topologie de référence ci-dessus. Ajustez proportionnellement si votre matériel est plus petit ou plus grand.
Dimensionnement de Base des Executors (Exemple Construit)
- Executors par worker : 3
- Cœurs par executor : 2 à 3
- Mémoire par executor : 12 à 16 Go
- Surcharge mémoire par executor : 2 à 3 Go
- Parallélisme total visé : environ 2x à 3x le total des cœurs entre executors (ex. 48 à 72 tâches si vous avez 24 cœurs d'executor au total)
Fonctionnalités d'Exécution de Base
spark.sql.adaptive.enabled=truepour les jobs DataFrame/SQLspark.dynamicAllocation.enabled=falsepour le premier pilote (activez plus tard une fois mesuré)spark.sql.shuffle.partitionsdéfini à 2x à 3x le total des cœurs d'executor (exemple construit : 64)spark.sql.autoBroadcastJoinThresholddéfini à une taille modeste (exemple construit : 64 Mo) pour permettre les jointures de hachage broadcast quand utilespark.speculation=falseinitialement ; activez plus tard si les tâches montrent des traînards persistants après réglage
Résumé de la Configuration
| Domaine | Paramètre | Exemple Construit |
|---|---|---|
| Executors | spark.executor.instances | 9 (3 par worker) |
| Executors | spark.executor.cores | 2 |
| Executors | spark.executor.memory | 14g |
| Executors | spark.executor.memoryOverhead | 3g |
| Parallélisme | spark.default.parallelism | 48 |
| Shuffles | spark.sql.shuffle.partitions | 64 |
| AQE | spark.sql.adaptive.enabled | true |
| Jointures | spark.sql.autoBroadcastJoinThreshold | 64m |
| Spéculation | spark.speculation | false |
| Alloc. dynamique | spark.dynamicAllocation.enabled | false |
Exemple Construit : spark-submit pour un Job DataFrame
spark-submit \
--class com.example.PilotJob \
--master yarn \
--deploy-mode cluster \
--conf spark.executor.instances=9 \
--conf spark.executor.cores=2 \
--conf spark.executor.memory=14g \
--conf spark.executor.memoryOverhead=3g \
--conf spark.default.parallelism=48 \
--conf spark.sql.shuffle.partitions=64 \
--conf spark.sql.adaptive.enabled=true \
--conf spark.sql.autoBroadcastJoinThreshold=64m \
--conf spark.speculation=false \
pilot-job-assembly.jar \
--input hdfs:///data/pilot/input \
--output hdfs:///data/pilot/output
Si vous préférez PySpark pour le pilote, définissez les mêmes clés de configuration et passez votre script :
spark-submit \
--master yarn \
--deploy-mode cluster \
--conf spark.executor.instances=9 \
--conf spark.executor.cores=2 \
--conf spark.executor.memory=14g \
--conf spark.executor.memoryOverhead=3g \
--conf spark.sql.shuffle.partitions=64 \
--conf spark.sql.adaptive.enabled=true \
pilot_job.py \
--input hdfs:///data/pilot/input \
--output hdfs:///data/pilot/output
Exemple Construit de Logique de Job (PySpark) Illustrant les Opérations Étroites et Larges
from pyspark.sql import SparkSession, functions as F
spark = (SparkSession.builder
.appName("pilot-job")
.getOrCreate())
src = spark.read.parquet("hdfs:///data/pilot/input")
ref = spark.read.parquet("hdfs:///data/pilot/ref_dim")
# Transformations étroites
clean = (src
.filter(F.col("status") == "active")
.withColumn("event_dt", F.to_date("event_ts")))
# Opération large : join (peut shuffler)
joined = clean.join(ref.hint("broadcast"), on="key", how="left")
# Opération large : agrégation (shuffle)
agg = (joined
.groupBy("event_dt")
.agg(F.count("*").alias("cnt")))
agg.write.mode("overwrite").partitionBy("event_dt").parquet("hdfs:///data/pilot/output")
Directive de Dimensionnement des Partitions (Exemple Construit)
- Si votre entrée fait 100 Go et que votre taille de partition cible est d'environ 128 Mo, attendez-vous à environ 800 partitions pour une lecture initiale. Après les opérations larges, utilisez
spark.sql.shuffle.partitionspour contrôler le parallélisme en aval (ex. 64 à 128). Profilez le temps d'exécution et le décalage avant de changer.
Vérification et Diagnostics
La vérification confirme que l'architecture se comporte comme prévu et que les résultats sont corrects. Effectuez ces contrôles dans l'ordre.
1. Confirmer que le Plan de Contrôle est Sain
UI Spark : Ouvrez l'UI de l'application (par défaut 4040 pour l'application active ; en mode cluster utilisez le lien dans l'UI du gestionnaire de ressources). Vérifiez :
- Onglet Jobs : Le job pilote doit afficher 1 à 3 jobs se terminant avec un petit nombre d'étapes.
- Onglet Stages : Les étapes pour les jointures et agrégations doivent avoir plusieurs tâches égales à vos partitions de shuffle.
- Onglet Executors : 9 executors actifs (exemple construit), avec une distribution de tâches relativement uniforme.
- Logs : Inspectez les logs du driver et des executors pour les WARN/ERROR. Attendez-vous à peu de WARN sur les tâches spéculatives (désactivées) et d'éventuels INFO sur les décisions AQE.
2. Valider le Chemin de Données avec des Contrôles de Cardinalité Simples
# Dans le shell PySpark ou un notebook connecté au driver
print(src.count()) # Lignes d'entrée
print(agg.count()) # Groupes de sortie (ex. nombre de dates)
- Vérification ponctuelle des clés : vérifiez que quelques clés de la table de référence apparaissent dans la sortie.
- Schéma : assurez-vous que les colonnes attendues existent et que le partitionnement est appliqué comme prévu.
3. Inspecter le Plan Physique
Utilisez explain dans DataFrame/SQL pour voir où les shuffles se produisent :
print(agg.explain(True)) # Cherchez BroadcastHashJoin ou SortMergeJoin et les nœuds Exchange
Résultats attendus :
- Si la table de référence est petite (moins que le seuil de broadcast), Spark doit choisir BroadcastHashJoin. L'AQE peut changer les types de jointures ou réduire les partitions de shuffle à l'exécution ; vous devriez voir des notes dans les infobulles de l'UI ou les détails d'étape.
4. Capacité et Utilisation des Ressources
- Onglet Executors : Le temps GC doit être une petite fraction du temps de tâche (objectif construit : < 10 %). Si le GC est élevé, réduisez la mémoire de l'executor ou les cœurs par executor pour ajuster la pression sur le tas et le parallélisme.
- Tâches par executor : relativement équilibrées. De gros déséquilibres suggèrent un décalage.
5. Exactitude de la Sortie et Idempotence
Relancez le job avec les mêmes entrées et le mode overwrite. La sortie doit être structurellement identique (mêmes partitions et comptages de lignes). Sinon, vérifiez les UDF non déterministes ou les timestamps d'ingestion qui s'infiltrent dans le chemin d'écriture.
Modes de Défaillance et Récupération
Même un pilote sûr rencontre des obstacles réels. Utilisez cette section pour identifier rapidement les symptômes et appliquer des correctifs précis.
1. Driver Out-of-Memory (OOM)
Symptômes : Le log du driver montre OutOfMemoryError ; le job échoue près d'actions comme collect(), toPandas(), ou création de broadcast volumineux.
Cause racine : Récupération de trop de données vers le driver ou agrégations massives côté driver.
Récupération :
- Supprimez
collect()sur les gros jeux de données ; utilisezshow()avec limites ou écrivez vers le stockage pour inspection. - Augmentez modestement la mémoire du driver si vraiment nécessaire (ex.
--driver-memory 8gà12g) et privilégiez le travail côté executor. - Si un broadcast est trop gros, abaissez
spark.sql.autoBroadcastJoinThresholdet utilisez une jointure par shuffle.
2. Executor OOM ou GC Excessif
Symptômes : Les executors meurent en milieu d'étape ; longs temps de GC ; OOM dans les logs.
Cause racine : Partitions trop grosses, UDF gourmands en mémoire, gros shuffles.
Récupération :
- Réduisez
spark.sql.shuffle.partitionspour éviter trop de reducers concurrents par executor, ou augmentez-le si les tâches individuelles sont trop lourdes. Ajustez selon le décalage. - Réduisez les cœurs par executor (ex. de 3 à 2) pour diminuer la pression mémoire concurrente.
- Augmentez
spark.executor.memoryOverheadsi les tampons de spill sont serrés. - Privilégiez les fonctions intégrées aux UDF Python ; envisagez
mapInPandasseulement quand nécessaire et sûr en mémoire.
3. Échecs de Récupération de Shuffle
Symptômes : FetchFailedException ; les étapes retentent répétitivement ; fichiers de spill de sortie manquants.
Cause racine : Executor ou nœud perdu pendant le shuffle, pression disque, incidents réseau.
Récupération :
- Assurez un disque local adéquat pour le shuffle et surveillez la saturation I/O.
- Relancez après avoir vérifié la santé des nœuds. Si fréquent, réduisez les partitions de shuffle ou ajustez la stabilité du cluster.
- Envisagez d'activer le service de shuffle externe sur les gestionnaires supportés et maintenez
spark.shuffle.service.enabledcohérent avec le mode de déploiement.
4. Décalage de Données et Traînards
Symptômes : Quelques tâches prennent beaucoup plus de temps ; blocage en tête de file d'étape ; comptages de tâches inégaux par executor.
Cause racine : Clés hautement décalées dans les jointures ou agrégations.
Récupération :
- Activez l'AQE (
spark.sql.adaptive.enabled=true) pour coalescer les partitions post-shuffle et appliquer la gestion des jointures décalées. - Salez les clés décalées (exemple construit : ajoutez un petit suffixe aléatoire aux clés avant
groupBy, puis agrégez à nouveau pour recombiner) quand l'AQE est insuffisant. - Utilisez les jointures broadcast quand un côté est assez petit.
5. Sorties Corrompues ou Partielles
Symptômes : Les jobs en aval échouent sur incompatibilité de schéma ou partitions manquantes ; fichiers partiels où le job a été interrompu.
Cause racine : Écritures non atomiques ou interruptions.
Récupération :
- Pour les écritures overwrite, privilégiez les modes qui remplacent des partitions entières atomiquement si votre stockage le supporte.
- Écrivez vers un chemin temporaire et renommez vers le final en cas de succès. Si une écriture échoue, supprimez le chemin temporaire et relancez.
- Rendez les puits idempotents : écritures partitionnées, clés déterministes, et pas d'effets de bord en milieu de tâche.
Conseils de Rollback
- Rollback de configuration : Gardez un petit ensemble de configs connues bonnes. Si un changement dégrade la stabilité ou la performance, revenez au dimensionnement d'executor et aux comptes de partitions de base et relancez.
- Rollback de code : Empaquetez le job pilote avec versioning sémantique. Si une nouvelle stratégie de jointure ou UDF cause des échecs, redéployez le dernier jar ou script connu bon et comparez les graphes d'étapes de l'UI Spark pour identifier la régression.
- Rollback de données : Les écritures overwrite sont les plus simples à annuler en restaurant depuis un snapshot ou en relançant avec la version de code précédente. Pour les puits JDBC, enveloppez les écritures dans des transactions si supporté, ou écrivez vers des tables de staging et échangez.
Exemples Pratiques Qui Enseignent l'Architecture
Exemple A : Vérifier le Comportement Broadcast et Shuffle
- Définissez
spark.sql.autoBroadcastJoinThreshold=64met chargez une table de dimension de 20 Mo. - Dans l'UI Spark, l'étape de jointure doit afficher un BroadcastHashJoin. Le plan physique contient un BroadcastExchange.
- Augmentez la table de dimension à 200 Mo. Relancez : le plan doit basculer vers SortMergeJoin ou ShuffledHashJoin ; l'onglet Stages montrera de grosses métriques de lecture/écriture de shuffle.
- Leçon : Le driver choisit les stratégies de jointure en utilisant les statistiques ; le broadcast peut réduire les shuffles, mais seulement quand c'est assez petit.
Exemple B : Bien Dimensionner les Partitions
- Commencez avec
spark.sql.shuffle.partitions=64(sur le cluster construit). Observez les durées de tâche et le GC dans l'UI. - Si les tâches sont constamment sous-seconde et que la surcharge d'ordonnancement domine, essayez 32.
- Si les tâches durent plusieurs secondes avec un GC élevé et une grande variance, essayez 96 ou 128.
- Leçon : Le nombre de partitions contrôle la concurrence et la granularité des tâches ; trouvez le point optimal en observant la latence de queue et le GC.
Liste de Contrôle Opérationnelle
Utilisez cette liste pour garder les exécutions prévisibles et les diagnostics rapides.
Planifier
- Confirmez les chemins d'entrée, les comptages de lignes attendus et les colonnes de partitionnement.
- Décidez des types de jointures attendus (broadcast vs shuffle) et vérifiez les tailles de tables.
- Enregistrez la config de base : instances d'executor, cœurs, mémoire, partitions de shuffle, AQE.
Configurer
- Appliquez la config de base de ce guide, en l'ajustant à votre matériel.
- Assurez-vous que les logs du driver et des executors sont conservés au moins pour la durée du job plus le temps de revue.
- Définissez le mode d'écriture sur overwrite vers un chemin temporaire, puis déplacez vers le final en cas de succès.
Exécuter
- Lancez le job et ouvrez l'UI Spark pour surveiller Jobs et Stages.
- Confirmez que le nombre d'executors et la distribution des tâches sont comme attendus.
Vérifier
- Validez les comptages de lignes et agrégats simples par rapport aux attentes.
- Inspectez
explain(True)pour les shuffles et décisions de broadcast. - Vérifiez le temps GC, lecture/écriture Shuffle, et décalage des tâches dans l'UI.
Récupérer (si nécessaire)
- Pour OOM : réduisez les cœurs par executor, augmentez la surcharge, ou décomposez les opérations larges.
- Pour le décalage : activez l'AQE, ajustez les partitions, envisagez le salage.
- Pour les problèmes de sortie : nettoyez le chemin temporaire et relancez, ou revenez au snapshot de code/config précédent.
Revoir
- Capturez le log d'événements de l'UI Spark pour l'exécution et annotez ce qui a fonctionné et ce qui n'a pas fonctionné.
- Ajustez deux boutons au maximum par itération (ex. partitions de shuffle et cœurs par executor) et relancez pour comparer.
Conclusion
Vous disposez maintenant d'un modèle mental fonctionnel de l'architecture de Spark et d'une méthode concrète pour l'exercer en toute sécurité : le driver construit et planifie les plans, les executors exécutent les tâches et échangent les données via les shuffles, et le gestionnaire de cluster alloue les ressources. Avec une configuration de base prudente, des étapes de vérification claires et des tactiques de récupération ciblées, vous pouvez exécuter un pilote étroit et mesurable et vous développer en toute confiance. Appliquez la liste de contrôle aux futurs jobs, ajustez en fonction du comportement observé dans l'UI Spark, et maintenez les chemins de rollback simples pour pouvoir vous adapter rapidement au fur et à mesure que les données, le code et les conditions du cluster évoluent.