1 - Architecture d’exécution et modèle interne d’Apache Spark
- Rôles et interactions entre le driver, les executors et le cluster manager
- Modèle d’exécution Spark : jobs, stages, tasks et mécanismes de shuffle
- Fonctionnement du DAG, suivi du lignage des données et évaluation différée avec la lazy evaluation
- Optimisation des requêtes avec Catalyst Optimizer : plan logique, plan logique optimisé et plan physique
- Fonctionnement de Tungsten : sérialisation binaire, génération de code et gestion de la mémoire off-heap
- Principes et bénéfices de l’Adaptive Query Execution (AQE)
- Architecture client-serveur de Spark Connect
- Évolutions introduites depuis Apache Spark 3.5 : améliorations de l’AQE, extension de la prise en charge d’ANSI SQL et optimisation du support Python
AtelierSoumettre un job Spark et analyser son DAG dans Spark UI
Afficher les plans logique et physique avec explain() et explain("formatted")
Identifier les opérations narrow et wide et localiser les shuffles
Activer et désactiver l’AQE afin d’observer son incidence sur le plan d’exécution
Tester une connexion avec Spark Connect au moyen d’une Remote Spark Session et comparer son fonctionnement avec le mode local classique
2 - RDD avancés : partitionnement et persistance
- Identifier les cas d’usage pour lesquels les RDD sont plus adaptés que les DataFrames
- Comparer les transformations narrow et wide et comprendre leur incidence sur le shuffle
- Utiliser persist et cache avec les différents niveaux de stockage : MEMORY_ONLY, MEMORY_AND_DISK, MEMORY_AND_DISK_DESER, DISK_ONLY et OFF_HEAP
- Comparer la sérialisation Java et Kryo et mesurer leur impact sur les performances
- Gérer manuellement le partitionnement avec repartition, coalesce et partitionBy
- Créer des partitionneurs personnalisés pour les RDD clé-valeur
- Utiliser les variables broadcast et les accumulateurs
AtelierMesurer l’impact de coalesce et de repartition sur un job de transformation
Mettre en place un partitionneur personnalisé sur un RDD clé-valeur
Comparer les performances de persist() avec plusieurs niveaux de stockage
Utiliser une variable broadcast pour optimiser une jointure avec une table de petite taille
3 - DataFrames, Spark SQL et optimisation avec Catalyst
- Utilisation avancée de l’API DataFrame : fonctions de fenêtrage, pivots et opérations complexes
- UDF, Pandas UDF avec Apache Arrow et agrégations personnalisées : performances et bonnes pratiques
- Utilisation de Spark SQL et des expressions de table communes CTE : intérêts et limites
- Optimisation des lectures avec le predicate pushdown et le projection pruning
- Partitionnement et bucketing des tables Hive, Apache Iceberg et Delta Lake
- Optimisation dynamique avec Adaptive Query Execution : dynamic coalesce, dynamic join et gestion du data skew
- Présentation de Pandas API on Spark, anciennement Koalas, pour les data scientists
- Fonctionnement du Cost-Based Optimizer (CBO) et collecte des statistiques
AtelierCréer des requêtes avec des window functions et comparer leur implémentation avec l’API DataFrame et Spark SQL
Comparer les performances d’une UDF Python classique et d’une Pandas UDF sur un même calcul
Configurer le bucketing d’une table et mesurer son incidence sur les performances d’une jointure
Activer le CBO et exécuter ANALYZE TABLE afin d’observer son impact sur les plans d’exécution
4 - Shuffle et stratégies de jointure dans Spark
- Fonctionnement du shuffle : map output, fetch, tri et spill to disk
- Coût du shuffle et analyse des métriques Shuffle Read et Shuffle Write
- Principales stratégies de jointure : Broadcast Hash Join, Sort-Merge Join, Shuffle Hash Join et Broadcast Nested Loop Join
- Sélection automatique de la stratégie de jointure par Catalyst et configuration des seuils
- Utilisation des hints de jointure : BROADCAST, MERGE, SHUFFLE_HASH et SHUFFLE_REPLICATE_NL
- Détection et traitement du data skew provoqué par des clés très fréquentes
- Techniques anti-skew : salting, partition split et AQE Skew Join
- Utilisation du bucketing comme optimisation structurelle des jointures
AtelierComparer Broadcast Hash Join et Sort-Merge Join sur un jeu de données
Imposer une stratégie de jointure avec un hint et mesurer le gain de performance
Identifier une jointure affectée par le data skew dans Spark UI
Appliquer la technique du salting afin de corriger le déséquilibre des données
Activer la gestion automatique du skew avec l’AQE et comparer les résultats obtenus
5 - Gestion des ressources et parallélisation des traitements Spark
- Architecture mémoire d’un executor : execution memory, storage memory et user memory
- Utilisation de la mémoire off-heap et de Tungsten : avantages et paramètres de configuration
- Configuration des executors : num-executors, executor-cores, executor-memory et memoryOverhead
- Gestion du parallélisme avec spark.sql.shuffle.partitions et spark.default.parallelism
- Activation de Dynamic Allocation et réglage de ses paramètres
- Présentation des principaux gestionnaires de clusters : Spark sur YARN, Kubernetes et Spark Standalone
- Soumission des jobs avec spark-submit et comparaison des modes cluster et client
- Stratégies de dimensionnement des ressources selon le volume et la complexité des traitements
AtelierCalculer les paramètres optimaux des executors pour un cluster donné
Reproduire et résoudre une erreur OutOfMemoryError du côté d’un executor
Activer Dynamic Allocation et observer la variation du nombre d’executors
Comparer les modes cluster et client sur une même application Spark
6 - Spark Structured Streaming en production
- Modèle de Structured Streaming : table infinie et traitement incrémental
- Modes de déclenchement : micro-batch, continuous processing et available-now
- Sources de données streaming : Kafka, fichiers Parquet, JSON ou CSV, Delta Lake, Apache Iceberg et socket
- Sinks disponibles : Kafka, fichiers, foreach, foreachBatch, console et mémoire
- Gestion des données tardives avec les watermarks
- Création de fenêtres temporelles : tumbling windows, sliding windows et session windows
- Opérations avec état : agrégations, déduplication, mapGroupsWithState et flatMapGroupsWithState
- Présentation de TransformWithState et des State Stores RocksDB et HDFSBackedStateStore
- Checkpointing et garanties de traitement exactly-once
- Monitoring et analyse des métriques de streaming
7 - Pipeline Structured Streaming avec watermark et gestion d’état
- Lire et analyser un flux d’événements provenant de Kafka
- Appliquer une agrégation par fenêtre temporelle avec un watermark
- Gérer la déduplication des événements avec dropDuplicates
- Écrire les résultats dans Delta Lake ou Parquet avec une garantie exactly-once
- Inspecter les métriques du flux de données dans Spark UI
8 - Débogage et profiling des applications Spark
- Utilisation des onglets Jobs, Stages, Storage, Executors, SQL et Streaming de Spark UI
- Analyse post-mortem des applications avec Spark History Server
- Identification des principaux goulets d’étranglement : tâches lentes, data skew, shuffle excessif et pression du garbage collector
- Diagnostic des erreurs OutOfMemoryError côté driver ou executor et analyse des messages associés
- Réglage de la JVM : choix du Garbage Collector, notamment G1GC ou ZGC, et configuration des paramètres clés
- Présentation des outils de profiling SparkMeasure, SparkLens et Dr. Elephant
- Configuration du logging applicatif et accès aux logs des executors
- Intégration avec Prometheus, Grafana et OpenTelemetry pour assurer un monitoring continu
AtelierNaviguer dans Spark UI afin d’identifier le stage le plus coûteux d’un job
Reproduire et résoudre un incident lié à un garbage collection excessif et à des pauses prolongées
Analyser avec SparkLens un rapport présentant les recommandations d’optimisation
Utiliser Spark History Server pour comparer deux exécutions d’un même job
Mettre en place une supervision d’un cluster Spark avec Prometheus et Grafana