Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
Concepts fondamentaux sous-jacents du dimensionnement, de l’optimisation et de la résolution des problèmes. Lisez ceci en premier si vous débutez avec Spark sur Fabric.
Générales Récommandations et Choses à éviter
Scénario : Vous débutez avec Spark. Que sont les Dos et Don’ts
| Cas d’utilisation | Meilleures pratiques |
|---|---|
| Utiliser des formats sérialisés optimisés | Faire : préférer des formats comme Avro, Parquet ou Optimized Row Columnar (ORC), car ils incorporent un schéma, sont compacts et optimisent le stockage et le traitement. Dans Fabric, utilisez le format Delta pour garantir l'atomicité, la cohérence, l'isolation et la durabilité (ACID), ainsi que pour bénéficier des avantages en termes de performance. |
| Soyez prudent avec XML/JSON | Ne vous fiez pas à l’inférence de schéma pour les fichiers JSON (JavaScript Object Notation) ou XML (Extensible Markup Language), car Spark lit l’ensemble du jeu de données pour déduire le schéma, ce qui ralentit le traitement et consomme beaucoup de mémoire. Fournissez un schéma principal statique lors de la lecture de JSON/XML ou utilisez-le .option("samplingRatio", 0.1) pour accélérer les lectures, mais sachez que si l’exemple ne représente pas le jeu de données complet, les lectures peuvent échouer. Une approche plus sûre déduit le schéma d’un exemple représentatif et la conserve pour toutes les lectures.Évitez d’analyser les fichiers XML volumineux. L’analyse XML s’exécute de façon inhérente plus lente en raison du traitement des balises et du cast de type. |
| Optimiser les jointures et le filtrage | Action : appliquez l’élagage des colonnes et le filtrage au niveau des lignes avant les jointures pour réduire l’utilisation du shuffle et la consommation de mémoire. L’optimiseur Catalyst assure la gestion automatique du pushdown de prédicat lorsque vous utilisez les API DataFrame. Évitez les API RDD (Resilient Distributed Dataset), car elles contournent les optimisations de Catalyst. |
| Préférer les DataFrames aux RDD | Utilisez des DataFrames plutôt que des RDD pour la plupart des opérations. Les DataFrames utilisent l’optimiseur Catalyst et le moteur d’exécution Tungsten pour une exécution efficace. |
| Activer l’exécution de requêtes adaptatives (AQE) | Faites : activez AQE pour optimiser dynamiquement les partitions de shuffle et gérer automatiquement les données asymétriques. |
Gestion de la mémoire de l’exécuteur
Scénario : vous souhaitez comprendre la gestion de la mémoire de l’exécuteur pour le réglage des performances.
Même si un exécuteur est configuré avec une mémoire de 56 Go, Spark n’autorise pas tout cela à être utilisé directement pour les données utilisateur. Spark Core divise et gère la mémoire de l’exécuteur :
Mémoire réservée : Partie fixe réservée au système et à la surcharge interne Spark (par exemple, machine virtuelle Java (JVM), interne).
Mémoire utilisateur : Stocke les fonctions définies par l’utilisateur (UDF), les variables locales, les structures de données (listes, cartes, dictionnaires) et les objets créés pendant le calcul.
Mémoire de stockage : Contient les données mises en cache/persistantes, les variables diffusées et les données de permutation qui peuvent être mises en cache.
Mémoire d’exécution : Utilisé pour le calcul intermédiaire (shuffles, jointures, tris, agrégations).
Partage de mémoire dynamique : La limite entre stockage et mémoire d’exécution est mobile. Spark peut emprunter de la mémoire d’une région à l’autre, ce qui permet une utilisation flexible de la mémoire.
Débordement : Se produit lorsque la demande de mémoire, qu’elle soit pour le stockage ou l’exécution, dépasse la mémoire disponible après l’emprunt. Cela force les données sur le disque, ce qui peut affecter les performances.
Erreurs de mémoire insuffisante (OOM)
Scénario : les travaux Spark échouent avec des erreurs de mémoire insuffisante (OOM).
Pilote OOM :
Les erreurs OOM du contrôleur se produisent lorsque le contrôleur Spark dépasse sa mémoire allouée.
Cause courante : opérations lourdes de pilotes telles que collect(), countByKey()ou appels volumineux toPandas() qui extrayent trop de données dans la mémoire du pilote.
Atténuation : évitez les opérations lourdes dans la mesure du possible. Si cela est inévitable, augmentez la taille du pilote et effectuez un benchmark pour trouver la configuration optimale.
Exécuteur hors mémoire (OOM) :
Les erreurs OOM de l’exécuteur se produisent lorsqu’un exécuteur Spark dépasse sa mémoire allouée.
Cause courante : transformations nécessitant beaucoup de mémoire et de calcul sur des jeux de données volumineux (par exemple, des jointures larges, des agrégations, des shuffles) ou des jeux de données mis en cache/persistants qui dépassent la mémoire disponible de l’exécuteur (exécution + régions de stockage).
Atténuation : augmentez la mémoire de l’exécuteur si nécessaire, ajustez les paramètres de mémoire Spark (spark.memory.fraction, spark.memory.storageFraction) et persistez de manière sélective. Assurez-vous que les données mises en cache s’intègrent dans la mémoire disponible.
Asymétries des données
Symptômes de biais :
- Quelques tâches prennent plus de temps que d’autres dans l’interface utilisateur Spark (les tâches intermédiaires affichent une queue lourde).
- Écart important entre les durées médianes et maximales des tâches dans les métriques d’étape.
- Étapes avec de grandes tailles de lecture ou d'écriture de shuffle pour quelques partitions.
Causes courantes :
- Distribution inégale des données pour les clés de jointure/de groupe (clés chaudes).
- Partitionnement incorrect ou trop peu de partitions pour le volume de données.
- Anomalies de données en amont qui produisent des enregistrements volumineux ou de nombreuses clés null/vides.
Atténuation:
- Répartir ou coaliser pour augmenter le parallélisme des partitions et équilibrer les tailles.
- Appliquez le sel de clé ou le partitionnement personnalisé pour répartir les clés chaudes entre les partitions.
- Utilisez AQE (Exécution de requête adaptative) pour fusionner les partitions post-shuffle et activer les optimisations de jointure asymétrique.
- Utilisez des jointures de diffusion pour les petites tables de recherche afin d’éviter de mélanger entièrement.
- Conservez les jeux de données intermédiaires équilibrés avant les étapes coûteuses et réexécutez le travail.
Meilleures pratiques UDF
Scénario : vous devez appliquer une logique personnalisée qui ne peut pas être exprimée via des fonctions DataFrame intégrées.
Utilisez les API DataFrame Spark dans la mesure du possible. L’optimiseur Catalyst optimise les fonctions intégrées et les exécute en mode natif sur la machine virtuelle JVM, afin qu’elles offrent les meilleures performances.
Si vous devez utiliser une fonction UDF (fonction définie par l’utilisateur), évitez les UDF PySpark Python standard. Au lieu de cela, tenez compte des alternatives suivantes :
Fonctions Définies par l’Utilisateur Pandas (également appelées fonctions vectorisées définies par l’utilisateur) : Utilisez Apache Arrow pour assurer un transfert de données efficace entre JVM et Python. Les UDF Pandas permettent des opérations vectorisées, ce qui améliore considérablement les performances par rapport aux UDF Python ligne par ligne.
Fonctions définies par l’utilisateur Scala/Java : exécutez directement sur la machine virtuelle JVM, ce qui évite la surcharge de sérialisation Python. Les UDF Scala/Java sont généralement plus performantes que les UDF Python.
Soyez prudent avec les fonctions UDF (définies par l’utilisateur) Python. Chaque exécuteur lance un processus Python distinct, nécessitant la sérialisation et la désérialisation des données entre la machine virtuelle JVM et Python. Cela crée un goulot d’étranglement des performances, en particulier à grande échelle.
Journalisation des erreurs
Scénario : Meilleures pratiques pour la journalisation des erreurs dans Fabric Spark
Utilisez
log4jplutôt queprint()ce qui charge fortement le conducteur. Aveclog4j, vous pouvez accéder aux journaux d’activité des pilotes et les rechercher (à l’aide du nom de l’enregistreur d’événements, par exemple : PySparkLogger).Encadrez les lectures, les écritures et les transformations dans les blocs try et except. Utiliser
logger.errorpour les exceptions etlogger.infopour les messages de progression.Journalisation Python : Idéal pour les opérations de journalisation, les mises à jour d’état ou le débogage d’informations à partir du code qui s’exécute uniquement sur le pilote Spark. Le module de journalisation de Python ne se propage pas aux journaux d’exécution. Consultez la documentation sur le développement, l’exécution et la gestion des notebooks.
Spark log4j : Standard pour la journalisation d’applications robustes au niveau de la production dans Spark, car elle s’intègre en mode natif avec les journaux de pilote/exécuteur spark.
Exemple d’utilisation de log4j dans PySpark :
import traceback # Get log4j logger log4jLogger = spark._jvm.org.apache.log4j logger = log4jLogger.LogManager.getLogger("PySparkLogger") logger.info("Application started.") try: # Create DataFrame with 20 records data = [(f"Name{i}", i) for i in range(1, 21)] # 20 records df = spark.createDataFrame(data, ["name", "age"]) logger.info("DataFrame created successfully with 20 records.") df.show(s) # 's' is not defined -> will throw error but the application will not fail except Exception as e: logger.error(f"Error while creating or showing DataFrame: {str(e)}\n{traceback.format_exc()}")Centraliser la surveillance des erreurs :
Utilisez l’extension d’émetteur de diagnostic (Surveiller les applications Apache Spark avec Azure Log Analytics) dans l’environnement et joignez-vous aux notebooks exécutant des applications Spark. L’émetteur peut envoyer des journaux d’événements, des journaux personnalisés (comme log4j) et des métriques à Azure Log Analytics/Stockage Azure/Azure Event Hubs. Passez le nom log4j à la propriété :
spark.synapse.diagnostic.emitter.\<destination\>.filter.loggerName.match.En outre, pour le débogage, vous pouvez également collecter des lignes/enregistrements ayant échoué dans des tables Lakehouse (LH) pour la capture des données incorrectes au niveau des enregistrements.