Spark : pourquoi votre job ne fait rien avant une action

map, filter, select ne lancent aucun calcul. count, collect, write si. Comprendre la frontière transformation / action évite 80 % des mauvaises surprises en Spark.

3 min de lecturesparkscalabig-data

Vous enchaînez cinq transformations Spark, vous regardez l’UI, et… rien. Pas d’étape, pas de tâche, pas de CPU. Ce n’est pas un bug. C’est le modèle d’exécution : Spark est paresseux. Il construit un plan. Il ne calcule que lorsqu’une action exige un résultat.

Tant que cette frontière n’est pas claire, on « optimise » des pipelines qui ne tournent même pas — ou on déclenche dix fois le même calcul sans le savoir.

Transformations vs actions

TypeExemplesEffet
Transformationmap, filter, select, join, groupByAjoute une étape au plan
Actioncount, collect, show, write, takeDéclenche l’exécution
scala
val logs = spark.read.parquet("s3a://logs/2026/")
val errors = logs.filter($"level" === "ERROR")
val byService = errors.groupBy("service").count()

// Jusqu'ici : aucun job Spark.
byService.show() // Première action → le cluster travaille.

Le DAG, pas la ligne de code

Quand l’action part, Spark assemble un DAG (graphe acyclique) des étapes nécessaires, le découpe en stages selon les shuffles, puis en tasks.

Conséquence pratique : une transformation coûteuse placée après un filtre agressif coûte beaucoup moins. Spark ne lit pas « dans l’ordre de vos habitudes », il lit le plan.

L’anti-pattern : plusieurs actions sur le même DataFrame

scala
val errors = logs.filter($"level" === "ERROR")
println(errors.count())          // job 1 : lit + filtre
errors.write.parquet("/out")     // job 2 : relit + refiltre

Sans cache, chaque action rejoue le plan depuis la source. Sur un lac de données, c’est la différence entre dix minutes et une heure.

scala
val errors = logs.filter($"level" === "ERROR").cache()
errors.count()                 // matérialise
errors.write.parquet("/out")   // réutilise

cache() n’est pas gratuit non plus : il occupe de la mémoire. On le met là où le DataFrame est réutilisé, pas par réflexe. Les niveaux de stockage, l’unpersist() que tout le monde oublie et le cas où seul checkpoint() règle le problème sont détaillés dans cache, persist et checkpoint.

Scala ou PySpark : la frontière est la même
python
errors = logs.filter(logs.level == "ERROR")
errors.count()  # action
errors.write.parquet("/out")  # autre action

La syntaxe change ; le contrat lazy/eager, non. C’est pour cela qu’un aide-mémoire qui aligne les deux langages est utile : les pièges sont partagés. Sur le choix du langage lui-même — et sur les trois seuls endroits où la performance diffère réellement — voir Scala ou PySpark.

Comment lire un plan

scala
byService.explain(true)

Cherchez :

  • les scans trop larges (colonnes ou partitions inutiles)
  • les Exchange (shuffles) — souvent le vrai coût
  • les filtres poussés (PushedFilters) — signe que la source fait le travail

Ce qu’il faut retenir

En Spark, écrire du code n’exécute rien. Demander un résultat exécute tout le plan nécessaire. Une fois ce réflexe acquis, vous arrêtez de debugger le cluster pour des jobs qui n’ont pas encore démarré — et vous commencez à compter les actions comme on compte les allers-retours réseau.

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