createOrReplaceTempView : les vues temporaires Spark

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.

7 min de lecturesparkpysparkscalaspark-sqlbig-data

Voici l’exercice classique : une liste Python, un RDD, un DataFrame, puis une requête SQL.

python
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.

Une vue temporaire ne contient pas de données

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.

python
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 Go

Deux 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 :

python
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.

Les quatre méthodes, et laquelle choisir

MéthodeVisible depuisDisparaîtSi le nom existe déjà
createTempViewla session courantefin de sessionlève une exception
createOrReplaceTempViewla session courantefin de sessionremplace en silence
createGlobalTempViewtoutes les sessions de l’applicationfin de l’applicationlève une exception
createOrReplaceGlobalTempViewtoutes les sessionsfin de l’applicationremplace

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 :

python
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.

Le piège global_temp

Les 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 :

python
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.

python
autre = spark.newSession()

autre.sql("SELECT * FROM personnes")               # ✗ invisible
autre.sql("SELECT * FROM global_temp.personnes")   # ✓ visible

Pour nettoyer, ne comptez pas sur DROP VIEW seul :

python
spark.catalog.dropTempView("personnes")
spark.catalog.dropGlobalTempView("personnes")
spark.catalog.listTables()          # inventorier ce qui est enregistré

L’équivalent Scala, ligne par ligne

L’API est la même, la syntaxe diffère sur trois points.

scala
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 :

  1. 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)].
  2. Le 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.
  3. 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 :

scala
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.

Deux pièges du même exercice

L’exemple de départ en contient deux autres, discrets mais coûteux.

master("local") n’utilise qu’un seul cœur

python
spark = 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.

ValeurThreads utilisés
local1
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 JSON

python
df = spark.read.json("/content/people.json")
df.printSchema()

Par défaut, Spark attend un objet JSON par ligne :

json
{"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 :

python
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 :

python
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.

Pourquoi `sc.stop()` échoue au premier lancement dans Colab
python
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 :

python
spark = SparkSession.builder.master("local[*]").appName("exemple1").getOrCreate()
sc = spark.sparkContext

getOrCreate() réutilise ce qui existe au lieu d’échouer, puisqu’il n’y a qu’un seul SparkContext par JVM.

Ce qu’il faut retenir

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.

Ce sujet fait partie d’un cours complet

Développement et déploiement de solutions de données — les premiers modules sont en accès libre.

Voir le plan du cours

Continuer sur le même sujet