E-NO
Dépannage Apache Spark 11 min de lecture

Dépannage d'Apache Spark avec exemples pratiques : guide d'implémentation

calendar_today Publié : 2026-08-13
update Dernière mise à jour : 2026-08-13
analytics Efficacité SEO : 100%
Illustration du guide technique pour « Dépannage d'Apache Spark avec exemples pratiques : guide d'implémentation ».

Les tâches Apache Spark échouent pour un nombre restreint de raisons récurrentes : pression sur les ressources, problèmes de shuffle, erreurs de sérialisation, dépendances manquantes, problèmes de sources de données ou dérive d'environnement. Ce guide vous propose un chemin pratique pour isoler la cause, modifier les paramètres en toute sécurité et récupérer sans amplifier le rayon d'impact. Vous allez inventorier votre environnement, exécuter une reproduction ciblée, lire les bons journaux, essayer de petits changements réversibles et vérifier la correction avant de passer à l'échelle.

Objectifs clés :

  • Réduire le temps de diagnostic grâce à un inventaire simple et une reproduction minimale.
  • Appliquer d'abord des changements réversibles au niveau de la tâche ; étendre seulement après vérification.
  • Associer les symptômes observables aux modes de défaillance courants de Spark et appliquer des correctifs sûrs.

Inventaire des versions et de l'environnement

Avant de changer quoi que ce soit, rassemblez ces faits. Ils éliminent des catégories entières d'erreurs et accélèrent chaque étape ultérieure.

  1. Confirmez les versions de Spark, Scala, Java et Hadoop.

Sur un nœud edge ou driver :

spark-submit --version
pyspark --version          # si vous utilisez PySpark
scala -version             # si vous utilisez les API Scala
java -version
hadoop version             # si vous vous connectez à HDFS/YARN

Notez les versions exactes et les numéros de build. Les incompatibilités (par exemple, Spark 3.3 compilé pour Scala 2.12 mais un jar compilé pour 2.13) se manifestent souvent par des ClassNotFoundException ou NoSuchMethodError.

  1. Notez le mode d'exécution et la topologie.
  • Local, cluster autonome ou mode cluster YARN.
  • Où s'exécute le driver (mode client vs mode cluster).
  • Sources de données en jeu (chemins HDFS, buckets de stockage objet, topics Kafka) et leur accessibilité réseau.
  1. Localisez les configurations et la journalisation.
  • Défauts et surcharges Spark :
  • /etc/spark/conf/spark-defaults.conf
  • /etc/spark/conf/spark-env.sh
  • /etc/spark/conf/log4j2.properties (ou l'ancien log4j.properties)
  • Les configurations au niveau de la tâche sont préférables pour des tests rapides et réversibles : spark-submit --conf cle=valeur.
  1. Capturez la SparkConf à l'exécution.

Dans votre application, affichez les paramètres effectifs dès le début :

Scala :

spark.sparkContext.getConf.getAll.foreach{ case (k, v) => println(s"CONF: $k=$v") }

Python :

for k, v in spark.sparkContext.getConf().getAll():
    print(f"CONF: {k}={v}")
  1. Vérifiez l'accessibilité des sources de données.

Exemple HDFS :

hdfs dfs -ls /data/input
hdfs dfs -test -e /data/input && echo existe || echo manquant

Exemple permissions système de fichiers :

hdfs dfs -ls /data | grep -E "^d|^-"   # visibilité rapide sur perms/propriétaires

Chemin de configuration sûr

Apportez le plus petit changement sûr en premier. Préférez les paramètres au niveau de la tâche et de minuscules entrées. Ne promouvez à une portée plus large qu'après avoir confirmé l'effet.

Principes :

  • Utilisez une exécution pilote étroite et mesurable qui reproduit l'échec rapidement. Par exemple, filtrez sur une seule partition ou un échantillon de 100 000 lignes.
  • Appliquez les surcharges via spark-submit --conf ou le constructeur de Session, pas en modifiant les fichiers du cluster dans un premier temps.
  • Augmentez la verbosité de la journalisation au driver et aux exécuteurs juste assez longtemps pour capturer l'échec, puis revenez en arrière.

Exemples d'ajustements sûrs au niveau de la tâche :

  1. Augmentez ou diminuez la verbosité des logs pour cette exécution seulement.

Scala :

spark.sparkContext.setLogLevel("INFO")  // ou WARN/ERROR pour contrôler le bruit

Python :

spark.sparkContext.setLogLevel("INFO")
  1. Surchargez le dimensionnement des exécuteurs avec prudence (valeurs d'exemple).
spark-submit \
  --executor-memory 4g \
  --conf spark.executor.cores=2 \
  --conf spark.executor.instances=4 \
  --class com.example.Job \
  app.jar
  1. Limitez la pression du shuffle pendant le diagnostic (valeurs d'exemple).
spark-submit \
  --conf spark.sql.shuffle.partitions=100 \
  --conf spark.default.parallelism=100 \
  --class com.example.Job \
  app.jar
  1. Ajoutez des dépendances explicitement pour cette exécution.
spark-submit \
  --jars deps/udf-lib.jar,libs/metrics.jar \
  --py-files deps/util.zip \
  --class com.example.Job \
  app.jar

Promotions après vérification :

  • Si un paramètre corrige systématiquement une classe d'échecs, intégrez-le dans votre empaquetage de tâche ou dans spark-defaults.conf.
  • Évitez les changements à l'échelle du cluster jusqu'à ce que plusieurs tâches aient validé la nouvelle ligne de base.

Vérification et diagnostics

Utilisez le plus petit jeu de données qui échoue encore pour accélérer les itérations. Observez sous trois angles : journaux, interface utilisateur Spark et validation de l'environnement.

Accès rapide aux journaux et à l'interface :

ModeJournaux driverJournaux exécutateursInterface Spark
Autonome$SPARK_HOME/logs sur nœuds master/worker$SPARK_HOME/logs sur nœuds workerhttp://hote-driver:4040 pendant l'exécution
YARNyarn logs -applicationId <app_id>yarn logs -applicationId <app_id>Serveur d'historique Spark si activé
  1. Identifiez l'application.
  • Depuis la sortie de spark-submit, capturez l'applicationId (YARN) ou le nom de l'application.
  • Si perdu, scannez les applications récentes par heure dans le Serveur d'historique Spark et faites correspondre votre nom de tâche et l'heure de début.
  1. Récupérez les journaux avec des filtres pertinents.
yarn logs -applicationId <app_id> | grep -i -E "exception|error|killed|memory|fetchfailed|timeout" -n
  1. Capturez la signature de l'échec.

Signatures courantes :

  • java.lang.OutOfMemoryError, GC overhead limit exceeded
  • org.apache.spark.shuffle.FetchFailedException
  • Task not serializable
  • ClassNotFoundException ou NoSuchMethodError
  • Container killed by YARN for exceeding memory limit
  • Py4JJavaError avec trace de pile Java imbriquée (lisez la cause racine près de "Caused by:")
  1. Inspectez l'interface utilisateur Spark (Jobs -> Stages -> Tasks).
  • Indicateur de déséquilibre : une ou quelques tâches s'exécutent beaucoup plus longtemps ou traitent bien plus d'entrées.
  • Les pics de taille de lecture de shuffle sont visibles dans le détail des Stages.
  • Nombre de tâches échouées/tuées, plus messages d'erreur par tâche.
  1. Validez que les entrées et sorties existent et sont accessibles en écriture.
hdfs dfs -ls /data/output_tmp && echo ok || echo manquant
hdfs dfs -test -w /data/output_tmp && echo accessible_en_ecriture || echo non_accessible_en_ecriture
  1. Reproduisez minimalement.

Si la tâche lit depuis une grande table, sélectionnez un petit sous-ensemble représentatif :

Scala :

val sampleDf = spark.read.parquet("/data/input").limit(100000)
sampleDf.repartition(16).groupBy("key").count().collect()

Python :

sample_df = spark.read.parquet("/data/input").limit(100000)
sample_df.repartition(16).groupBy("key").count().collect()

Attendez-vous à voir la même classe d'échec (par exemple, FetchFailedException) plus rapidement. Si l'échec disparaît, la cause racine peut être liée à l'échelle (déséquilibre, pression mémoire ou timeouts).

Modes de défaillance et récupération

Cette section associe les erreurs Spark fréquentes aux causes probables, correctifs sûrs, étapes de vérification et conseils de retour arrière.

1) OutOfMemoryError ou GC overhead limit exceeded

Symptômes :

  • Le driver ou les exécuteurs meurent avec OutOfMemoryError.
  • Message YARN : Container killed by YARN for exceeding memory limit.

Causes probables :

  • Partitions trop grandes ; les transformations larges matérialisent de gros shuffles.
  • Les actions DataFrame collectent trop de données vers le driver.
  • Mise en cache de grandes tables non filtrées sans assez de mémoire.

Correctifs sûrs (appliquez de façon incrémentale) :

  • Réduisez la taille des partitions ; augmentez le nombre de partitions pour les shuffles.
  • Évitez collect() vers le driver sur de gros jeux de données ; utilisez take(n) ou show(n, truncate=false) avec parcimonie.
  • Augmentez modérément la mémoire et l'overhead des exécuteurs ; vérifiez les limites des conteneurs.

Exemples d'ajustements (valeurs d'exemple) :

spark-submit \
  --conf spark.sql.shuffle.partitions=400 \
  --conf spark.executor.memoryOverhead=1024 \
  --executor-memory 6g \
  app.jar

Vérification :

  • L'interface Spark montre moins de spill de shuffle et une mémoire de tâche stable.
  • Aucune nouvelle entrée OOM dans les journaux des exécuteurs.

Retour arrière :

  • Si l'utilisation du cluster grimpe en flèche ou si d'autres tâches meurent de faim, annulez les augmentations de mémoire et privilégiez les corrections de partitionnement.

2) FetchFailedException (échec de récupération du shuffle)

Symptômes :

  • Les stages échouent répétitivement avec FetchFailedException.
  • Les tâches réessaient de nombreuses fois ; finalement le stage avorte.

Causes probables :

  • Exécuteur perdu qui a supprimé des blocs en cours de stage.
  • Service de shuffle qui ne sert pas les blocs de manière fiable.
  • Instabilité réseau ou timeouts lors de la récupération de gros blocs.
  • Déséquilibre sévère concentrant les données de shuffle sur quelques exécutateurs.

Correctifs sûrs :

  • Augmentez les paramètres de retry et de timeout du shuffle (valeurs d'exemple) :
spark-submit \
  --conf spark.shuffle.io.maxRetries=10 \
  --conf spark.shuffle.io.retryWait=5s \
  --conf spark.network.timeout=600s \
  app.jar
  • Activez et vérifiez le service de shuffle externe lors de l'utilisation de l'allocation dynamique sur YARN ou en mode autonome :
# niveau tâche
--conf spark.shuffle.service.enabled=true \
--conf spark.dynamicAllocation.enabled=true
  • Réduisez le déséquilibre en augmentant les partitions de shuffle et en pré-agrégant.

Vérification :

  • Les stages progressent sans échecs de récupération répétés ; moins de tâches perdues.
  • L'interface montre des durées de tâche plus uniformes après atténuation du déséquilibre.

Retour arrière :

  • Si des timeouts plus longs masquent simplement des problèmes réseau, réduisez-les et investiguez le nœud défaillant ; excluez temporairement les hôtes instables au niveau du gestionnaire de ressources si supporté.

3) Task not serializable (Tâche non sérialisable)

Symptômes :

  • Échec peu après le démarrage de la tâche.
  • Trace de pile mentionnant java.io.NotSerializableException.

Causes probables :

  • Les fermetures capturent des objets non sérialisables (par exemple, clients de base de données, loggers, SparkSession) à l'intérieur de map/flatMap.

Correctifs sûrs :

  • Sortez les objets non sérialisables des fermetures ; passez des paramètres primitifs.
  • Utilisez des variables de diffusion (broadcast) pour les grandes structures de données en lecture seule.

Exemple de correction en Scala :

val bcConfig = spark.sparkContext.broadcast(heavyConfig)
rdd.map(x => compute(x, bcConfig.value))

Vérification :

  • Le stage s'exécute ; plus de NotSerializableException dans les journaux.

Retour arrière :

  • Aucun nécessaire au-delà de la réversion des changements de code si le comportement régresse ; les tests devraient couvrir les limites des fermetures.

4) ClassNotFoundException / NoSuchMethodError

Symptômes :

  • Échecs à l'exécution lors de l'accès aux classes ou UDF.

Causes probables :

  • Jars manquants ou mauvaise version binaire de Scala.
  • Dépendance ombrée (shaded) sous un paquet différent ou conflit sur le classpath.

Correctifs sûrs :

  • Fournissez tous les jars nécessaires au moment de la soumission :
spark-submit --jars libs/udfs.jar,libs/dep.jar --class com.example.Job app.jar
  • Assurez-vous que la version binaire de Scala correspond au build Spark (par exemple, 2.12).
  • Pour PySpark, transmettez les fichiers zip pour les modules :
spark-submit --py-files deps/util.zip job.py

Vérification :

  • La classe en échec se charge ; l'enregistrement de l'UDF réussit.

Retour arrière :

  • Si un nouveau jar a introduit des conflits, supprimez-le et isolez avec l'ombrage (shading) dans un build ultérieur.

5) Déséquilibre de données causant des queues longues et timeouts

Symptômes :

  • 1 % des tâches s'exécutent 10 à 100 fois plus longtemps.
  • Pics de taille de lecture de shuffle sur quelques tâches.

Causes probables :

  • Clés hautement déséquilibrées dans groupBy/join.

Correctifs sûrs :

  • Augmentez les partitions de shuffle (exemple : 200 -> 800) pour réduire la charge par tâche.
  • Appliquez une pré-agrégation (combination côté map) avant les opérations larges.
  • Salez les clés pour un déséquilibre extrême ; par exemple, dupliquez les clés avec un petit suffixe aléatoire pour répartir la charge, puis agrégez les résultats.

Exemple d'ajustement (valeurs d'exemple) :

spark-submit --conf spark.sql.shuffle.partitions=800 app.jar

Vérification :

  • Les durées de tâche sont plus uniformes ; aucune tâche unique ne domine le temps du stage.

Retour arrière :

  • Si un nombre de partitions plus élevé augmente la surcharge, revenez à une valeur équilibrée et combinez avec la pré-agrégation.

6) Py4JJavaError et incompatibilités d'environnement Python

Symptômes :

  • Py4JJavaError enveloppant une trace de pile Java.
  • Erreurs sur l'incompatibilité de version Python ou modules manquants sur les exécutants.

Causes probables :

  • Versions Python différentes sur le driver et les exécutants.
  • Dépendances PySpark manquantes sur les workers.

Correctifs sûrs :

  • Définissez l'interpréteur Python de manière cohérente :
export PYSPARK_PYTHON=python3
export PYSPARK_DRIVER_PYTHON=python3
  • Expédiez les modules Python avec --py-files et importez-les dans la tâche.

Vérification :

  • Les imports de modules réussissent sur les exécutants ; les erreurs disparaissent.

Retour arrière :

  • Annulez les changements d'interpréteur s'ils cassent d'autres tâches ; limitez les changements d'environnement au shell de soumission ou au script wrapper.

7) Erreurs de chemin ou permissions HDFS

Symptômes :

  • FileNotFoundException ou AccessControlException.

Causes probables :

  • Mauvais chemin, répertoire de partition manquant ou permissions insuffisantes.

Correctifs sûrs :

  • Validez avec les commandes HDFS :
hdfs dfs -ls /data/input/date=2024-01-01
  • Demandez ou définissez les bonnes permissions ; évitez d'exécuter en superutilisateur pour les tâches de routine.

Vérification :

  • Les chemins se résolvent ; les écritures réussissent vers un chemin de sortie temporaire.

Retour arrière :

  • Aucun, à part la réversion de tout changement de permission non aligné avec la politique.

8) Tâches bloquées et progression arrêtée

Symptômes :

  • Aucune progression de stage ; les exécutants semblent inactifs.

Causes probables :

  • Ressources insuffisantes dans la file d'attente.
  • Blocages sur systèmes externes (par exemple, metastore lent ou système de fichiers distant).
  • Attente de localité très longue retardant l'ordonnancement des tâches.

Correctifs sûrs :

  • Réduisez spark.locality.wait pour déplacer le travail sans attendre la localité idéale (valeur d'exemple) :
--conf spark.locality.wait=1s
  • Dimensionnez correctement les instances d'exécutants ; des exécutants trop grands causent une mauvaise utilisation des créneaux.

Vérification :

  • Les stages démarrent rapidement ; pas de longs intervalles avant le lancement des tâches.

Retour arrière :

  • Si le trafic réseau grimpe, restaurez les attentes de localité par défaut et investiguez le système de stockage.

Carte rapide erreur-action

Utilisez cette carte compacte pendant un incident pour accélérer les premières actions.

Signature d'erreurCause probableÀ vérifierPremier correctif sûr
OutOfMemoryErrorGrosses partitions, collect vers driverJournaux exécutants ; métriques de stage ; spillAugmenter partitions ; augmentation mémoire modérée
FetchFailedExceptionBlocs perdus, timeouts, déséquilibreExécutants perdus ; retries shuffleAugmenter retries/timeouts ; activer service shuffle
Task not serializableFermeture capture non-sérialisableCode dans map/flatMapDiffuser (broadcast) ou refactorer fermetures
ClassNotFoundExceptionJar manquant ou incompatibilité versionArguments spark-submit ; version ScalaAjouter --jars/--py-files ; aligner versions
Tâche bloquéeAttente localité, famine ressourcesInterface montre exécutants inactifsBaisser attente localité ; dimensionner exécutants

Patterns de récupération et retour arrière

Les écritures partielles ou échouées peuvent laisser un état désordonné. Privilégiez les patterns qui permettent de confirmer avant de rendre visible.

  1. Sortie atomique avec temporaire-et-renommage.
  • Écrivez vers un chemin temporaire, puis renommez en cas de succès :

Scala :

val tmp = "/data/output_tmp/run_20240101"
df.write.mode("overwrite").parquet(tmp)
// après écriture réussie
import org.apache.hadoop.fs.{FileSystem, Path}
val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration)
fs.delete(new Path("/data/output"), true)
fs.rename(new Path(tmp), new Path("/data/output"))

Python :

tmp = "/data/output_tmp/run_20240101"
df.write.mode("overwrite").parquet(tmp)
from py4j.java_gateway import java_import
jconf = spark.sparkContext._jsc.hadoopConfiguration()
fs = spark._jvm.org.apache.hadoop.fs.FileSystem.get(jconf)
fs.delete(spark._jvm.org.apache.hadoop.fs.Path("/data/output"), True)
fs.rename(spark._jvm.org.apache.hadoop.fs.Path(tmp), spark._jvm.org.apache.hadoop.fs.Path("/data/output"))

Retour arrière :

  • Si la validation échoue, supprimez le chemin temporaire ; le chemin visible reste intact.
  1. Points de contrôle (checkpoints) et réexécutions incrémentielles.
  • Utilisez les checkpoints DataFrame pour couper la lignée quand les échecs répétés proviennent de plans profonds :
spark.sparkContext.setCheckpointDir("/checkpoints/jobA")
val stabilized = df.checkpoint(eager = true)
  • Réexécutez depuis le dernier bon checkpoint après un échec transitoire.

Retour arrière :

  • Si un changement augmente le temps d'exécution ou le coût, annulez la config et réutilisez le checkpoint précédent.
  1. Retour arrière de configuration.
  • Gardez les configs au niveau de la tâche dans un wrapper de soumission ou un fichier de définition de tâche avec un bloc précédent-connu-bon que vous pouvez restaurer instantanément.

Extrait de fichier de soumission d'exemple (pseudo-shell) :

# connu-bon
CONF=(
  "spark.sql.shuffle.partitions=200"
  "spark.network.timeout=120s"
)
# expérimental
# CONF+=("spark.dynamicAllocation.enabled=true")

Un pilote minimal et mesurable

Validez chaque changement avec une exécution pilote que vous pouvez inspecter rapidement. Cela réduit le retravail et l'ambiguïté quand plusieurs boutons bougent à la fois.

  • Reproduisez l'échec sur un petit jeu de données (par exemple, limit(100k) avec clés représentatives).
  • Changez un seul paramètre.
  • Capturez deux mesures : temps d'exécution et taux d'erreur (0 vs 1 échec) pour ce stage.
  • Si le pilote passe, multipliez l'entrée par 5-10x et observez à nouveau avant le volume complet.

Liste de contrôle opérationnelle

Exécutez ceci quand une tâche Spark échoue ou se comporte mal.

Prévol (2-5 minutes) :

  • Identifiez l'applicationId et le nom de la tâche.
  • Enregistrez les versions Spark, Scala, Java, Hadoop.
  • Confirmez le mode d'exécution (client vs cluster) et les sources de données impliquées.

Journaux et interface (5-10 minutes) :

  • Récupérez les journaux driver et exécutants ; grep pour les signatures d'erreur.
  • Ouvrez l'interface Spark ; notez le déséquilibre, tailles de shuffle et tâches échouées.
  • Validez que les chemins d'entrée/sortie existent et sont accessibles en écriture.

Classifiez l'échec (2-5 minutes) :

  • Pression mémoire : OutOfMemoryError, conteneur tué pour mémoire.
  • Problème de shuffle : FetchFailedException, exécuteur perdu.
  • Sérialisation : NotSerializableException.
  • Dépendance : ClassNotFoundException.
  • Environnement : erreurs Py4J, incompatibilité Python.
  • Source/permission : FileNotFound, AccessControl.

Appliquez le premier correctif sûr (5-15 minutes) :

  • Mémoire : augmentez partitions ; augmentation modérée mémoire/overhead.
  • Shuffle : augmentez retries/timeouts ; vérifiez service shuffle ; atténuez déséquilibre.
  • Sérialisation : refactorisez fermetures ; diffusez configuration/données.
  • Dépendance : ajoutez --jars/--py-files ; alignez version Scala.
  • Environnement : définissez Python cohérent ; expédiez modules.
  • Source : corrigez chemin/perms ; validez avec hdfs dfs.

Pilote et vérification (5-20 minutes) :

  • Réexécutez avec un jeu de données étroit.
  • Confirmez depuis les journaux et l'interface que la signature a disparu.
  • Si amélioré, passez à l'échelle au jeu de données complet.

Stabilisez et documentez (5-10 minutes) :

  • Promouvez les surcharges au niveau de la tâche dans le code ou la config de tâche après succès.
  • Ajoutez un test unitaire ou d'intégration qui attraperait la régression la prochaine fois.

Conclusion

La plupart des incidents Spark relèvent de motifs reconnaissables que vous pouvez diagnostiquer rapidement avec un inventaire d'environnement précis, les bons journaux et une reproduction minimale. Apportez d'abord de petits changements réversibles, vérifiez le résultat sur un pilote étroit, puis promouvez les corrections avec prudence. Les tableaux et le runbook ici vous donnent un chemin reproductible pour isoler la pression mémoire, l'instabilité du shuffle, les erreurs de sérialisation, les lacunes de dépendances et la dérive d'environnement, et pour récupérer en toute sécurité sans dommages collatéraux.

Recherches connexes

Score de qualité de l’article

Utilité pour le lecteur 100%
  • check_circle Guide prêt à lire
  • check_circle Exemples pratiques inclus
  • check_circle URL d’article optimisée pour le SEO