Les tâches Spark peuvent échouer de manière mystérieuse lorsque la production est en feu et que l'horloge tourne. Ce guide explique les messages d'erreur Apache Spark courants, pourquoi ils se produisent, et les moyens sûrs de les corriger avec des exemples pratiques et reproductibles. Vous apprendrez à :
- Inventorier votre environnement pour ne pas poursuivre le mauvais correctif.
- Choisir d'abord des changements de configuration sûrs et limités.
- Vérifier le résultat en utilisant l'interface utilisateur Spark, les journaux et des vérifications ciblées.
- Comprendre les modes de défaillance et annuler proprement si nécessaire.
- Garder une liste de contrôle opérationnelle concise que vous pouvez suivre sous pression.
Les exemples utilisent du code d'échantillon construit et des nombres hypothétiques que vous pouvez adapter à votre pile. Appliquez chaque changement à un pilote étroit et mesurable et inspectez les résultats localement ou sur un espace de travail de développement dédié avant le déploiement.
Inventaire des Versions et de l'Environnement
De nombreuses erreurs Spark sont sensibles aux versions ou à l'environnement. Capturez ces faits en premier et joignez-les à tout incident ou enregistrement de changement.
Version et informations de construction Spark
CLI :
spark-submit --version
spark-shell --version
Depuis le REPL :
// Scala
spark.version
# Python
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
print(spark.version)
Langages d'exécution et Java
- Java :
java -version - Scala :
scala -version(si pertinent) - Python :
python --versionoupython3 --version
Gestionnaire de cluster et forme des ressources
- Déterminer le master :
- Scala :
spark.sparkContext.master - Python :
spark.sparkContext.master - Afficher la configuration actuelle (assainie) :
- Scala :
spark.sparkContext.getConf.getAll.foreach(println) - Python :
for k, v in spark.sparkContext.getConf().getAll():
print(k, v)
- Enregistrer le nombre d'exécuteurs, cœurs, mémoire. Si vous soumettez des tâches, consignez les indicateurs utilisés, par exemple :
--num-executors 8 --executor-cores 4 --executor-memory 6g --driver-memory 2g(exemple construit)
Points de terminaison de stockage et de système de fichiers
- HDFS :
hdfs dfsadmin -reportet un exemplehdfs dfs -ls /pathde vos racines d'entrée et de sortie. - Stockages objet : les noms de bucket/endpoint (ne pas journaliser les secrets), et les configurations Hadoop ou Spark utilisées pour l'authentification (masquées).
Bibliothèques et connecteurs
- Lister les JARs et packages supplémentaires fournis avec
--packagesou--jars. - Noter les paramètres de sérialisation (Java vs Kryo) et les configurations SQL pertinentes.
Avoir cet inventaire raccourcit le chemin du symptôme à la cause racine et rend les correctifs reproductibles entre environnements.
Chemin de Configuration Sûr
Traitez la configuration comme un scalpel, pas un marteau. Changez d'abord la plus petite portée, pour le temps le plus court, et vérifiez. Voici une séquence pratique que vous pouvez suivre pour la plupart des erreurs Spark.
- Reproduire sur un pilote étroit et mesurable
- Filtrer sur une petite partition ou plage de dates représentative.
- Garder les chemins d'entrée et de sortie séparés de la production.
- Privilégier d'abord les paramètres au niveau session/tâche
- Utiliser
--confsurspark-submitouspark.conf.set("clé", valeur)en session. - Éviter de modifier les valeurs par défaut du cluster avant que le correctif ne soit prouvé.
- Changer une variable à la fois
- Par exemple, augmenter d'abord
spark.sql.shuffle.partitions(ou diminuer), puis évaluer ; seulement ensuite toucher à la mémoire ou aux seuils de broadcast.
- Observer les effets immédiatement
- Suivre les temps d'exécution des étapes, les tailles de lecture/écriture de shuffle, les échecs de tâches, le temps GC, et les indicateurs de déséquilibre dans l'interface utilisateur Spark.
- Codifier le correctif
- Une fois stable, épingler la configuration dans le code (là où approprié) ou documenter les indicateurs de soumission exacts. Garder un paramètre de retour prêt.
Erreurs Courantes et Premiers Correctifs Sûrs
Le tableau ci-dessous associe les erreurs fréquentes à la cause principale et à un premier correctif sûr à essayer sur votre tâche pilote.
| Message d'erreur (abrégé) | Cause probable | Premier correctif sûr |
|---|---|---|
SparkException: Job aborted due to stage failure (FetchFailedException) | Shuffle volumineux, réseau instable, ou exécuteur perdu en plein shuffle | Réduire la taille du shuffle (repartitionner par clé, augmenter modérément les partitions), activer les tentatives, valider la stabilité des exécuteurs |
OutOfMemoryError: Java heap space (exécuteur) | Partitions trop grandes ou transformations larges matérialisant d'énormes lignes | Augmenter les partitions, persister sélectivement sur DISK ou MEMORY_AND_DISK, éviter .collect() sur de gros RDD/DataFrame |
OutOfMemoryError: Java heap space (driver) | Driver collectant ou conservant de grosses structures | Remplacer collect() par take(), écrire vers le stockage, déplacer la logique vers les exécuteurs |
SparkException: Task not serializable | Objet non sérialisable capturé dans une fermeture | Utiliser une case class sérialisable, diffuser les petites tables de correspondance, utiliser mapPartitions avec fabrique locale |
AnalysisException: Path already exists when saving | Chemin de sortie non nettoyé ou mode de sauvegarde non défini | Utiliser .mode("overwrite") intentionnellement, ou supprimer le chemin cible avant l'écriture |
ClassNotFoundException (ex: source Kafka) | JAR/package de connecteur manquant | Ajouter le bon package via --packages correspondant à la version Spark/Scala |
Py4JJavaError (wrapper générique) | Exception JVM sous-jacente remontée vers Python | Développer la cause dans les journaux, lire la cause la plus interne pour choisir un correctif ciblé |
| Déséquilibre : queues longues dans quelques tâches | Clés chaudes ou partitions déséquilibrées | Salage/indices de jointure déséquilibrés, AQE activé, ou partitionneur personnalisé |
Vérification et Diagnostics
Transformez un échec vague en diagnostic concret avec ces étapes et vérifications.
1) Inspecter l'Interface Utilisateur Spark et les Journaux
- Onglet Stages : rechercher des lectures/écritures de shuffle très grandes, ou des tâches avec des durées extrêmes comparées aux pairs (déséquilibre).
- Onglet Executors : vérifier les compteurs d'exécuteurs perdus, le % de temps GC, et la mémoire maximale.
- Journaux driver et exécuteurs : trouver la première exception de cause racine (pas le wrapper). Pour PySpark, développer la pile Py4J jusqu'à l'exception Java la plus interne.
2) Reproduire avec une Requête ou Fonction Minimale
- Créer un petit jeu de données représentatif et exécuter la transformation problématique.
- Laisser un seul changement suspecté à la fois et comparer les métriques avant/après.
3) Diagnostics Concrets et Exemples
A) FetchFailedException pendant une jointure ou agrégation
Symptôme d'exemple construit :
- Erreur :
org.apache.spark.SparkException: Job aborted due to stage failure: FetchFailed ...après un shuffle large.
Vérifications et correctifs :
- Vérifier la taille du shuffle : en SQL, lire l'onglet SQL ou en code afficher les comptes de partitions.
- Augmenter modérément les partitions de shuffle pour SQL :
spark.conf.set("spark.sql.shuffle.partitions", 400) # exemple construit
- Si vous utilisez l'API RDD, utiliser
rdd.repartition(400)ou, pour des clés connues, utiliserpartitionBy(numPartitions)sur les PairRDDs. - Activer l'exécution de requête adaptative (AQE) si disponible :
spark.conf.set("spark.sql.adaptive.enabled", True)
- Vérifier : relancer l'étape ; comparer la lecture de shuffle maximale par tâche et le nombre d'échecs de fetch.
B) OutOfMemoryError de l'Exécuteur
Symptôme d'exemple construit :
- Erreur :
java.lang.OutOfMemoryError: Java heap spacedans les journaux d'exécuteur pendant ungroupBy().agg()sur des lignes larges.
Étapes sûres :
- Réduire la taille des partitions en augmentant le parallélisme prudemment :
spark.conf.set("spark.sql.shuffle.partitions", 400)
- Persister sélectivement en utilisant un niveau de stockage qui peut déborder :
df_to_reuse = df_heavy.transform(some_fn)
df_to_reuse.persist(storageLevel="MEMORY_AND_DISK")
- Éviter
.collect()ou.toPandas()sur de gros jeux de données. Utiliser.limit().collect()pour l'échantillonnage ou écrire vers le stockage. - Si l'échec persiste, augmenter modérément la mémoire d'exécuteur (exemple construit) :
--executor-memory 6g --executor-cores 4- Vérifier : surveiller le temps GC de l'exécuteur et les tâches échouées ; confirmer l'absence d'OOM dans les journaux.
C) OutOfMemoryError du Driver
Symptôme d'exemple construit :
- Erreur sur le driver :
java.lang.OutOfMemoryError: Java heap spaceaprès.collect()sur un gros DataFrame.
Correctifs :
- Remplacer
.collect()par.take(1000)pour l'échantillonnage, ou écrire les résultats via.writevers le stockage. - Si vous devez conserver des données modérées sur le driver, augmenter modérément la mémoire du driver (construit) :
--driver-memory 4g. - Vérifier : la tâche se termine ; les journaux du driver montrent un tas stable sans thrash GC.
D) Tâche Non Sérialisable
Exemple Scala construit du problème :
case class Item(id: Long, v: Double)
class NonSerializable(val factor: Double)
val helper = new NonSerializable(2.0)
val rdd = sc.parallelize(Seq(Item(1, 3.0), Item(2, 4.0)))
val out = rdd.map(x => x.v * helper.factor).collect() // échoue
Options de correctif :
- Rendre l'état sérialisable ou éviter de le capturer dans la fermeture :
case class SerializableHelper(factor: Double) extends Serializable
val helper = SerializableHelper(2.0)
val out = rdd.map(x => x.v * helper.factor).collect()
- Ou instancier par partition :
val out = rdd.mapPartitions { it =>
val helper = new NonSerializable(2.0)
it.map(x => x.v * helper.factor)
}.collect()
- Vérifier : la tâche s'exécute ; pas d'erreur
Task not serializabledans les journaux d'exécuteur.
E) AnalysisException : Le Chemin Existe Déjà
Exemple construit :
df.write.mode("errorifexists").parquet("/data/out/2024-01-01")
Correctifs :
- Utiliser
.mode("overwrite")intentionnellement quand c'est sûr :
df.write.mode("overwrite").parquet("/data/out/2024-01-01")
- Ou supprimer d'abord le chemin cible (exemple HDFS) :
hdfs dfs -rm -r /data/out/2024-01-01
- Vérifier : le chemin affiche les nouveaux fichiers et la tâche a réussi.
F) ClassNotFoundException pour Connecteurs Externes (ex: Kafka)
Symptôme d'exemple construit :
- Erreur :
java.lang.ClassNotFoundException: org.apache.spark.sql.kafka010.KafkaSourceProvider
Correctif :
- Ajouter le bon package correspondant à votre version Spark et Scala au moment de la soumission. Modèle (construit) :
--packages org.apache.spark:spark-sql-kafka-0-10_2.12:<version_spark>
- Vérifier : la session démarre et
readStream.format("kafka")ne produit plus d'erreur.
G) Déséquilibre Causant des Queues Longues
Diagnostic :
from pyspark.sql import functions as F
(df.groupBy("key").count()
.orderBy(F.desc("count"))
.limit(10)
.show())
Correctifs :
- Activer AQE pour gérer les jointures déséquilibrées quand supporté :
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", True)
- Appliquer le salage (construit) :
import pyspark.sql.functions as F
salt_buckets = 10
left_salted = left.withColumn("salt", F.rand(seed=42) * salt_buckets).withColumn("salt", F.floor(F.col("salt")))
right_salted = right.withColumn("salt", F.lit(0))
joined = left_salted.join(right_salted, ["key", "salt"], "left")
- Vérifier : les tâches de queue raccourcissent ; moins de traînards dans l'onglet Stages.
Modes de Défaillance et Récupération
Même un changement raisonnable peut avoir des effets secondaires. Utilisez cette section pour anticiper les risques et garder un plan de retour.
- Augmenter
spark.sql.shuffle.partitionstrop haut - Risque : Trop de minuscules tâches et petits fichiers de sortie ; surcharge d'ordonnancement plus élevée.
- Récupération : Revenir à la valeur précédente ; compacter les petits fichiers avec un
coalesce/repartitionciblé à l'écriture.
- Désactiver les jointures broadcast globalement
- Risque : Gros shuffles remplaçant des broadcasts efficaces ; temps d'exécution plus longs.
- Récupération : Restaurer le seuil de broadcast ; utiliser des indices de jointure seulement là où nécessaire.
- Sur-allouer la mémoire d'exécuteur
- Risque : Moins d'exécuteurs tiennent sur le cluster ; files d'attente plus longues ; possible arrêt de conteneur par le gestionnaire de ressources.
- Récupération : Revenir à la mémoire précédente ; envisager d'augmenter les partitions ou les niveaux de stockage favorables au débordement à la place.
- Mise en cache agressive sans unpersist
- Risque : Le cache remplit la mémoire d'exécuteur ; thrashing d'éviction et OOMs.
- Récupération : Appeler
unpersist()sur les DataFrames/RDDs plus nécessaires ; redémarrer l'application si la mémoire est fragmentée.
- Écrasement non intentionnel des chemins de sortie
- Risque : Perte de données.
- Récupération : Arrêter la tâche, restaurer depuis la sauvegarde ou la partition précédente si disponible. Imposer une garde qui compare les compteurs d'enregistrements ou vérifie les filigranes avant l'écrasement.
- FetchFailed fréquents dus à des nœuds instables
- Risque : Nouvelles tentatives répétées et queues longues.
- Récupération : Mettre en quarantaine les mauvais exécuteurs ou nœuds ; resoumettre après la stabilisation de la santé des nœuds.
Liste de Contrôle Opérationnelle
Utilisez cette courte liste de contrôle lors de la gestion des erreurs Spark en astreinte ou pendant les post-mortems.
- Inventaire
- Enregistrer les versions Spark, Java/Scala/Python.
- Capturer le master, exécuteurs, cœurs, mémoire, et indicateurs
--confclés. - Noter les chemins d'entrée/sortie et packages de connecteurs.
- Reproduire sur un Pilote
- Rétrécir le jeu de données par date/partition.
- Sauvegarder les journaux, captures d'écran de l'interface utilisateur Spark, et le code/requête minimal qui échoue.
- Diagnostiquer
- Interface utilisateur Spark : trouver la première étape échouante et la plus grande lecture/écriture de shuffle.
- Journaux : extraire l'exception la plus interne et la pile d'appels.
- Vérifier le déséquilibre via les comptes des clés principales ou durées de tâches extrêmes.
- Appliquer un Correctif Limité
- Privilégier
--confouspark.conf.setplutôt que les valeurs par défaut du cluster. - Changer une variable à la fois ; documenter les effets attendus et observés.
- Vérifier
- Comparer les métriques avant/après : temps d'étape, tailles de shuffle, % GC, échecs.
- Valider les compteurs de lignes ou sommes de contrôle des sorties.
- Retour
- Garder les valeurs de configuration précédentes à portée de main ; revenir rapidement si des effets secondaires apparaissent.
- Nettoyer les sorties partielles et relancer idempotemment.
- Clôture
- Codifier le correctif dans le code ou les scripts de soumission.
- Mettre à jour les manuels d'exploitation et ajouter un test de garde là où c'est faisable.
Référence de Configuration Pratique
Utilisez ces boutons avec précaution et privilégiez d'abord les changements au niveau tâche. Les valeurs ci-dessous sont des exemples construits ; ajustez selon les mesures.
| Paramètre | Portée | Changement de départ sûr | Effet |
|---|---|---|---|
spark.sql.shuffle.partitions | Session/tâche SQL | 200 → 400 | Tailles de partition plus petites réduisent la pression mémoire par tâche ; peut augmenter le nombre de tâches |
spark.sql.adaptive.enabled | Session/tâche SQL | false → true | Laisse Spark optimiser les jointures/shuffles à l'exécution ; aide avec le déséquilibre |
spark.sql.autoBroadcastJoinThreshold | Session/tâche SQL | défaut → 50MB | Ajuster pour permettre/éviter les jointures broadcast ; éviter de broadcaster d'énormes tables |
spark.serializer | App/cluster | Java → Kryo (avec enregistrement) | Sérialisation plus rapide pour pipelines RDD lourds ; valider la compatibilité |
spark.executor.memory | Tâche/cluster | +1-2g | Plus de tas par exécuteur ; réduit les OOMs si les partitions restent trop grandes |
spark.executor.cores | Tâche/cluster | 4 → 2-4 | Ajuster pour équilibrer CPU vs mémoire par tâche ; éviter la sur-souscription |
spark.memory.fraction | Tâche/cluster | défaut | Rarement changer en premier ; privilégier le partitionnement et les niveaux de stockage favorables au débordement |
Exemple Travaillé de Bout en Bout (Construit)
Scénario : Une jointure quotidienne entre une table de faits de 120M lignes et une table de dimension de 200K lignes échoue avec FetchFailedException et OOM occasionnels d'exécuteur.
- Inventaire
- Spark 3.3.x, master sur YARN, 12 exécuteurs, 4 cœurs chacun, 6g mémoire d'exécuteur.
spark.sql.shuffle.partitions=200(défaut), AQE désactivé.
- Reproduire sur un Pilote
- Filtrer les faits sur une seule journée (~5% des données) et exécuter la jointure dans un espace de développement.
- Diagnostiquer
- L'interface utilisateur Spark montre 1,2 To de lecture de shuffle sur la grosse jointure ; quelques tâches lisant 20+ Go chacune ; queues longues.
- Correctifs Limités
- Activer AQE et gestion du déséquilibre :
spark.conf.set("spark.sql.adaptive.enabled", True)
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", True)
- Augmenter les partitions de shuffle à 400 pour réduire la charge par tâche :
spark.conf.set("spark.sql.shuffle.partitions", 400)
- Vérifier que la dimension de 200K lignes fait du broadcast au lieu de shuffler :
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", 50 * 1024 * 1024)
- Vérifier
- La lecture de shuffle chute par tâche ; plus de
FetchFailed; le temps d'exécution passe de 45 min à 18 min sur le pilote. Pas d'OOM dans les journaux d'exécuteur.
- Déploiement
- Cuire ces configurations de session dans la soumission de la tâche ; laisser les valeurs par défaut du cluster inchangées.
- Surveiller la prochaine exécution complète ; garder une option de retour qui ramène les partitions de shuffle à 200 et désactive AQE si des régressions apparaissent.
- Clôture
- Documenter les nouveaux paramètres et l'impact mesuré. Ajouter une entrée de manuel d'exploitation pour la gestion future des jointures déséquilibrées.
Conclusion
La plupart des erreurs Spark se réduisent à quelques causes racines : trop de travail par tâche, partitions déséquilibrées, dépendances manquantes, ou motifs dangereux comme la collecte de gros jeux de données sur le driver. En inventoriant votre environnement, en appliquant d'abord des changements limités et mesurables, en vérifiant avec l'interface utilisateur Spark et les journaux, et en gardant des étapes de retour claires, vous pouvez transformer la réponse chaotique aux incidents en une pratique délibérée à faible risque. Commencez petit, mesurez, et promouvez le correctif seulement quand vous pouvez montrer qu'il fonctionne sur un pilote. Cette approche réduit régulièrement les incidents répétés et rend les opérations quotidiennes plus calmes et plus rapides.