map, filter, select ne lancent aucun calcul. count, collect, write si. Comprendre la frontière transformation / action évite 80 % des mauvaises surprises en Spark.
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.
| Type | Exemples | Effet |
|---|---|---|
| Transformation | map, filter, select, join, groupBy | Ajoute une étape au plan |
| Action | count, collect, show, write, take | Déclenche l’exécution |
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.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.
val errors = logs.filter($"level" === "ERROR")
println(errors.count()) // job 1 : lit + filtre
errors.write.parquet("/out") // job 2 : relit + refiltreSans 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.
val errors = logs.filter($"level" === "ERROR").cache()
errors.count() // matérialise
errors.write.parquet("/out") // réutilisecache() 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.
errors = logs.filter(logs.level == "ERROR")
errors.count() # action
errors.write.parquet("/out") # autre actionLa 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.
byService.explain(true)Cherchez :
PushedFilters) — signe que la source fait le travailEn 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.
Développement et déploiement de solutions de données — les premiers modules sont en accès libre.