Une vue temporaire n’est ni une table, ni un cache. Sa portée, sa durée de vie, le préfixe global_temp qui piège tout le monde, et l’équivalent Scala de l’exemple PySpark.
Voici l’exercice classique : une liste Python, un RDD, un DataFrame, puis une requête SQL.
x = [(1, 'toto', 'yoyo'), (2, 'titi', 'jiji'),
(3, 'tata', 'gogo'), (4, 'tutu', 'nono')]
rdd = sc.parallelize(x)
dataframe = rdd.toDF(("id", "nom", "prenom"))
dataframe.createOrReplaceTempView("personnes")
spark.sql("SELECT * FROM personnes WHERE id = 1").show()Ça marche du premier coup, et c’est précisément le problème : createOrReplaceTempView a l’air d’un détail de syntaxe. Ce n’en est pas un. Il existe quatre méthodes voisines dont les portées diffèrent, et deux malentendus sur ce qu’une vue contient réellement coûtent cher en production.
C’est le malentendu principal, et il vaut la peine d’être posé avant tout le reste.
createOrReplaceTempView("personnes") n’exécute rien et ne stocke rien. Il enregistre un nom dans le catalogue de la session, associé au plan logique du DataFrame. Chaque spark.sql("SELECT … FROM personnes") réexécute ce plan depuis le début.
df = spark.read.parquet("s3a://bucket/gros-fichier") # 400 Go
df.createOrReplaceTempView("commandes")
spark.sql("SELECT COUNT(*) FROM commandes").show() # relit 400 Go
spark.sql("SELECT SUM(total) FROM commandes").show() # relit 400 GoDeux lectures complètes du lac, facturées deux fois. Une vue n’est pas un cache — c’est un alias. Si vous interrogez la même vue plusieurs fois, c’est le DataFrame sous-jacent qu’il faut matérialiser :
df.cache() # ou persist(StorageLevel.MEMORY_AND_DISK)
df.createOrReplaceTempView("commandes")Et encore faut-il le faire correctement : les niveaux de stockage, le unpersist oublié et le cas où seul checkpoint fonctionne sont détaillés dans cache, persist et checkpoint : quoi garder.
| Méthode | Visible depuis | Disparaît | Si le nom existe déjà |
|---|---|---|---|
createTempView | la session courante | fin de session | lève une exception |
createOrReplaceTempView | la session courante | fin de session | remplace en silence |
createGlobalTempView | toutes les sessions de l’application | fin de l’application | lève une exception |
createOrReplaceGlobalTempView | toutes les sessions | fin de l’application | remplace |
Dans 95 % des cas, createOrReplaceTempView est le bon choix — c’est celui qui rend un notebook réexécutable sans erreur. createTempView a un usage précis : détecter une collision de nom dans un pipeline où deux modules pourraient écraser la vue de l’autre sans prévenir.
Aucune des quatre ne survit à l’arrêt de l’application. Pour une table qui persiste entre deux exécutions, il faut le metastore :
df.write.mode("overwrite").saveAsTable("catalogue.schema.commandes")Là, les données sont réellement écrites, et la table reste interrogeable après redémarrage — par Spark, mais aussi par Trino, DuckDB ou Athena si le format le permet. C’est un tout autre engagement, traité dans écrire en Spark : Parquet, Delta Lake, petits fichiers.
global_tempLes vues globales sont stockées dans une base système, et il faut la préfixer. C’est l’erreur la plus fréquente sur le sujet :
df.createGlobalTempView("personnes")
spark.sql("SELECT * FROM personnes").show()
# AnalysisException: Table or view not found: personnes
spark.sql("SELECT * FROM global_temp.personnes").show() # ✓Le nom de cette base est configurable via spark.sql.globalTempDatabase, mais changez-le et plus aucun exemple trouvé en ligne ne fonctionnera chez vous. Laissez global_temp.
Le seul cas où une vue globale est vraiment utile : partager un jeu de données entre plusieurs sessions du même SparkContext, typiquement dans un serveur Spark multi-utilisateurs.
autre = spark.newSession()
autre.sql("SELECT * FROM personnes") # ✗ invisible
autre.sql("SELECT * FROM global_temp.personnes") # ✓ visiblePour nettoyer, ne comptez pas sur DROP VIEW seul :
spark.catalog.dropTempView("personnes")
spark.catalog.dropGlobalTempView("personnes")
spark.catalog.listTables() # inventorier ce qui est enregistréL’API est la même, la syntaxe diffère sur trois points.
import org.apache.spark.sql.SparkSession
val spark = SparkSession.builder()
.master("local[*]")
.appName("exemple1")
.getOrCreate()
import spark.implicits._ // 1. indispensable pour toDF
val sc = spark.sparkContext // 2. le contexte vient de la session
val x = Seq((1, "toto", "yoyo"), (2, "titi", "jiji"),
(3, "tata", "gogo"), (4, "tutu", "nono"))
val rdd = sc.parallelize(x)
val dataframe = rdd.toDF("id", "nom", "prenom") // 3. varargs, pas un tuple
dataframe.show()
dataframe.createOrReplaceTempView("personnes")
spark.sql("SELECT * FROM personnes WHERE id = 1").show()Les trois écarts à retenir :
import spark.implicits._ conditionne .toDF(), .as[T] et la syntaxe $"colonne". Sans lui, l’erreur est déroutante : value toDF is not a member of RDD[(Int, String, String)].SparkContext se récupère depuis la session (spark.sparkContext) au lieu d’être construit à part. La version PySpark de l’exercice fait sc = SparkContext() puis crée la session — ça fonctionne, mais c’est redondant et la distinction entre sc et spark explique pourquoi.toDF prend des varargs en Scala : toDF("id", "nom"), pas toDF(("id", "nom")). PySpark accepte les deux formes, Scala non.En Scala, un RDD n’est convertible en DataFrame que s’il contient un type Product — tuple ou case class. Pour du code lisible, la case class est nettement préférable :
case class Personne(id: Int, nom: String, prenom: String)
val ds = Seq(Personne(1, "toto", "yoyo"), Personne(2, "titi", "jiji")).toDS()
ds.createOrReplaceTempView("personnes")Le schéma est alors déduit de la classe, et vous obtenez un Dataset[Personne] typé à la compilation — l’option la plus sûre, comme détaillé dans RDD, DataFrame ou Dataset.
L’exemple de départ en contient deux autres, discrets mais coûteux.
master("local") n’utilise qu’un seul cœurspark = SparkSession.builder.master("local").appName("exemple1").getOrCreate()local signifie un thread. Pas « en local avec toutes les ressources » — un seul. Tout s’exécute en série, quel que soit le nombre de cœurs de la machine.
| Valeur | Threads utilisés |
|---|---|
local | 1 |
local[4] | 4 |
local[*] | tous les cœurs disponibles |
Sur quatre lignes de données, aucune importance. Sur un fichier d’un gigaoctet en développement, c’est la différence entre huit secondes et une minute — et cela masque tous les problèmes de parallélisme que vous auriez vus autrement. Écrivez local[*]. Le lien entre cette valeur et le nombre de partitions initiales est développé dans defaultParallelism : combien de cœurs Spark utilise-t-il.
read.json attend du JSON Lines, pas un tableau JSONdf = spark.read.json("/content/people.json")
df.printSchema()Par défaut, Spark attend un objet JSON par ligne :
{"name":"Michael"}
{"name":"Andy","age":30}Donnez-lui un tableau formaté sur plusieurs lignes — la sortie normale de n’importe quelle API REST — et le résultat n’est pas une erreur, c’est un schéma à une seule colonne :
root
|-- _corrupt_record: string (nullable = true)Le symptôme est déroutant parce que le job réussit. La correction tient en une option :
df = spark.read.option("multiLine", True).json("/content/people.json")Second point sur read.json : l’inférence de schéma déclenche une passe complète sur les données avant même votre première requête. Acceptable sur un fichier d’exemple, très cher sur un lac. En production, déclarez le schéma :
from pyspark.sql.types import StructType, StructField, StringType, LongType
schema = StructType([
StructField("name", StringType(), True),
StructField("age", LongType(), True),
])
df = spark.read.schema(schema).json("s3a://bucket/people/")Vous y gagnez la passe évitée, et surtout des types stables : sans schéma explicite, un champ vide dans le premier fichier et rempli dans le suivant peut changer de type d’une exécution à l’autre. Le mécanisme et ses conséquences — dont des comparaisons faussées qui ne lèvent aucune erreur — sont détaillés dans ce que l’option inferSchema coûte.
sc.stop() # NameError: name 'sc' is not defined
sc = SparkContext()Contrairement au spark-shell, un notebook Colab ne crée aucun sc automatiquement : il n’existe qu’après votre premier SparkContext(). La ligne fonctionne donc à la deuxième exécution de la cellule, jamais à la première — ce qui explique le nombre de tutoriels où elle apparaît sans commentaire, écrits par quelqu’un qui avait déjà lancé la cellule.
La forme robuste, valable partout :
spark = SparkSession.builder.master("local[*]").appName("exemple1").getOrCreate()
sc = spark.sparkContextgetOrCreate() réutilise ce qui existe au lieu d’échouer, puisqu’il n’y a qu’un seul SparkContext par JVM.
Une vue temporaire est un nom donné à un plan logique, pas un stockage : chaque requête réexécute tout, et il faut cache() pour matérialiser. createOrReplaceTempView couvre presque tous les besoins ; les vues globales servent uniquement à traverser plusieurs sessions du même contexte, et exigent le préfixe global_temp. Pour survivre au redémarrage, il faut saveAsTable, qui écrit réellement les données.
Côté portage PySpark vers Scala, trois réflexes suffisent : import spark.implicits._, spark.sparkContext plutôt qu’un contexte construit à part, et toDF en varargs. Et dans les deux langages, écrivez local[*] — pas local.
Développement et déploiement de solutions de données — les premiers modules sont en accès libre.