Introduction
Apache Spark peut traiter d'énormes volumes de données, mais la fiabilité en production repose sur la capacité à observer les bons signaux et à agir vite. Ce guide explique exactement quoi surveiller, comment le collecter, l'alerter et le vérifier, avec des extraits de configuration concrets, des règles d'alerte pragmatiques, des résultats attendus, des modes de panne courants, des étapes de retour arrière et une checklist d'exploitation hebdomadaire. Pour le référencement et la clarté, nous mentionnons les expressions courantes du domaine (Apache Spark monitoring, Apache Spark alerts, Apache Spark metrics, Apache Spark dashboard, Apache Spark incident response) là où elles s'intègrent naturellement.
Utilisez ce guide pour :
- Lancer un pilote ciblé et à faible risque avant un déploiement généralisé.
- Mettre en évidence la santé du driver et des exécutors, la pression mémoire, les performances de shuffle, les goulets d'étranglement SQL/Streaming, et les frontières d'intégration (Kafka, HDFS, déclencheurs Airflow/NiFi).
- Construire des alertes et des tableaux de bord actionnables avec des étapes de vérification et de remédiation claires.
Inventaire des versions et de l'environnement
Commencez par consigner précisément ce que vous exécutez. Cela évite les configurations incompatibles et offre un point de retour arrière propre.
À enregistrer :
- Distribution et version de Spark (ex. : Spark 3.3.x ou 3.4.x)
- Version de la JVM (ex. : OpenJDK 11 ou 17)
- Gestionnaire de cluster (Standalone, YARN, ou Mesos)
- Mix de charges (batch, Structured Streaming, Spark SQL)
- Allocation dynamique activée/désactivée, tailles et nombres typiques d'exécuteurs
- Sources/cibles de données couvertes (ex. : Kafka, HDFS, stockage objet)
- Chaîne de métriques (ex. : agent JMX exporter, Prometheus, outil de dashboard)
Attendus de base :
- Spark expose des métriques via son système de métriques et la JMX JVM. Attachez un sink de métriques ou un exporter JMX au driver et aux exécuteurs.
- Centralisez les logs (driver/executor) via un agent ou les journaux du gestionnaire de cluster. Utilisez des motifs de logs précis pour déclencher des incidents et faciliter le diagnostic.
Chemin de configuration sûre
Déployez par incréments, via un pilote étroit et inspectable localement :
- Choisissez une application Spark non critique qui représente votre charge courante (ex. : ETL quotidien avec jointures et shuffle).
- Instrumentez les JVM du driver et des exécuteurs avec un JMX exporter ou activez un sink de métriques Spark éprouvé.
- Collectez un minimum utile : mémoire JVM du driver, mémoire/GC des exécuteurs, tâches actives, tâches échouées, durée de stage, shuffle read/write, durée des micro-batchs en streaming (si applicable).
- Construisez un tableau de bord pilote et au plus trois alertes. Gardez un périmètre étroit et mesurable.
- Vérifiez en induisant des conditions de test à faible risque (cf. section Vérification) et confirmez la bonne émission/résolution des alertes.
- Documentez le retour arrière (suppression des flags d'agent ou du sink) et les configs validées.
- Étendez à deux apps critiques supplémentaires après 1-2 semaines de stabilité.
Que surveiller et alerter
Spark fournit plusieurs couches de signaux. Commencez par l'essentiel, puis ajoutez de la profondeur sélectivement.
Composants clés et signaux
| Composant | Métrique/Signal | Interprétation | Première action |
|---|---|---|---|
| Driver JVM | Heap utilisé %, temps GC, nombre de threads | Pression mémoire et pauses | Capturer la tendance de heap/GC, ajuster la taille du driver |
| Exécuteurs | Tâches actives, taux d'échec, compte d'exécuteurs perdus | Santé du parallélisme et stabilité | Vérifier santé des nœuds, logs, localité des données |
| Stages/Jobs | Durée, skew (max vs p50), échecs | Points chauds et retries | Inspecter partitions/jointures déséquilibrées |
| Shuffle | Read/write MB, fetch wait, erreurs de transfert | Pression réseau/disque | Valider shuffle service et I/O disque |
| SQL | Durée de requête, taille de broadcast, métriques de spill | Qualité du plan et adéquation mémoire | Revoir le plan, seuils de broadcast |
| Streaming | Taux d'entrée/traitement, durée batch, taille state store | Backpressure et croissance d'état | Régler trigger, TTL d'état, santé du checkpoint |
| Intégrations | Kafka lag, erreurs HDFS/FS | Santé amont/aval | Coordonner avec les propriétaires plateforme |
Règles d'alerte pratiques (démarrage)
| Signal et condition | Rationale | Sévérité | Première action |
|---|---|---|---|
| Heap driver > 85% sur 5 min OU temps GC > 20% sur 5 min | Risque OOM/pauses longues | Haute | Capturer métriques, envisager +mémoire ou correction de plan |
| Taux d'échec des tâches > 2% sur 10 min OU hausse des exécuteurs perdus | Instabilité données/nœuds | Haute | Inspecter logs tâches, santé nœuds, politique de retry |
| p95 tâche de stage à 3× p50 sur 10 min | Skew/point chaud | Moyenne | Vérifier partitionnement, jointures skew |
| p95 fetch wait shuffle > 2 s sur 10 min | Problème réseau/shuffle service | Moyenne | Valider réseau, I/O shuffle |
| Durée batch streaming > intervalle de trigger sur 3 cycles | Backpressure | Haute | Réduire travail/batch, ajuster ressources |
Notes :
- Ce sont des exemples construits. Ajustez aux normes de vos charges après la première semaine.
- Utilisez des métriques en taux/ratio pour stabiliser selon la taille variable des traitements.
Signaux de logs utiles
Complétez les métriques par des motifs de logs pour accélérer le triage :
- OutOfMemoryError ou GC overhead limit exceeded
- ExecutorLostFailure ou FetchFailedException
- TaskKilled dû à la spéculation ou à des pics de préemption
- Avertissements StateStore pour le streaming
- Timeouts et échecs d'autorisation clients HDFS/Kafka
Convertissez les motifs répétés en alertes silencieuses ou liens de runbook plutôt qu'en pages immédiates.
Mise en œuvre : collecte des métriques et des logs
Deux approches principales pour extraire des métriques numériques du driver et des exécuteurs : le système de métriques Spark (sinks) et la JMX JVM. La voie JMX est largement compatible et simple à vérifier.
Option A : JMX exporter en agent Java (driver et exécuteurs)
- Obtenez un jar de JMX exporter et choisissez des ports (ex. : 7071 pour le driver, 7072 pour les exécuteurs).
- YAML minimal pour JMX exporter (exemple construit) :
rules:
- pattern: ".*"
name: jmx_$0
type: GAUGE
labels: {}
- Ajoutez l'agent au driver et aux exécuteurs. Exemple de flags spark-submit (chemins et ports construits) :
spark-submit \
--class com.example.YourJob \
--conf "spark.driver.extraJavaOptions=-javaagent:/opt/jmx/jmx_exporter.jar=7071:/opt/jmx/jmx.yaml" \
--conf "spark.executor.extraJavaOptions=-javaagent:/opt/jmx/jmx_exporter.jar=7072:/opt/jmx/jmx.yaml" \
--conf "spark.executor.instances=4" \
your-job-assembly.jar
- Vérification locale :
- Vérifiez que l'endpoint HTTP JMX du driver répond :
curl http://<driver-host>:7071/metricsrenvoie du texte. - Pendant un run, confirmez qu'au moins un endpoint JMX d'exécuteur répond sur son port.
- Scrapez ces endpoints dans votre magasin de métriques et construisez un tableau de bord pilote avec les signaux listés plus haut.
Atouts : changements Spark minimaux, couvre JVM et sous-systèmes via MBeans. Contraintes : ports à ouvrir et uniques par processus.
Option B : Système de métriques Spark (exemple de sink)
Configurez spark.metrics.conf pour envoyer vers un sink (ex. CSV ou de type Graphite). Exemple construit :
*.sink.csv.class=org.apache.spark.metrics.sink.CsvSink
*.sink.csv.period=10
*.sink.csv.unit=seconds
*.sink.csv.directory=/tmp/spark-metrics
master.source.jvm.class=org.apache.spark.metrics.source.JvmSource
worker.source.jvm.class=org.apache.spark.metrics.source.JvmSource
executor.source.jvm.class=org.apache.spark.metrics.source.JvmSource
driver.source.jvm.class=org.apache.spark.metrics.source.JvmSource
Soumettez avec :
spark-submit \
--conf spark.metrics.conf=/path/to/spark.metrics.conf \
your-job-assembly.jar
Vérification :
- Confirmez la création de fichiers sous
/tmp/spark-metricssur les hôtes driver et exécuteurs. - Inspectez les jauges JVM et les métriques de tâches pour valider l'activité.
Atouts : pas de ports supplémentaires. Contraintes : prévoir un forwarder vers le magasin central de métriques.
Logs
- Assurez une rétention centralisée des logs driver/executor.
- Ajoutez des parseurs pour exceptions courantes et tags Spark (executor ID, stage ID, job ID).
- Indexez par application ID et numéro d'attempt pour distinguer les retries.
Tableaux de bord qui guident l'action
Démarrez avec un tableau de bord focalisé par type de charge. Évitez les pages tentaculaires. Regroupez par action opérateur.
Sections recommandées :
- Santé du driver : heap utilisé %, temps GC %, threads, CPU JVM.
- Exécuteurs : total vs actifs, exécuteurs perdus, tâches actives, échecs, p95 temps de tâche.
- Stages et shuffle : durée p50/p95, skew (p95/p50), MB/s lecture/écriture shuffle, p95 fetch wait.
- SQL/ETL : enregistrements en entrée/sortie, métriques de spill, taille de broadcast, top-N requêtes par durée.
- Streaming (si applicable) : taux d'entrée vs traitement, durée de batch, backlog/lag, taille state store, latence d'écriture du checkpoint.
- Bords : Kafka lag ou erreurs d'écriture HDFS (panneaux construits si vos métriques les incluent).
Gardez la vue lisible : 12-16 panneaux maximum sur la page principale ; liez des approfondissements.
Vérification et diagnostics
Avant le déploiement large, vérifiez que vos signaux sont fiables et non bruyants.
- Accessibilité des endpoints
- Driver : curl l'endpoint JMX/métriques et attendez une réponse non vide.
- Exécuteur : au moins un endpoint durant un run.
- Si sink, validez la création de fichiers ou l'émission réseau au rythme attendu.
- Forme et labels des métriques
- Confirmez la présence de application_id, executor_id, stage_id, job_id dans les labels ou les noms.
- Surveillez la cardinalité : pas de labels par partition/tâche.
- Dry-runs d'alertes (exemples construits)
- Risque mémoire : chargez davantage pour pousser le heap driver > 85 % temporairement ; vérifiez déclenchement et résolution.
- Échecs de tâches : injectez un jeu d'entrées erronées contrôlé ; l'alerte basée sur un taux doit se déclencher sans spammer.
- Pression shuffle : exécutez une large jointure pour augmenter read/write shuffle ; observez la distribution fetch wait.
- Backpressure streaming : réduisez temporairement les ressources pour dépasser l'intervalle de trigger ; validez l'alerte.
- Recoupements
- Comparez les tendances du dashboard avec l'UI Spark sur la même attempt.
- Vérifiez que durées de stage et comptes de tâches concordent dans des marges raisonnables.
- Répétition de runbook
- Pour chaque alerte, suivez les premières actions et ajustez le runbook avec hôtes, chemins et commandes réellement utilisés.
Résultats attendus :
- Les métriques arrivent en 10-30 s après un changement.
- Les alertes se déclenchent dans une fenêtre d'évaluation, se résolvent proprement, sans battement.
- Les opérateurs corrèlent alerte, UI Spark et logs en moins de 2 minutes.
Modes de panne et reprise
Anticipez ces problèmes courants et préparez les remèdes.
- Exporter non attaché ou mauvais port
- Symptôme : pas de métriques driver/exécuteurs ; erreurs de scrape.
- Correctif : vérifier flags extraJavaOptions, chemin du jar et ports ; redémarrer le job.
- Retour arrière : retirer extraJavaOptions et redéployer.
- Labels à forte cardinalité
- Symptôme : magasin saturé, dashboards lents.
- Correctif : relabeller/agréger au niveau stage/job/app.
- Retour arrière : revenir au mapping/règles précédent(e)s.
- Churn d'allocation dynamique
- Symptôme : exécuteurs éphémères, séries manquantes, alertes bruyantes.
- Correctif : agréger au niveau application ; définir un minimum d'exécuteurs pour stabiliser durant le pilote.
- Retour arrière : désactiver les alertes concernées jusqu'à stabilisation.
- Trous d'acheminement de logs
- Symptôme : logs manquants pendant un incident.
- Correctif : vérifier permissions de l'agent et rotation ; augmenter la rétention durant le déploiement.
- Retour arrière : s'appuyer sur les logs du gestionnaire de cluster.
- Blocages réseau/pare-feu
- Symptôme : le scraper n'atteint pas les ports JMX.
- Correctif : ouvrir les ports nécessaires du scraper vers les nœuds driver/exécuteurs ; envisager des exporters au niveau nœud.
- Retour arrière : basculer sur le sink de métriques Spark utilisant des sorties existantes.
- Fatigue d'alertes
- Symptôme : trop d'alertes « Moyenne », ignorées.
- Correctif : déclasser en ticket/email ; resserrer conditions et durées.
- Retour arrière : revenir au set d'alertes initial réduit.
Checklist de reprise après un changement de monitoring :
- Si la nouvelle config perturbe, redéployez immédiatement la dernière config valide.
- Confirmez que les métriques et alertes précédentes fonctionnent comme avant.
- Enregistrez le changement avec horodatage, configs et raison du rollback.
Runbook d'actions pratiques
Quand une alerte tombe, agissez avec des étapes concrètes.
Pression mémoire du driver :
- Capturer heap courant, % temps GC et jobs/stages récents.
- Ouvrir l'onglet Storage de l'UI Spark ; évincer les datasets mis en cache si possible.
- Envisager +mémoire driver au prochain run ou refactorer des collect() volumineux.
Échecs de tâches ou perte d'exécuteur :
- Inspecter les logs d'exécuteurs (disque, réseau, OOM).
- Valider la cohérence des données d'entrée pour la fenêtre traitée.
- Si un sous-ensemble de partitions échoue, relancer le stage après correction des entrées.
Goulets de shuffle :
- Vérifier la santé du shuffle service et l'I/O disque des nœuds affectés.
- Augmenter le nombre de partitions de shuffle ou activer l'adaptive query execution pour une meilleure granularité.
Backpressure en streaming :
- Comparer taux d'entrée vs traitement ; ajuster l'intervalle de trigger.
- Régler la compaction/TTL du state store ; vérifier la latence et la santé du stockage de checkpoint.
Checklist d'exploitation
À exécuter chaque semaine et lors des déploiements.
- Inventaire
- Vérifier que les versions Spark/JVM/outils de métriques sont stables ou documentées.
- Confirmer les flags d'agents driver/executor ou la conf du sink en contrôle de version.
- Smoke checks
- Lancer un petit job et vérifier l'arrivée des métriques et logs en < 30 s.
- Vérifier qu'au moins un endpoint métrique d'exécuteur répond en cours d'exécution.
- Dashboards
- Confirmer l'affichage des panneaux driver/exécuteurs/stages/shuffle.
- Revoir le top-N des jobs/requêtes les plus lents.
- Alertes
- Passer en revue la semaine écoulée ; ajuster seuils/durations pour réduire le bruit.
- Déclencher manuellement une alerte de test sûre par mois (pic de charge court et contrôlé).
- Signaux de capacité
- Inspecter les tendances du heap driver et du temps GC sur 7 jours.
- Examiner les comptes et causes d'échecs d'exécuteurs.
- Bords de données
- Surveiller Kafka lag et erreurs HDFS pour anomalies.
- Runbooks
- Mettre à jour ports/chemins/outils si nécessaire.
- Vérifier l'accès et les permissions de l'astreinte.
- Sauvegarde et retour arrière
- Archiver configs actuelles et dernière version stable.
- Confirmer la suppression des flags d'agent et le redéploiement en < 10 min.
Conclusion
Vous disposez désormais d'un chemin pragmatique et à faible risque pour implémenter la surveillance et les alertes Spark :
- Démarrez par un pilote étroit et un set de métriques à forte valeur.
- Attachez un JMX exporter ou configurez les sinks de métriques Spark côté driver et exécuteurs.
- Construisez un tableau de bord ciblé et un set compact d'alertes actionnables.
- Vérifiez endpoints, labels de métriques et comportement des alertes via des tests sûrs avant d'étendre au cluster.
- Anticipez les modes de panne et gardez prêts les pas de retour arrière.
Après 1-2 semaines de fiabilité du pilote, déployez sur davantage d'applications et adaptez les seuils par charge. Gardez des changements petits, mesurables et faciles à inspecter localement ; cela réduit le rework et accélère une réponse aux incidents plus sûre et plus efficace.