Une UDF est une boîte noire pour Catalyst : plus de codegen, plus de pushdown. Ce que vous perdez, et comment la remplacer par des fonctions natives.
Vous écrivez une petite fonction Scala, vous l’enregistrez comme UDF, elle fait exactement ce qu’il faut. Puis le job passe de quatre minutes à trente-cinq. Ce n’est pas votre logique qui est lente — c’est que vous venez de rendre aveugle l’optimiseur de Spark sur cette colonne.
Une UDF n’est pas interdite. Mais elle a un coût précis, mesurable, et évitable dans la grande majorité des cas. Voici lequel, et comment s’en passer.
Avec une fonction native, Spark lit l’expression :
import org.apache.spark.sql.functions._
df.withColumn("majuscule", upper($"nom"))Catalyst sait que c’est upper, connaît son comportement sur les null, peut la fusionner avec les opérations voisines dans une seule boucle compilée, et peut décider de pousser le filtre associé jusqu’au fichier Parquet.
Avec une UDF :
val majusculeUdf = udf((s: String) => s.toUpperCase)
df.withColumn("majuscule", majusculeUdf($"nom"))Catalyst voit une boîte noire. Il ne peut plus rien en déduire. Concrètement, vous perdez quatre choses :
| Ce que vous perdez | Conséquence |
|---|---|
| Whole-stage codegen | Plus de bytecode fusionné : un appel de fonction par ligne, au lieu d’une boucle compilée |
| Format Tungsten | Chaque valeur est désérialisée en objet JVM, passée à la fonction, puis re-sérialisée |
| Predicate pushdown | Un WHERE udf(col) = x ne peut pas être poussé vers la source : Spark lit tout |
Gestion des null | À votre charge — une UDF qui reçoit null lève une NullPointerException |
Le troisième point est souvent le plus coûteux : un filtre qui aurait éliminé 99 % des données au niveau du fichier doit maintenant charger 100 % des lignes pour les évaluer une par une.
Sur Databricks, la facture est double : une UDF fait aussi retomber l’étape sur Spark JVM au lieu du moteur Photon, que vous payez pourtant. Le détail est dans Snowflake vs Databricks : Spark, Photon et le moteur.
La plupart des UDF écrites en production ont un équivalent natif. Voici les cas les plus fréquents :
| Ce que fait votre UDF | Fonction native |
|---|---|
if / else sur une valeur | when(...).otherwise(...) |
Remplacer les null | coalesce(col, lit(défaut)), nvl |
| Extraire par expression régulière | regexp_extract(col, motif, groupe) |
| Remplacer par regex | regexp_replace |
| Découper une chaîne | split, puis element_at ou getItem |
| Parser une date | to_date, to_timestamp, date_format |
| Concaténer | concat, concat_ws |
| Tester l’appartenance à une liste | isin(...), ou un broadcast join |
| Manipuler un tableau | array_contains, transform, filter, explode |
| Parser du JSON | from_json avec un schéma, get_json_object |
| Hasher | md5, sha2, xxhash64 |
| Arithmétique conditionnelle | Expressions directes, least, greatest |
Exemple typique, avant :
val categoriser = udf((montant: Double) =>
if (montant > 1000) "grand" else if (montant > 100) "moyen" else "petit"
)
df.withColumn("categorie", categoriser($"montant"))Après — même résultat, entièrement optimisable :
df.withColumn("categorie",
when($"montant" > 1000, "grand")
.when($"montant" > 100, "moyen")
.otherwise("petit")
)Depuis Spark 3, les fonctions d’ordre supérieur couvrent aussi les cas qui obligeaient historiquement à écrire une UDF sur les tableaux :
// Doubler chaque élément d'un tableau, sans UDF
df.withColumn("doubles", transform($"valeurs", x => x * 2))
// Filtrer un tableau
df.withColumn("positifs", filter($"valeurs", x => x > 0))Il reste des cas légitimes, et il faut savoir les reconnaître :
Dans ce cas, réduisez les dégâts :
// 1. Typer explicitement le retour : évite une inférence coûteuse
val monUdf = udf[String, String]((s: String) => traiter(s))
// 2. Filtrer AVANT l'UDF, jamais après
df.filter($"actif" === true) // ← pushdown préservé
.withColumn("resultat", monUdf($"donnee"))
// 3. Gérer les null soi-même
val sûr = udf((s: String) => if (s == null) null else traiter(s))L’ordre du point 2 est le levier le plus important : appliquer l’UDF sur les 2 % de lignes qui survivent au filtre plutôt que sur 100 % change tout, et ne demande qu’un réarrangement de lignes.
Une UDF Python ne tourne pas dans la JVM. Chaque ligne est sérialisée depuis la JVM, envoyée à un processus Python, traitée, puis re-sérialisée en retour. Le surcoût est d’un ordre de grandeur au-dessus d’une UDF Scala.
La réponse est pandas_udf, qui utilise Apache Arrow pour transférer des lots de lignes en format colonnaire :
from pyspark.sql.functions import pandas_udf
import pandas as pd
@pandas_udf("double")
def normaliser(s: pd.Series) -> pd.Series:
return (s - s.mean()) / s.std()
df.withColumn("norm", normaliser("valeur"))Typiquement 10 à 100 fois plus rapide qu’une UDF Python ligne à ligne, parce que la sérialisation est amortie sur des milliers de lignes et que le calcul est vectorisé par NumPy. Ça ne rend pas la fonction visible à Catalyst pour autant — mais ça élimine l’essentiel du surcoût.
Le plan explain montre littéralement l’endroit où l’optimisation s’arrête :
df.withColumn("x", monUdf($"col")).explain(true)
// == Physical Plan ==
// *(1) Project [..., UDF(col#12) AS x#20]
// ^^^ boîte noire : pas de codegen sur cette expressionComparez avec la version native, où l’expression apparaît en clair (upper(col#12), CASE WHEN ...) et est intégrée dans l’étage compilé. Et si vous hésitez entre deux écritures, mesurez sur un échantillon représentatif — pas sur dix lignes, où tout paraît instantané.
Une UDF échange de la lisibilité contre de l’optimisation, et le taux de change est mauvais : vous perdez le codegen, le format binaire, le pushdown et la gestion des null, souvent pour du code qui existait déjà dans org.apache.spark.sql.functions. Le réflexe à installer : chercher la fonction native d’abord, l’UDF seulement quand on peut expliquer pourquoi aucune ne convient.
Et quand l’UDF est inévitable, deux gestes suffisent à limiter la casse : filtrer en amont pour réduire le nombre d’appels, et gérer les null explicitement. Le contexte général de cette optimisation — pourquoi Catalyst est le vrai moteur de performance de Spark — est développé dans RDD ou DataFrame.
Développement et déploiement de solutions de données — les premiers modules sont en accès libre.