Apache Spark

30 questions en accès libre

Partitions, shuffle, RDD contre DataFrame, cache, joins et UDF : les questions posées en entretien de data engineering, avec la réponse qu’attend l’intervieweur.

Niveau

Cliquez sur une question pour dérouler la réponse attendue.

  1. 01Qu’est-ce qu’un RDD ?Juniorrddfondamentaux

    Un RDD (Resilient Distributed Dataset) est une collection d’objets répartie en partitions sur plusieurs machines, immuable, et reconstructible en cas de panne.

    Les trois mots du sigle portent la réponse complète :

    • Distributed : les données sont découpées en partitions traitées en parallèle.
    • Resilient : Spark mémorise la suite d’opérations qui a produit le RDD — sa lignée — donc il peut recalculer une partition perdue sans relancer le job.
    • Dataset : c’est une collection d’objets JVM, sans schéma de colonnes.

    Ce qu’il faut ajouter pour montrer qu’on a compris la limite : un RDD n’a pas de schéma, donc Catalyst ne peut pas l’optimiser. C’est pourquoi on écrit du DataFrame en pratique, et du RDD seulement pour du contrôle fin de partition ou des données sans structure tabulaire.

  2. 02Quelle différence entre transformation et action ?Juniorlazy-evaluationfondamentaux

    Une transformation décrit un calcul sans l’exécuter ; une action déclenche l’exécution. C’est ce qui permet à Spark de voir toute la chaîne avant de la lancer, donc de l’optimiser.

    Lire la réponse détaillée
  3. 03À quoi servent sc et spark dans un shell Spark ?Juniorsparksessionfondamentaux

    sc est le SparkContext : le point d’entrée de l’API RDD (parallelize, textFile, variables diffusées, accumulateurs, dossier de checkpoint).

    spark est la SparkSession : le point d’entrée unifié depuis Spark 2.0, qui couvre les DataFrames et le SQL (read, sql, createDataFrame, conf, catalog).

    La relation compte autant que la définition : la session contient le contexte, récupérable par spark.sparkContext. Les deux sont créés automatiquement par un spark-shell, mais dans une application soumise avec spark-submit, c’est à vous de construire la session.

    scala
    val spark = SparkSession.builder().appName("job").getOrCreate()
    val sc = spark.sparkContext

    Deux points qui font bonne impression : il n’existe qu’un seul SparkContext par JVM — d’où getOrCreate() plutôt que new — et si un tutoriel utilise sqlContext ou HiveContext, il précède Spark 2.0, où ces deux objets ont été fusionnés dans la session.

  4. 04Quelle différence entre RDD, DataFrame et Dataset ?Intermédiairerdddataframedataset

    Les trois sont immuables, distribués et tolérants aux pannes. Ce qui les sépare tient en deux axes : le schéma et le typage.

    RDDDataFrameDataset[T]
    Schéma de colonnesnonouioui
    Types vérifiés à la compilationouinonoui
    Optimiseur Catalystnonouioui
    Disponible en Pythonouiouinon

    La phrase qui résume : le RDD est une collection distribuée bas niveau — vous décrivez comment traiter ; le DataFrame est une table distribuée haut niveau — vous décrivez ce que vous voulez, et Spark choisit comment.

    C’est précisément parce qu’on cède le « comment » que Catalyst peut réécrire le plan, ce qui rend le DataFrame nettement plus rapide dans la quasi-totalité des cas.

    Deux précisions qui montrent la maîtrise : DataFrame est Dataset[Row], donc passer de l’un à l’autre par .as[T] ne coûte rien ; et Dataset[T] n’existe pas en PySpark, puisque le typage à la compilation suppose un compilateur.

  5. 05Qu’est-ce qu’une partition dans Spark ?Juniorpartitionsfondamentaux

    Une partition est un morceau des données traité par une seule tâche, sur un seul cœur, de façon séquentielle. C’est l’unité de parallélisme de Spark : le nombre de partitions détermine combien de tâches peuvent tourner en même temps.

    La formulation qui montre qu’on a compris : une partition est l’unité de travail, une tâche est son exécution. Un RDD de 200 partitions produit 200 tâches par étape.

    Deux réglages à citer, parce que ce sont ceux qu’on ajuste en pratique :

    • sc.defaultParallelism fixe le nombre de partitions initiales, dérivé du nombre de cœurs disponibles.
    • spark.sql.shuffle.partitions (200 par défaut) fixe le nombre de partitions après un shuffle — c’est presque toujours cette valeur qu’il faut changer.

    La règle de dimensionnement usuelle : viser 100 à 200 Mo de données par partition, et un nombre de partitions égal à deux ou trois fois le nombre de cœurs. Trop peu de partitions laisse des cœurs inactifs ; trop de partitions noie le job dans le coût d’ordonnancement des tâches.

  6. 06Qu’est-ce qu’un shuffle et pourquoi est-ce coûteux ?Intermédiaireshufflepartitionsperformance

    La redistribution des données entre partitions, déclenchée par un groupBy, un join ou un repartition. Coûteuse parce qu’elle cumule écriture disque, transfert réseau et sérialisation.

    Lire la réponse détaillée
  7. 07Pourquoi préférer reduceByKey à groupByKey ?Intermédiaireshufflereducebykeyperformance

    Parce que reduceByKey agrège avant le shuffle et ne transfère que des résultats partiels, là où groupByKey déplace chaque valeur. Trois ordres de grandeur sur le volume réseau.

    Lire la réponse détaillée
  8. 08Quelle différence entre coalesce et repartition ?Intermédiairepartitionscoalescerepartition

    Les deux changent le nombre de partitions ; une seule fait un shuffle.

    repartition(n) shuffle et équilibre les partitions. Il peut augmenter ou réduire le nombre de partitions, et il coûte un transfert réseau complet.

    coalesce(n) fusionne des partitions voisines sans shuffle — c’est quasi gratuit. Deux contreparties : il ne peut que réduire (demander plus de partitions qu’il n’en existe est ignoré en silence), et il hérite du déséquilibre de départ.

    La réponse qui distingue un candidat qui a pratiqué : le bon usage de coalesce est juste avant l’écriture, pour éviter de cracher 200 petits fichiers.

    scala
    resultat.coalesce(10).write.parquet("s3a://bucket/sortie/")

    Le piège classique, à mentionner : coalesce placé en amont d’une grosse agrégation laisse les tailles inégales, et une seule tâche supporte tout le poids. Pire, comme coalesce ne shuffle pas, il remonte dans le plan et réduit aussi le parallélisme des étapes précédentes. Si vous avez besoin de moins de partitions mais de tâches équilibrées, c’est repartition qu’il faut, malgré son coût.

  9. 09Combien de partitions Spark crée-t-il par défaut ?Intermédiairepartitionsdefaultparallelismconfiguration

    Cela dépend de l’origine des données, et c’est le piège de la question : il y a deux réglages différents.

    À la création, sc.defaultParallelism décide. En mode local, il vaut le nombre de cœurs de la machine ; sur un cluster, le total des cœurs des exécuteurs, avec un plancher de 2.

    À la lecture d’un fichier, c’est la taille des blocs qui décide : Spark crée environ une partition par bloc de 128 Mo (spark.sql.files.maxPartitionBytes). Un fichier Parquet de 1 Go donne donc environ 8 partitions.

    Après un shuffle, c’est spark.sql.shuffle.partitions, qui vaut 200 par défaut — indépendamment du volume, de la taille du cluster et du bon sens. C’est la valeur qu’on ajuste le plus souvent.

    scala
    sc.defaultParallelism                                  // parallélisme de base
    spark.conf.get("spark.sql.shuffle.partitions")          // "200"
    df.rdd.getNumPartitions                                 // le compte réel

    Le détail à ajouter pour montrer l’expérience : spark.default.parallelism doit être posé avant la création du contexte pour être pris en compte, alors que spark.sql.shuffle.partitions se règle à chaud, requête par requête.

  10. 10Avec 100 partitions, ai-je besoin de 100 processeurs ?Seniorpartitionstachesexecution

    Non. Une partition est une tâche potentielle, pas un cœur réservé.

    Avec 100 partitions et 10 cœurs disponibles, Spark exécute 10 tâches à la fois et enchaîne 10 vagues successives. Rien n’échoue, rien n’attend anormalement : l’ordonnanceur distribue les tâches aux cœurs libres au fur et à mesure.

    C’est pour cela qu’avoir plus de partitions que de cœurs est recommandé, contrairement à l’intuition. Deux à trois fois le nombre de cœurs est la règle usuelle, pour deux raisons :

    • L’équilibrage de charge. Si une partition est plus lourde, les autres cœurs enchaînent d’autres tâches pendant ce temps au lieu de rester inactifs.
    • La reprise sur panne. Recalculer une petite partition perdue coûte moins qu’une énorme.

    Les deux cas à éviter, qu’il faut nommer :

    • Moins de partitions que de cœurs : des cœurs payés et inactifs.
    • Partitions énormes : une tâche traîne — le straggler — pendant que les 99 autres sont finies, et l’étape suivante attend la barrière du shuffle.

    Le détail de l’exécution par vagues est dans pourquoi 100 partitions ne font pas 100 processeurs.

  11. 11cache, persist ou checkpoint : lequel utiliser ?Intermédiairecachepersistcheckpoint

    cache est un persist au niveau par défaut, persist laisse choisir le stockage, et checkpoint écrit sur un stockage fiable en coupant la lignée. Seul le dernier règle les traitements itératifs.

    Lire la réponse détaillée
  12. 12Pourquoi une UDF est-elle plus lente qu’une fonction native ?Seniorudfcatalystperformance

    Parce qu’une UDF est une boîte noire pour Catalyst. L’optimiseur ne peut pas lire ce qu’elle fait, donc il perd trois leviers d’un coup :

    1. Le whole-stage codegen : l’étape ne peut plus être fusionnée avec ses voisines en une seule boucle compilée.
    2. Le predicate pushdown : un filtre exprimé dans une UDF ne peut pas être poussé vers Parquet ou la base source.
    3. Le format Tungsten : les lignes doivent être désérialisées en objets JVM pour appeler la fonction, puis resérialisées.

    Le troisième point est le plus coûteux dans les faits : un filtre qui aurait éliminé 99 % des données au niveau du fichier doit maintenant charger 100 % des lignes pour les évaluer une à une.

    En PySpark, un facteur s’ajoute : chaque ligne traverse la frontière JVM ↔ Python. C’est là que l’écart avec Scala devient réel, et qu’il faut mentionner les UDF vectorisées (pandas_udf), qui transfèrent par lots via Arrow et réduisent fortement le surcoût sans le supprimer.

    La conclusion à donner : avant d’écrire une UDF, chercher dans org.apache.spark.sql.functionswhen, regexp_extract, coalesce, split, date_format couvrent l’essentiel des cas où l’on est tenté d’en écrire une.

    Sur Databricks, la facture est double : une UDF fait aussi retomber l’étape sur Spark JVM au lieu du moteur Photon, que l’on paie pourtant.

  13. 13Qu’est-ce que Catalyst ?Intermédiairecatalystdataframeperformance

    Catalyst est l’optimiseur de requêtes de Spark SQL. Il transforme ce que vous écrivez en un plan d’exécution efficace, en quatre étapes :

    1. Plan logique non résolu — l’arbre de votre requête, colonnes non vérifiées.
    2. Analyse — résolution des noms de colonnes et de tables contre le catalogue. C’est ici qu’une faute comme select("agge") échoue.
    3. Optimisation logique — réécritures à base de règles : remontée des filtres au plus près de la source, élimination des colonnes inutiles, simplification des expressions constantes, fusion des projections.
    4. Planification physique — choix des algorithmes concrets, notamment la stratégie de join, en s’appuyant sur des statistiques de coût.

    Les deux optimisations à nommer, parce que ce sont celles qui changent tout en pratique : le predicate pushdown (le filtre descend jusqu’au fichier Parquet, qui saute des blocs entiers) et le column pruning (seules les colonnes demandées sont lues).

    Le point clé de la réponse : Catalyst ne fonctionne que sur des expressions qu’il peut lire. Un RDD ou une UDF est opaque pour lui, ce qui explique d’un coup pourquoi le DataFrame est plus rapide que le RDD et pourquoi une UDF annule le gain.

  14. 14Qu’est-ce qu’un broadcast join et quand l’utiliser ?Intermédiairejoinbroadcastshuffle

    Il copie la petite table sur chaque exécuteur pour joindre localement, sans shuffle de la grande. La stratégie la plus rapide quand la table tient en mémoire, avec un risque d’OOM sur le driver.

    Lire la réponse détaillée
  15. 15Comment gérez-vous un déséquilibre de données (data skew) ?Seniorskewshufflejoin

    Une ou quelques clés portent l’essentiel des lignes, donc une seule tâche traite presque tout. On mesure dans la Spark UI, puis on choisit entre AQE, isolation des sentinelles, broadcast ou salting.

    Lire la réponse détaillée
  16. 16Scala est-il plus rapide que PySpark ?Intermédiairescalapysparkperformance

    Ça dépend de ce que vous écrivez — et répondre « oui » sans nuance est une erreur fréquente.

    Sur l’API DataFrame, les performances sont équivalentes. Votre code Python ne fait que construire un plan logique, qui est envoyé à la JVM et exécuté par Catalyst. Le langage d’écriture disparaît avant l’exécution : deux jobs identiques en Scala et en PySpark produisent le même plan physique.

    Scala gagne réellement dans trois cas :

    1. Les UDF classiques, où chaque ligne traverse la frontière JVM ↔ Python. C’est l’écart le plus visible, et il se réduit fortement avec les UDF vectorisées (pandas_udf) qui transfèrent par lots via Arrow.
    2. L’API RDD, où les objets Python doivent être sérialisés et désérialisés en permanence.
    3. Le Dataset[T] typé, qui n’existe pas en Python faute de compilateur.

    La conclusion à donner, parce que c’est celle qui compte en équipe : le choix se fait sur l’écosystème et les compétences, pas sur la vitesse. PySpark donne accès à pandas, scikit-learn et aux bibliothèques de science des données ; Scala donne le typage à la compilation, un seul artefact JVM et un accès direct aux API internes. Écrivez du DataFrame dans les deux cas et la question de performance ne se pose plus.

  17. 17Pourquoi préférer Parquet à CSV ?Juniorparquetformatsperformance

    Quatre raisons, et la première est la plus importante : Parquet est orienté colonnes, le CSV est orienté lignes.

    1. Lecture sélective des colonnes. Une requête sur 3 colonnes d’une table de 80 ne lit que ces 3 colonnes. En CSV, il faut lire chaque ligne entière pour en extraire trois champs.
    2. Le schéma est dans le fichier. Le pied de fichier Parquet contient les noms et les types, donc aucune inférence n’est nécessaire — et aucun risque qu’un type change d’une exécution à l’autre.
    3. Compression bien plus efficace. Des valeurs d’un même type et souvent répétées sont stockées ensemble, ce qui permet l’encodage par dictionnaire et RLE. Un facteur 5 à 10 par rapport au CSV est courant.
    4. Statistiques par bloc. Chaque row group stocke le min et le max de chaque colonne, donc un filtre WHERE date > '2026-01-01' peut sauter des blocs entiers sans les décompresser — c’est le predicate pushdown.

    La formule à retenir et à énoncer : le CSV est un format d’échange, pas un format de stockage. On le convertit une fois à l’arrivée, avec un schéma explicite, puis on ne relit plus que du Parquet ou du Delta.

    Le cas où le CSV reste justifié, à mentionner pour ne pas paraître dogmatique : un échange avec un système qui ne lit rien d’autre, ou un fichier destiné à être ouvert à la main.

  18. 18Quel est le coût de l’option inferSchema ?Intermédiaireschemacsvperformance

    Deux coûts, et le second est celui qui compte.

    Le coût visible : Spark effectue une passe complète supplémentaire sur les données, uniquement pour déterminer les types. Sur 500 Go, vous payez deux balayages là où un seul était nécessaire — temps de lecture doublé et deux fois les requêtes GET facturées par le stockage objet.

    Le coût invisible : un schéma deviné dépend des données du jour. Une colonne d’entiers est typée Integer aujourd’hui ; demain, un système en amont écrit un seul N/A, et elle devient String. Votre code ne change pas, mais $"age" > 18 passe en comparaison lexicographique — et "9" > "18" renvoie true. Le job réussit et produit des chiffres faux.

    La bonne pratique à énoncer :

    scala
    spark.read
      .option("header", "true")
      .schema("nom STRING, age INT, salaire DOUBLE")
      .csv("s3a://bucket/clients.csv")

    Le schéma déclaré supprime la passe et fait échouer immédiatement toute donnée non conforme, ce qui est le comportement souhaitable.

    Le piège à éviter en réponse : ne pas dire « je vais simplement enlever l’option ». Sans inferSchema et sans schéma, tout est lu en String, donc le bug de comparaison devient permanent.

    Et la remarque qui clôt le sujet : tout ce débat ne concerne que le CSV. Parquet porte son schéma dans son pied de fichier, donc il n’y a rien à deviner.

  19. 19Qu’est-ce que le problème des petits fichiers ?Intermédiairepetits-fichiersparquetpartitions

    C’est l’accumulation de milliers de fichiers de quelques kilo-octets dans un lac de données, alors que le stockage et le moteur sont dimensionnés pour des fichiers de 100 Mo à 1 Go.

    Le coût se paie à trois endroits :

    • La liste des fichiers. Sur S3, énumérer 200 000 objets prend plus de temps que lire les données utiles. Cette phase se déroule sur le driver, donc elle ne se parallélise pas.
    • Une tâche par fichier. Spark crée au moins une partition par fichier, donc 200 000 tâches de 20 ms chacune. Le coût d’ordonnancement dépasse largement le calcul.
    • La compression perdue. Un fichier Parquet trop petit ne contient pas assez de valeurs pour que l’encodage par dictionnaire soit efficace, et ses statistiques par bloc ne servent presque à rien.

    La cause la plus fréquente : un job qui écrit avec 200 partitions par défaut, exécuté toutes les heures. Cela fait 4 800 fichiers par jour, dont l’essentiel est minuscule.

    Les corrections à citer :

    scala
    // À l'écriture : contrôler le nombre de fichiers produits
    resultat.coalesce(10).write.mode("append").parquet(chemin)
    
    // Ou, mieux sur du partitionné, viser une taille de fichier
    resultat.repartition($"date").write.partitionBy("date").parquet(chemin)

    Et surtout : ne pas partitionner sur une colonne à forte cardinalité. partitionBy("user_id") sur un million d’utilisateurs crée un million de dossiers — c’est la façon la plus rapide de créer le problème.

    Le point qui montre l’expérience : Delta Lake règle cela avec OPTIMIZE, qui compacte les petits fichiers en arrière-plan, et le compactage automatique sur Databricks.

  20. 20Qu’est-ce qu’une vue temporaire et que contient-elle ?Juniorspark-sqlvues

    Une vue temporaire est un nom enregistré dans le catalogue de la session, associé au plan logique d’un DataFrame. C’est ce qui permet de l’interroger en SQL.

    scala
    df.createOrReplaceTempView("personnes")
    spark.sql("SELECT * FROM personnes WHERE id = 1").show()

    La réponse qui compte, et que beaucoup de candidats manquent : elle ne contient aucune donnée. Ce n’est pas un cache, c’est un alias. Chaque requête sur la vue réexécute le plan depuis le début, donc deux requêtes sur une vue posée sur 400 Go relisent 400 Go deux fois. Pour matérialiser, il faut cache() sur le DataFrame sous-jacent.

    Les portées, si on demande la différence entre les variantes :

    MéthodeVisible depuisDisparaît
    createOrReplaceTempViewla session courantefin de session
    createGlobalTempViewtoutes les sessions de l’applicationfin de l’application
    saveAsTablele metastorejamais — les données sont écrites

    Le piège à mentionner : une vue globale vit dans une base système, donc il faut la préfixer — SELECT * FROM global_temp.personnes. Sans le préfixe, l’erreur est Table or view not found.

  21. 21Qu’est-ce que l’Adaptive Query Execution ?Senioraqecatalystshuffle

    L’AQE, introduite en Spark 3 et activée par défaut depuis la 3.2, permet à Spark de réviser son plan pendant l’exécution, à partir des statistiques réelles mesurées à chaque fin d’étape.

    C’est la réponse à une limite structurelle de Catalyst : il planifie avant d’avoir vu les données, donc à partir d’estimations qui peuvent être largement fausses — surtout sur des tables sans statistiques, ou après plusieurs transformations.

    Trois optimisations à nommer :

    1. Fusion des partitions de shuffle. Le fameux réglage à 200 devient largement inutile : Spark mesure la taille réelle après le shuffle et fusionne les partitions trop petites.
    2. Conversion de stratégie de join. Un sort-merge join planifié devient un broadcast join si la taille réellement mesurée passe sous le seuil.
    3. Gestion du skew. Les partitions anormalement grosses sont découpées en sous-partitions traitées en parallèle, ce qui traite le cas du straggler sans intervention.
    scala
    spark.conf.set("spark.sql.adaptive.enabled", "true")
    spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
    spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")

    La nuance qui montre la maturité : l’AQE ne dispense pas de comprendre les partitions. Elle corrige après coup ce qu’elle observe, mais elle ne peut agir qu’aux frontières de shuffle — elle ne rattrapera pas un partitionnement d’entrée catastrophique, ni un partitionBy sur une colonne à cardinalité démesurée.

  22. 22Comment lisez-vous un plan d’exécution Spark ?Seniorcatalystexplainperformance

    Avec explain(), et en le lisant de bas en haut : les feuilles de l’arbre sont les sources de données, la racine est le résultat final.

    scala
    df.explain(true)          // les quatre plans : parsé, analysé, optimisé, physique
    df.explain("formatted")   // le plan physique, lisible, avec les détails numérotés

    Les cinq éléments qu’on cherche en priorité :

    • PushedFilters sur le nœud de lecture. S’il est vide alors que votre requête filtre, le filtre s’applique après lecture de tout le fichier — le signe d’une UDF ou d’une expression que Catalyst ne peut pas traduire.
    • ReadSchema. Il doit ne contenir que les colonnes utiles. S’il contient tout, le column pruning n’a pas eu lieu.
    • La stratégie de join : BroadcastHashJoin, SortMergeJoin, ShuffleHashJoin, ou l’alarmant BroadcastNestedLoopJoin, qui signale une jointure sans condition d’égalité et une complexité quadratique.
    • Exchange. Chaque occurrence est un shuffle. Les compter donne le nombre d’étapes, et c’est le premier levier d’optimisation.
    • WholeStageCodegen. Les nœuds regroupés sous un même identifiant sont compilés en une seule boucle. Une étape qui en sort — typiquement à cause d’une UDF — perd cette fusion.

    Le complément à mentionner, parce qu’il distingue la théorie de la pratique : le plan dit ce que Spark prévoit, la Spark UI dit ce qui s’est passé. C’est là qu’on voit la durée réelle par étape, la distribution des durées de tâches, le volume de shuffle et les déversements sur disque. Avec l’AQE, le plan effectivement exécuté peut d’ailleurs différer du plan initial — l’onglet SQL de l’UI montre le plan final.

  23. 23Quels sont les rôles du driver et des executors ?Juniorarchitecturefondamentauxcluster

    Le driver est le processus qui exécute votre main. Il construit le plan, le découpe en étapes et en tâches, demande des ressources au gestionnaire de cluster, distribue les tâches et collecte les résultats. C’est lui qui détient le SparkContext.

    Les executors sont les processus JVM lancés sur les nœuds du cluster. Ils exécutent les tâches, stockent les données mises en cache, et servent les fichiers de shuffle aux autres executors.

    Entre les deux, le gestionnaire de cluster (YARN, Kubernetes, Mesos ou le mode standalone) attribue les ressources. Il n’exécute aucun calcul applicatif.

    Les deux conséquences pratiques qui font une bonne réponse :

    Le driver est un point de défaillance unique. S’il tombe, le job entier meurt. C’est aussi lui qui souffre de collect() sur un gros DataFrame — les données remontent dans sa mémoire — et du broadcast, qui passe par lui avant redistribution.

    Le code exécuté dans une transformation tourne sur les executors. D’où l’erreur classique : un println dans un map n’apparaît pas dans vos journaux, il part dans ceux de l’executor. Et toute variable capturée par une closure doit être sérialisable, sinon Task not serializable — c’est précisément ce que résolvent les variables diffusées et les accumulateurs.

  24. 24Qu’est-ce qu’un job, un stage et une tâche ?Intermédiairearchitectureshuffleexecution

    Trois niveaux, du plus grand au plus petit :

    • Job : tout le travail déclenché par une action. Deux count() produisent deux jobs.
    • Stage (étape) : un ensemble de transformations exécutables sans déplacer de données. Les frontières sont les shuffles.
    • Task (tâche) : l’exécution d’une étape sur une partition. C’est l’unité que l’ordonnanceur envoie à un cœur.

    La règle à énoncer, parce qu’elle prouve la compréhension : le nombre d’étapes d’un job est le nombre de shuffles plus un. Et le nombre de tâches d’une étape est le nombre de partitions qu’elle traite.

    action → 1 job
       ├── stage 0 : lecture + filter + map   (200 tâches)
       │        ⇅ shuffle
       └── stage 1 : agrégation + écriture    (200 tâches)

    Le point important sur l’enchaînement : les étapes sont séquentielles quand elles dépendent l’une de l’autre, à cause de la barrière du shuffle — l’étape 1 ne démarre pas avant que la dernière tâche de l’étape 0 soit terminée. C’est pour cela qu’une seule tâche lente retient tout le job, et c’est le mécanisme derrière le problème de skew.

    Deux étapes indépendantes, en revanche — deux branches d’un join, par exemple — peuvent s’exécuter en parallèle si les ressources le permettent.

  25. 25Executor en OutOfMemory : comment diagnostiquer ?Seniormemoirediagnosticperformance

    D’abord identifier le processus tombé : driver ou exécuteur, les causes n’ont rien en commun. Puis mesurer avant de rallonger la mémoire, qui est le dernier levier et pas le premier.

    Lire la réponse détaillée
  26. 26Qu’est-ce que Tungsten et le whole-stage codegen ?Seniortungstencatalystperformance

    Tungsten est le moteur d’exécution de Spark SQL. Là où Catalyst optimise quoi exécuter, Tungsten optimise comment, sur trois axes :

    1. La gestion mémoire hors tas. Les lignes sont stockées dans un format binaire compact, en dehors du tas JVM. Il n’y a plus un objet Java par ligne, donc plus de surcoût d’en-tête ni de pression sur le ramasse-miettes.
    2. Le cache locality. Les données sont disposées pour tenir dans le cache du processeur, ce qui réduit les défauts de cache.
    3. Le whole-stage code generation. C’est le point le plus important.

    Le codegen prend une étape entière — lecture, filtre, projection, agrégation partielle — et génère à la volée une seule fonction Java qui fait tout dans une seule boucle, compilée par le JIT.

    Sans codegen, chaque opérateur est un objet qui appelle le suivant par une méthode virtuelle, ligne par ligne : c’est le modèle volcano, où le coût des appels dépasse souvent le calcul utile. Avec codegen, ces appels disparaissent, remplacés par du code aussi direct qu’une boucle écrite à la main.

    C’est ce qu’on voit dans le plan physique sous WholeStageCodegen, avec un identifiant partagé par les nœuds fusionnés.

    La conséquence à énoncer, qui relie tout : dès qu’une UDF ou une opération RDD s’intercale, l’étape sort du codegen et retombe sur le modèle ligne par ligne, avec désérialisation vers des objets JVM. C’est la raison technique précise pour laquelle une UDF est dix fois plus lente qu’une fonction native.

  27. 27Spark SQL peut-il remplacer une base de données ?Intermédiairespark-sqlarchitecture

    Non, et la nuance à éviter est de répondre « parce que Spark est distribué » : les bases modernes le sont aussi. La différence est une question de cible de conception.

    Base transactionnelleSpark SQL
    Optimisé pourrequêtes ponctuellesbalayages de gros volumes
    Latence typiquemillisecondessecondes à minutes
    Indexcentraux (B-tree)aucun — élagage de partitions et statistiques
    ACIDcœur du produitvia Delta Lake ou Iceberg seulement
    UPDATE par lignenaturelcoûteux, voire impossible
    Stockageinterne, coupléfichiers sur un lac

    Concrètement : SELECT * FROM commandes WHERE id = 84213 est une mauvaise idée en Spark SQL — sans index, il balaiera les fichiers. Et demander à PostgreSQL d’agréger huit milliards de lignes en scannant 4 To est également une mauvaise idée.

    L’écart structurel à mentionner : Spark travaille sur des fichiers que vous possédez, pas sur un stockage propriétaire. La même table Parquet est lisible par Spark, Trino, DuckDB ou Athena. C’est la promesse du lac de données, et cela explique pourquoi Spark alimente les bases plutôt qu’il ne les remplace.

    La formule qui conclut bien : Spark SQL sert à préparer les données, la base sert à les servir. Un pipeline typique fait les deux — Spark transforme et écrit en Parquet, une base ou un entrepôt sert les requêtes de l’application.

  28. 28Qu’apporte Delta Lake par rapport à Parquet ?Intermédiairedelta-lakeparquetformats

    Delta Lake est du Parquet, plus un journal de transactions. Les données restent des fichiers Parquet ; ce qui s’ajoute est un dossier _delta_log qui décrit quels fichiers composent la table à chaque version.

    Ce journal apporte cinq choses que Parquet seul ne peut pas offrir :

    1. Les transactions ACID. Une écriture est atomique : un lecteur ne voit jamais un état partiel. En Parquet nu, un job interrompu laisse des fichiers orphelins que les lecteurs prennent pour des données.
    2. UPDATE, DELETE et MERGE. Indispensables pour un flux de type upsert, ou pour une suppression conforme au RGPD. En Parquet, il faut réécrire la partition entière à la main.
    3. Le voyage dans le temps. VERSION AS OF 12 permet de relire l’état d’hier, donc d’auditer et de revenir en arrière après un mauvais chargement.
    4. L’évolution de schéma contrôlée. Une colonne ajoutée est déclarée, pas subie ; une incompatibilité est refusée au lieu de corrompre la table.
    5. La lecture concurrente pendant l’écriture. Le journal isole les lecteurs de l’écriture en cours.

    S’y ajoute OPTIMIZE, qui compacte les petits fichiers, et ZORDER, qui coïmplante les valeurs proches pour rendre l’élagage plus efficace.

    Le contrepoint honnête, qui évite de paraître vendeur : Delta ajoute une dépendance et un journal à entretenir — VACUUM doit être passé pour supprimer les anciens fichiers, sinon le stockage gonfle indéfiniment. Pour une table écrite une fois puis lue en lecture seule, du Parquet simple suffit très bien.

  29. 29Comment dimensionnez-vous les executors d’un job Spark ?Seniorclusterconfigurationperformance

    La réponse attendue n’est pas un chiffre, c’est un raisonnement — et deux règles empiriques bien connues.

    4 à 5 cœurs par executor. Au-delà, le débit vers HDFS ou le stockage objet plafonne à cause de la contention sur les entrées-sorties. En dessous, on multiplie les JVM et on perd les bénéfices du partage de mémoire pour le cache et le broadcast.

    Pas plus de 32 Go de tas. Au-delà, la JVM perd la compression des pointeurs d’objets, et les pauses du ramasse-miettes deviennent une source de latence. Mieux vaut plusieurs executors moyens qu’un énorme.

    Le calcul concret, sur un nœud de 16 cœurs et 64 Go :

    1 cœur et ~1 Go réservés au système et au démon du gestionnaire
    → 15 cœurs, 63 Go utilisables
    → 3 executors de 5 cœurs
    → 63 / 3 = 21 Go par executor
    → moins l'overhead (~10 %) : spark.executor.memory ≈ 19g

    Le paramètre qu’on oublie systématiquement, et qu’il faut citer : spark.executor.memoryOverhead. Il couvre les tampons de shuffle, les buffers réseau et les processus Python en PySpark. Un conteneur tué par YARN ou Kubernetes avec le code 143 vient presque toujours de là, et pas du tas.

    Le dernier point, qui montre qu’on relie les réglages entre eux : le nombre total de cœurs détermine le parallélisme utile, donc le nombre de partitions doit suivre. Trois executors de 5 cœurs donnent 15 tâches en parallèle ; viser 30 à 45 partitions par étape est cohérent, et laisser spark.sql.shuffle.partitions à 200 ne l’est pas.

    Et pour finir : sur un cluster partagé, l’allocation dynamique (spark.dynamicAllocation.enabled) est souvent préférable à un dimensionnement fixe, puisqu’elle rend les executors inactifs au lieu de les réserver.

  30. 30Snowflake ou Databricks : comment choisir ?Seniordatabrickssnowflakearchitecture

    La distinction technique de base : Databricks exécute avec Spark, accéléré par le moteur vectorisé Photon ; Snowflake exécute avec son propre moteur propriétaire, sur des virtual warehouses dimensionnés par taille de T-shirt.

    Mais répondre uniquement sur le moteur passe à côté de la vraie question, qui est celle du type de charge :

    • Snowflake est excellent en SQL analytique sur données structurées, avec une administration quasi nulle et des performances très prévisibles. C’est le choix par défaut d’une équipe orientée SQL et BI.
    • Databricks est plus fort dès qu’on sort du SQL : traitement de fichiers bruts, machine learning, streaming, code Python ou Scala arbitraire, et données semi-structurées.

    Deux critères pratiques qui tranchent souvent plus vite que le débat technique :

    1. Où sont déjà vos données ? Si tout est en Parquet ou Delta sur S3, Databricks lit en place. Si tout est déjà chargé dans Snowflake, l’en sortir coûte cher.
    2. Quelles sont les compétences de l’équipe ? Un moteur excellent qu’une équipe ne sait pas régler produit de moins bons résultats qu’un moteur correct qu’elle maîtrise.

    Le piège à connaître sur Databricks, parce que c’est un vrai coût de production : Photon retombe silencieusement sur Spark JVM pour les opérations qu’il ne prend pas en charge — une UDF, certains types complexes. Vous payez le supplément Photon sans en bénéficier, et rien ne l’annonce sinon le plan d’exécution. C’est aussi pour cela qu’une UDF coûte deux fois sur cette plateforme.

    Et la nuance de contexte : les deux produits convergent. Snowflake a ajouté Snowpark et le support d’Iceberg ; Databricks a ajouté un entrepôt SQL sérieux. Le choix se joue de moins en moins sur les capacités et de plus en plus sur l’écosystème en place.